ee4cad3ada
The bridge now holds each session's `AgentHandle` disposer in its `SessionRecord` and runs it on teardown (client disconnect or fiber dispose) instead of the old `abort()` + `whenIdle()` drain that left agents registered. A bare client disconnect now leaves NO registered agent and NO session-store entry — not an idled-but-still-registered one. The queue-aware `cancel()` inside the disposer also closes the former pre-step best-effort window (a turn about to start is dropped), so teardown reaches true quiescence. The `session/load`-races-teardown leak is fixed: if the bridge closed while `resume()` was pending, the just-resumed handle is disposed before throwing, so it leaves no orphan (it has no SessionRecord, so quiesce() never sees it). Tests: the disconnect test now asserts (through the SAME memoized teardown) that the agent is unregistered AND its session removed; a durability test re-loads the persisted log after dispose and asserts the closing turn/end is on disk (guards the teardown-order contract); a sibling-isolation test proves one handle's dispose() leaves other agents untouched. Docs: agent / agent-loop / acp READMEs, architecture.md, and the stale in-code quiesce() ownership comment updated to the per-agent disposal model; the now-resolved TODO(rfc010-agent-disposal) / TODO(rfc010-cancel-prestep) teardown notes removed.
215 lines
12 KiB
TypeScript
215 lines
12 KiB
TypeScript
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
|
import { mkdtemp, rm } from 'node:fs/promises'
|
|
import { tmpdir } from 'node:os'
|
|
import { join } from 'node:path'
|
|
import { PROTOCOL_VERSION } from '@agentclientprotocol/sdk'
|
|
import { SessionId } from '@deepseek-ai/dsh-session'
|
|
import { makeBridgeHarness, textResponse } from './harness.ts'
|
|
|
|
describe('acp bridge — disposal & HMR safety', () => {
|
|
let storageDir: string
|
|
|
|
beforeEach(async () => { storageDir = await mkdtemp(join(tmpdir(), 'acp-dispose-')) })
|
|
afterEach(async () => { await rm(storageDir, { recursive: true, force: true }) })
|
|
|
|
it('disposal reaches quiescence: a running turn is aborted and awaited before dispose returns', async () => {
|
|
const harness = await makeBridgeHarness({ storageDir, script: ['hang'] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
const agent = harness.ctx.agents.get(sessionId)!
|
|
|
|
// Start a prompt that hangs in the model stream.
|
|
const promptDone = harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })
|
|
await new Promise(r => setTimeout(r, 30))
|
|
expect(agent.status).toBe('running')
|
|
|
|
// Dispose the whole context. The bridge's teardown must abort the agent and
|
|
// AWAIT whenIdle() — so right after dispose resolves, the agent is settled
|
|
// (not still running). Proves disposal waited, not just requested.
|
|
await harness.ctx.fiber.dispose()
|
|
expect(agent.status).not.toBe('running')
|
|
|
|
// The in-flight prompt settled (cancelled) rather than hanging forever.
|
|
const res = await promptDone
|
|
expect(res.stopReason).toBe('cancelled')
|
|
})
|
|
|
|
it('after an ACP-only HMR dispose, a late session/new creates no orphan agent (closed guard)', async () => {
|
|
// Dispose JUST the bridge's fiber (an HMR reload) while agents/agent-loop
|
|
// stay up and the transport is still live. A late session/new must hit the
|
|
// `closed` guard and reject — NOT create an agent the disposed bridge can no
|
|
// longer stream or settle. Verify the world: no agent appeared.
|
|
const harness = await makeBridgeHarness({ storageDir, script: [] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const before = harness.ctx.agents.list().length
|
|
await harness.acpFiber.dispose() // tear down ONLY the bridge
|
|
await expect(harness.client.newSession({ cwd: process.cwd(), mcpServers: [] }))
|
|
.rejects.toThrow(/disposed/)
|
|
expect(harness.ctx.agents.list().length).toBe(before)
|
|
await harness.dispose()
|
|
})
|
|
|
|
it('an agent created through the bridge is unregistered when ONLY the bridge fiber is disposed', async () => {
|
|
// The factory (`ctx.agents.create`) is reached through the bridge's
|
|
// traceable service proxy, so `AgentLoop.start`'s `this.ctx.effect(...)`
|
|
// registration binds to the CALLER context — the bridge fiber — not the
|
|
// AgentLoop fiber. Disposing JUST the bridge fiber (an ACP-only HMR reload)
|
|
// must therefore reclaim the agent's registry entry, even though agents/
|
|
// agent-loop stay up. This pins the fiber-ownership the bridge's teardown
|
|
// doc comment relies on; if a refactor rebinds the registration to the
|
|
// AgentLoop fiber, the agent would survive bridge dispose and this fails.
|
|
const harness = await makeBridgeHarness({ storageDir, script: [] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
expect(harness.ctx.agents.get(sessionId)).toBeDefined()
|
|
|
|
await harness.acpFiber.dispose() // tear down ONLY the bridge
|
|
expect(harness.ctx.agents.get(sessionId)).toBeUndefined()
|
|
await harness.dispose()
|
|
})
|
|
|
|
it('no agent is created by a session/new after the bridge has closed (closed guard)', async () => {
|
|
// After teardown (here a client disconnect sets `closed`), a late
|
|
// `session/new` must NOT create an orphan agent the bridge can no longer
|
|
// drive/settle. The transport is gone so the RPC rejects; assert the world:
|
|
// no new agent appeared in the registry.
|
|
const harness = await makeBridgeHarness({ storageDir, script: [] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const before = harness.ctx.agents.list().length
|
|
await harness.closeClientTransport() // teardown → closed = true
|
|
await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] }).catch(() => {})
|
|
await new Promise(r => setTimeout(r, 10))
|
|
expect(harness.ctx.agents.list().length).toBe(before)
|
|
await harness.dispose()
|
|
})
|
|
|
|
it('a client disconnect mid-prompt disposes the session (no registered agent left)', async () => {
|
|
// The ACP transport closes (editor quits) while a turn runs. The bridge must
|
|
// settle the in-flight prompt cancelled and DISPOSE the agent (PR D's
|
|
// per-agent AgentHandle teardown) rather than leaving an orphaned running —
|
|
// or even idled-but-still-registered — agent whose updates are swallowed.
|
|
const harness = await makeBridgeHarness({ storageDir, script: ['hang'] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
const agent = harness.ctx.agents.get(sessionId)!
|
|
// Start a prompt that hangs in the model stream. The prompt RPC will never
|
|
// return (its transport is severed), so do not await it.
|
|
void harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }).catch(() => {})
|
|
await new Promise(r => setTimeout(r, 30))
|
|
expect(agent.status).toBe('running')
|
|
|
|
// Sever the transport — the bridge's conn.closed teardown runs and drives the
|
|
// agent's AgentHandle dispose to quiescence on its OWN (before any dispose()).
|
|
await harness.closeClientTransport()
|
|
await agent.whenIdle()
|
|
// The agent's loop has stopped: status `disposed`.
|
|
expect(agent.status).toBe('disposed')
|
|
|
|
// Await the bridge teardown to completion WITHOUT tearing down the root
|
|
// agents/sessions services (so we can still query them). acpFiber.dispose()
|
|
// invokes the SAME memoized quiesce() the disconnect started and awaits its
|
|
// promise — which resolves only after every rec.dispose() (loop exit +
|
|
// session removal) has finished, closing the whenIdle()/owned.dispose()
|
|
// microtask race. The AgentHandle dispose has run: the agent is unregistered
|
|
// and its session removed from the store, not merely idled (the old
|
|
// behavior). The services live on the root ctx, so they survive this.
|
|
await harness.acpFiber.dispose()
|
|
expect(harness.ctx.agents.get(sessionId)).toBeUndefined()
|
|
expect(harness.ctx.sessions.get(sessionId)).toBeUndefined()
|
|
await harness.dispose()
|
|
})
|
|
|
|
it('a client disconnect racing fiber dispose both reach quiescence (shared teardown)', async () => {
|
|
// conn.closed teardown and ctx.fiber.dispose() can fire near-simultaneously.
|
|
// They must share one teardown promise: dispose() must NOT return before the
|
|
// disconnect teardown's whenIdle() has settled (a `record === undefined`-only
|
|
// guard would let the second caller return early mid-teardown).
|
|
const harness = await makeBridgeHarness({ storageDir, script: ['hang'] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
const agent = harness.ctx.agents.get(sessionId)!
|
|
void harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }).catch(() => {})
|
|
await new Promise(r => setTimeout(r, 30))
|
|
expect(agent.status).toBe('running')
|
|
|
|
// Fire both teardown paths without awaiting the first, then await both.
|
|
const close = harness.closeClientTransport()
|
|
const dispose = harness.ctx.fiber.dispose()
|
|
await Promise.all([close, dispose])
|
|
// After BOTH settle, the agent has fully drained (not still running).
|
|
expect(agent.status).not.toBe('running')
|
|
})
|
|
|
|
it('after dispose, session/update listeners are gone (no further updates emitted)', async () => {
|
|
const harness = await makeBridgeHarness({ storageDir, script: [] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
const session = harness.ctx.agents.get(sessionId)!.session
|
|
|
|
await harness.ctx.fiber.dispose()
|
|
const before = harness.updates.length
|
|
// Append an event directly to the (now-detached) session: the bridge's
|
|
// session/event listener should have been disposed, so no update fires.
|
|
session.append('turn/start', { turn: 99, trigger: { kind: 'message', source: { kind: 'user' } } })
|
|
await new Promise(r => setTimeout(r, 10))
|
|
expect(harness.updates.length).toBe(before)
|
|
})
|
|
|
|
it('the final turn closing events are persisted across an AgentHandle dispose (durability)', async () => {
|
|
// The teardown-ORDER guarantee: a per-agent dispose must stop the loop,
|
|
// AWAIT its exit (so the loop's final `turn/end` + `session/flush` fire
|
|
// through the still-attached `session.onAppend` → `session/event`), and only
|
|
// THEN detach onAppend + remove the session. If the order were inverted
|
|
// (detach first), the closing events would never reach persistence. Drive a
|
|
// CLEAN turn to completion, dispose JUST the bridge, then re-load the
|
|
// persisted log from disk and assert the closing turn/end is on disk — the
|
|
// world, not the agent's self-report.
|
|
const harness = await makeBridgeHarness({ storageDir, script: [textResponse('done')] })
|
|
await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
|
|
const { sessionId } = await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })
|
|
await harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })
|
|
const liveEvents = harness.ctx.agents.get(sessionId)!.session.events.length
|
|
expect(liveEvents).toBeGreaterThan(0)
|
|
|
|
// Tear down JUST the bridge (the AgentHandle dispose runs to quiescence).
|
|
await harness.acpFiber.dispose()
|
|
expect(harness.ctx.agents.get(sessionId)).toBeUndefined()
|
|
|
|
// Re-load the session from disk: every live event (incl. the closing
|
|
// turn/end) was flushed before the session was detached.
|
|
const reloaded = await harness.ctx.sessionPersistence.load(SessionId(sessionId))
|
|
expect(reloaded.events.length).toBe(liveEvents)
|
|
const last = reloaded.events.at(-1)!
|
|
expect(last.type).toBe('turn/end')
|
|
await harness.dispose()
|
|
})
|
|
|
|
it('per-session AgentHandle dispose leaves sibling agents untouched', async () => {
|
|
// The factory returns a per-agent AgentHandle whose dispose() tears down
|
|
// EXACTLY that agent + its session — RFC 011 isolation. Create two agents
|
|
// directly through the registry factory (the same path the ACP bridge uses),
|
|
// dispose one handle, and assert the other survives, registered and
|
|
// queryable, with its session still in the store.
|
|
const harness = await makeBridgeHarness({ storageDir, script: [] })
|
|
const handleA = harness.ctx.agents.create({
|
|
agentId: 'sib-a', sessionId: 'sib-a', agentOptions: { model: 'mock' },
|
|
})
|
|
const handleB = harness.ctx.agents.create({
|
|
agentId: 'sib-b', sessionId: 'sib-b', agentOptions: { model: 'mock' },
|
|
})
|
|
expect(harness.ctx.agents.get('sib-a')).toBe(handleA.agent)
|
|
expect(harness.ctx.agents.get('sib-b')).toBe(handleB.agent)
|
|
|
|
await handleA.dispose()
|
|
// A is gone — unregistered AND its session removed from the store.
|
|
expect(harness.ctx.agents.get('sib-a')).toBeUndefined()
|
|
expect(harness.ctx.sessions.get('sib-a')).toBeUndefined()
|
|
expect(handleA.agent.status).toBe('disposed')
|
|
// B is wholly unaffected.
|
|
expect(harness.ctx.agents.get('sib-b')).toBe(handleB.agent)
|
|
expect(harness.ctx.sessions.get('sib-b')).toBeDefined()
|
|
expect(handleB.agent.status).not.toBe('disposed')
|
|
await harness.dispose()
|
|
})
|
|
})
|