mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 16:02:29 +00:00
Keep guest listeners detached during restart and simplify byte streams
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
import { runCoalescedDaemonRestart } from '../../../daemon/daemon-restart-state'
|
||||
import { afterEach, expect, it, vi } from 'vitest'
|
||||
import { bindProviderListeners } from './bind-listeners'
|
||||
import { setRebindProviderListeners, unbindLocalProviderListeners } from './listener-lifecycle'
|
||||
@@ -67,3 +68,31 @@ it('binds added guests once, excludes SSH and removes listeners across reload an
|
||||
unbindLocalProviderListeners()
|
||||
expect(native.listeners.size).toBe(0)
|
||||
})
|
||||
|
||||
it('does not rebind a retiring local provider when guests change during restart', async () => {
|
||||
const rebind = vi.fn()
|
||||
setRebindProviderListeners(rebind)
|
||||
let finish!: () => void
|
||||
const pending = runCoalescedDaemonRestart(async () => {
|
||||
await new Promise<void>((resolve) => {
|
||||
finish = resolve
|
||||
})
|
||||
return { killedCount: 0 }
|
||||
})
|
||||
try {
|
||||
const release = registerWslPtyProvider(
|
||||
{ distro: 'Ubuntu', relayBuildId: 'during-restart' },
|
||||
provider().value
|
||||
)
|
||||
releases.push(release)
|
||||
release()
|
||||
expect(rebind).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
finish()
|
||||
await pending
|
||||
}
|
||||
releases.push(
|
||||
registerWslPtyProvider({ distro: 'Ubuntu', relayBuildId: 'after-restart' }, provider().value)
|
||||
)
|
||||
expect(rebind).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { isDaemonRestartInFlight } from '../../../daemon/daemon-restart-state'
|
||||
import { parseAppWslPtyId, type WslPtyOwner } from '../../../../shared/wsl-pty-id'
|
||||
import { wslPtyOwnerKey } from '../../../../shared/wsl-pty-consumer-recovery'
|
||||
import { relayProvidersByGeneration } from '../../../providers/relay-pty-generation-registry'
|
||||
@@ -50,11 +51,15 @@ export function registerWslPtyProvider(owner: WslPtyOwner, provider: IPtyProvide
|
||||
throw new Error('WSL terminal owner already registered')
|
||||
}
|
||||
wslProviders.set(key, provider)
|
||||
rebindLocalProviderListeners()
|
||||
if (!isDaemonRestartInFlight()) {
|
||||
rebindLocalProviderListeners()
|
||||
}
|
||||
return () => {
|
||||
if (wslProviders.get(key) === provider) {
|
||||
wslProviders.delete(key)
|
||||
rebindLocalProviderListeners()
|
||||
if (!isDaemonRestartInFlight()) {
|
||||
rebindLocalProviderListeners()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -164,9 +169,7 @@ export function getSshPtyProvider(connectionId: string): IPtyProvider | undefine
|
||||
return sshProviders.get(connectionId)
|
||||
}
|
||||
|
||||
/** Get the installed PTY provider (for direct access in tests/runtime).
|
||||
* After daemon init this may be a DaemonPtyAdapter/DaemonPtyRouter, not LocalPtyProvider;
|
||||
* callers needing LocalPtyProvider-specific methods must type-narrow or import the class. */
|
||||
/** Get the installed daemon provider for runtime operations and tests. */
|
||||
export function getLocalPtyProvider(): IPtyProvider {
|
||||
return localProvider
|
||||
}
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { DaemonConnectionLostError } from '../daemon/daemon-errors'
|
||||
import { afterEach, expect, it } from 'vitest'
|
||||
import { createServer } from 'node:net'
|
||||
import { mkdtempSync, rmSync, writeFileSync, symlinkSync } from 'node:fs'
|
||||
@@ -51,6 +52,8 @@ unix(
|
||||
},
|
||||
AbortSignal.timeout(10_000)
|
||||
)
|
||||
expect(stream.readableObjectMode).toBe(false)
|
||||
expect(stream.writableObjectMode).toBe(false)
|
||||
const chunks: Buffer[] = []
|
||||
stream.on('data', (chunk: Buffer) => chunks.push(chunk))
|
||||
const ended = once(stream, 'end')
|
||||
@@ -72,7 +75,7 @@ unix('rejects changed guest identity before connecting', async () => {
|
||||
},
|
||||
AbortSignal.timeout(10_000)
|
||||
)
|
||||
).rejects.toThrow('closed before connection')
|
||||
).rejects.toBeInstanceOf(DaemonConnectionLostError)
|
||||
})
|
||||
|
||||
it('cancels a connector that never reaches readiness', async () => {
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { Duplex, Readable, Writable } from 'node:stream'
|
||||
import { DaemonConnectionLostError } from '../daemon/daemon-errors'
|
||||
import { Duplex } from 'node:stream'
|
||||
import { spawnProcess, type ProcessSpec } from '../../shared/child-process/run-process'
|
||||
import { WSL_DAEMON_CONNECTOR_READY } from './wsl-daemon-connector-script'
|
||||
|
||||
@@ -9,10 +10,7 @@ export function openWslDaemonConnectorStream(
|
||||
): Promise<Duplex> {
|
||||
signal.throwIfAborted()
|
||||
const child = spawnProcess(spec)
|
||||
const stream = Duplex.fromWeb(
|
||||
{ readable: Readable.toWeb(child.stdout), writable: Writable.toWeb(child.stdin) },
|
||||
{ objectMode: false }
|
||||
)
|
||||
const stream = Duplex.from({ readable: child.stdout, writable: child.stdin })
|
||||
// A child can fail before the awaiting protocol consumer has installed its listener.
|
||||
stream.on('error', () => {})
|
||||
stream.once('close', () => child.kill())
|
||||
@@ -37,7 +35,7 @@ export function openWslDaemonConnectorStream(
|
||||
stream.destroy()
|
||||
reject(error)
|
||||
}
|
||||
const closed = () => fail(new Error('WSL daemon connector closed before connection'))
|
||||
const closed = () => fail(new DaemonConnectionLostError('Connection lost'))
|
||||
const abort = () =>
|
||||
fail(new Error('WSL daemon connector connection canceled', { cause: signal.reason }))
|
||||
const onStderr = (bytes: Buffer) => {
|
||||
|
||||
Reference in New Issue
Block a user