Merge pull request #3384 from deepseek-harness/worktree-connerr
fix(connection): avoid disconnecting during host stalls
This commit is contained in:
commit
8b4b815e7e
5 changed files with 56 additions and 17 deletions
|
|
@ -111,7 +111,7 @@ describe.skipIf(MODE === 'record')('web e2e: background job list', () => {
|
|||
onTestFailed(() => saveFailureShot(page, 'web-e2e-background-job-settled'))
|
||||
expect(scaffold.ctx.jobs.kill(jobId, agent, 'web e2e cancellation')).toBe('requested')
|
||||
|
||||
const idle = page.getByRole('button', { name: '1 background job' })
|
||||
const idle = page.getByRole('button', { name: '1 background job', exact: true })
|
||||
await idle.waitFor({ timeout: 20_000 })
|
||||
|
||||
const snapshot = await captureStableAria(page, '[class*="menu"]', scaffold.workspaceCwd)
|
||||
|
|
|
|||
|
|
@ -19,11 +19,13 @@ export type RemoteStreamOpener = (
|
|||
/** Convert an invocation or carrier failure to a stable wire value. */
|
||||
export type RemoteStreamFailureMapper = (error: unknown) => RemoteStreamFailure
|
||||
|
||||
const MAX_MISSED_HEARTBEATS = 2
|
||||
|
||||
/** Own the no-server WebSocket acceptor and every active logical stream. */
|
||||
export class RemoteStreamMuxServer {
|
||||
private readonly server = new WebSocketServer({ noServer: true })
|
||||
private readonly connections = new Set<Promise<void>>()
|
||||
private readonly heartbeatAlive = new WeakMap<WebSocket, boolean>()
|
||||
private readonly missedHeartbeats = new WeakMap<WebSocket, number>()
|
||||
private heartbeatTimer: NodeJS.Timeout | undefined
|
||||
|
||||
/**
|
||||
|
|
@ -45,8 +47,8 @@ export class RemoteStreamMuxServer {
|
|||
*/
|
||||
handleUpgrade(req: IncomingMessage, socket: Duplex, head: Buffer): void {
|
||||
this.server.handleUpgrade(req, socket, head, (websocket) => {
|
||||
this.heartbeatAlive.set(websocket, true)
|
||||
websocket.on('pong', () => { this.heartbeatAlive.set(websocket, true) })
|
||||
this.missedHeartbeats.set(websocket, 0)
|
||||
websocket.on('pong', () => { this.missedHeartbeats.set(websocket, 0) })
|
||||
this.startHeartbeat()
|
||||
const connection = new RemoteStreamMuxConnection(websocket, this.open, this.failure)
|
||||
const done = connection.run()
|
||||
|
|
@ -75,11 +77,16 @@ export class RemoteStreamMuxServer {
|
|||
this.heartbeatTimer = setInterval(() => {
|
||||
for (const socket of this.server.clients) {
|
||||
if (socket.readyState !== WebSocket.OPEN) continue
|
||||
if (this.heartbeatAlive.get(socket) === false) {
|
||||
const missed = this.missedHeartbeats.get(socket) as number
|
||||
if (missed >= MAX_MISSED_HEARTBEATS) {
|
||||
setImmediate(() => {
|
||||
if ((this.missedHeartbeats.get(socket) as number) >= MAX_MISSED_HEARTBEATS) {
|
||||
socket.terminate()
|
||||
}
|
||||
})
|
||||
continue
|
||||
}
|
||||
this.heartbeatAlive.set(socket, false)
|
||||
this.missedHeartbeats.set(socket, missed + 1)
|
||||
socket.ping()
|
||||
}
|
||||
}, this.heartbeatIntervalMs)
|
||||
|
|
|
|||
|
|
@ -50,7 +50,7 @@ describe('Remote stream mux server carrier lifecycle', () => {
|
|||
await closed
|
||||
})
|
||||
|
||||
it('terminates a socket that does not answer the previous heartbeat', async () => {
|
||||
it('requires two missed heartbeats before terminating an unresponsive socket', async () => {
|
||||
const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal), 20)
|
||||
const client = await connect(entry.url)
|
||||
const serverSocket = acceptedSocket(entry.mux)
|
||||
|
|
@ -58,10 +58,39 @@ describe('Remote stream mux server carrier lifecycle', () => {
|
|||
const terminated = vi.spyOn(serverSocket, 'terminate')
|
||||
const closed = once(client, 'close')
|
||||
|
||||
await once(client, 'ping')
|
||||
await once(client, 'ping')
|
||||
expect(terminated).not.toHaveBeenCalled()
|
||||
await vi.waitFor(() => { expect(terminated).toHaveBeenCalledOnce() })
|
||||
await closed
|
||||
})
|
||||
|
||||
it('keeps the socket when a delayed Pong arrives before the final check', async () => {
|
||||
const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal), 20)
|
||||
const client = await connect(entry.url, false)
|
||||
const serverSocket = acceptedSocket(entry.mux)
|
||||
const terminated = vi.spyOn(serverSocket, 'terminate')
|
||||
let finalCheck: (() => void) | undefined
|
||||
const immediate = vi.spyOn(globalThis, 'setImmediate').mockImplementation((callback) => {
|
||||
finalCheck = callback
|
||||
return 0 as unknown as NodeJS.Immediate
|
||||
})
|
||||
|
||||
try {
|
||||
await once(client, 'ping')
|
||||
await once(client, 'ping')
|
||||
await vi.waitFor(() => { expect(finalCheck).toBeDefined() })
|
||||
serverSocket.emit('pong', Buffer.alloc(0))
|
||||
finalCheck?.()
|
||||
expect(terminated).not.toHaveBeenCalled()
|
||||
} finally {
|
||||
immediate.mockRestore()
|
||||
const closed = once(client, 'close')
|
||||
client.close()
|
||||
await closed
|
||||
}
|
||||
})
|
||||
|
||||
it('rejects binary, malformed, and duplicate logical-stream messages', async () => {
|
||||
const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal))
|
||||
|
||||
|
|
@ -223,8 +252,8 @@ async function startMux(open: RemoteStreamOpener, heartbeatIntervalMs = 2_000):
|
|||
return entry
|
||||
}
|
||||
|
||||
async function connect(url: string): Promise<WebSocket> {
|
||||
const socket = new WebSocket(url)
|
||||
async function connect(url: string, autoPong = true): Promise<WebSocket> {
|
||||
const socket = new WebSocket(url, { autoPong })
|
||||
await once(socket, 'open')
|
||||
return socket
|
||||
}
|
||||
|
|
|
|||
|
|
@ -278,8 +278,7 @@ export class ConnectionController {
|
|||
this.callSink(() => { this.sinks.onConnected?.(host) })
|
||||
}
|
||||
} catch {
|
||||
// Transport failure: treat as generation failure, then enter the shared retry path.
|
||||
if (!ac.signal.aborted) ac.abort()
|
||||
// Source settlement and controller cancellation already abort the generation.
|
||||
}
|
||||
|
||||
await failed
|
||||
|
|
@ -306,12 +305,12 @@ export class ConnectionController {
|
|||
}
|
||||
}
|
||||
|
||||
/** Await source readiness without letting a stalled carrier wedge startup forever. */
|
||||
/** Await source readiness while reporting, but not cancelling, a slow Host. */
|
||||
function waitForReady<T>(ready: Promise<T>, timeoutMs: number, signal: AbortSignal): Promise<T> {
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
let settled = false
|
||||
const timeout = setTimeout(() => {
|
||||
finish({ error: new Error(`connection generation was not ready within ${String(timeoutMs)}ms`) })
|
||||
console.warn(`[connection] generation is still not ready after ${String(timeoutMs)}ms`)
|
||||
}, timeoutMs)
|
||||
const aborted = (): void => {
|
||||
finish({ error: new Error('connection generation aborted', { cause: signal.reason }) })
|
||||
|
|
|
|||
|
|
@ -568,7 +568,8 @@ describe('connection lifecycle', () => {
|
|||
}
|
||||
})
|
||||
|
||||
it('rejects and retries a generation whose source never reports ready', async () => {
|
||||
it('reports but retains a generation whose source is slow to report ready', async () => {
|
||||
vi.useFakeTimers()
|
||||
const source = new FakeGenerationSource()
|
||||
source.suppressReady = true
|
||||
let connected = 0
|
||||
|
|
@ -580,13 +581,16 @@ describe('connection lifecycle', () => {
|
|||
)
|
||||
controller.start()
|
||||
try {
|
||||
await Promise.resolve()
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(source.activeCount).toBe(1)
|
||||
await new Promise(resolve => setTimeout(resolve, 45))
|
||||
await vi.advanceTimersByTimeAsync(20)
|
||||
expect(connected).toBe(0)
|
||||
expect(source.activeCount).toBe(1)
|
||||
expect(warnSpy).toHaveBeenCalledWith('[connection] generation is still not ready after 20ms')
|
||||
} finally {
|
||||
controller.stop()
|
||||
warnSpy.mockRestore()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue