From e784e4dce5f5178b5e22a3d3376599144d8bea1b Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Tue, 14 Jul 2026 09:12:32 +0800 Subject: [PATCH 1/3] fix: await ACP client callbacks during shutdown --- packages/support/acp-snapshot/src/launcher.ts | 59 ++++++++++++------- .../acp-snapshot/tests/harness.spec.ts | 36 +++++++++++ 2 files changed, 75 insertions(+), 20 deletions(-) diff --git a/packages/support/acp-snapshot/src/launcher.ts b/packages/support/acp-snapshot/src/launcher.ts index 815753ae5f..30702a5888 100644 --- a/packages/support/acp-snapshot/src/launcher.ts +++ b/packages/support/acp-snapshot/src/launcher.ts @@ -65,7 +65,7 @@ export interface LaunchedAcpTestAgent { stderr(): string /** Resolve when a future session update matches the predicate. */ waitForUpdate(match: (update: SessionNotification['update']) => boolean): Promise - /** Gracefully close stdin, or send a signal, then wait for process exit, inherited stdio closure, and ACP parser drain. */ + /** Gracefully close stdin, or send a signal, then wait for process exit, inherited stdio closure, ACP parsing, and client callbacks. */ close(signal?: NodeJS.Signals): Promise } @@ -137,29 +137,41 @@ export function launchAcpTestAgent(options: AcpTestLaunchOptions): LaunchedAcpTe Writable.toWeb(child.stdin) as WritableStream, Readable.toWeb(passthrough) as ReadableStream, ) + const inFlightClientCallbacks = new Set>() + const trackClientCallback = (callback: () => T | PromiseLike): Promise => { + const pending = Promise.resolve().then(callback) + inFlightClientCallbacks.add(pending) + void pending.then( + () => { inFlightClientCallbacks.delete(pending) }, + () => { inFlightClientCallbacks.delete(pending) }, + ) + return pending + } + const requestPermission = options.requestPermission + ?? (() => Promise.resolve({ outcome: { outcome: 'cancelled' as const } })) const makeClient = (_agent: AcpAgent): Client => ({ sessionUpdate(params: SessionNotification): Promise { - updates.push(params.update) - for (let index = updateWaiters.length - 1; index >= 0; index--) { - const waiter = updateWaiters[index] - /* v8 ignore next 1 -- index is bounded by the array length */ - if (waiter === undefined) continue - let matches: boolean - try { - matches = waiter.match(params.update) - } catch (error: unknown) { + return trackClientCallback(() => { + updates.push(params.update) + for (let index = updateWaiters.length - 1; index >= 0; index--) { + const waiter = updateWaiters[index] + /* v8 ignore next 1 -- index is bounded by the array length */ + if (waiter === undefined) continue + let matches: boolean + try { + matches = waiter.match(params.update) + } catch (error: unknown) { + updateWaiters.splice(index, 1) + waiter.reject(error) + continue + } + if (!matches) continue updateWaiters.splice(index, 1) - waiter.reject(error) - continue + waiter.resolve(params.update) } - if (!matches) continue - updateWaiters.splice(index, 1) - waiter.resolve(params.update) - } - return Promise.resolve() + }) }, - requestPermission: options.requestPermission - ?? (() => Promise.resolve({ outcome: { outcome: 'cancelled' } })), + requestPermission: params => trackClientCallback(() => requestPermission(params)), }) const client = new ClientSideConnection(makeClient, stream) // `exit` only reports the parent process's status. Descendants may retain @@ -168,7 +180,14 @@ export function launchAcpTestAgent(options: AcpTestLaunchOptions): LaunchedAcpTe // `closed` follows parser exhaustion. Capture both eagerly so a caller that // invokes close after process exit still joins the complete drain boundary. const stdioClosed = new Promise(resolve => child.once('close', () => { resolve() })) - const drained = Promise.all([stdioClosed, client.closed]).then(() => undefined) + const drained = Promise.all([stdioClosed, client.closed]).then(async () => { + // The ACP SDK's readable loop dispatches client callbacks without awaiting + // them. Once `closed` settles no new callbacks can start, but callbacks + // already in flight still belong to this launch's teardown boundary. + while (inFlightClientCallbacks.size > 0) { + await Promise.allSettled([...inFlightClientCallbacks]) + } + }) // A caller may await a pending update without calling close(). Make natural // stream exhaustion terminal for those waiters too, but only after the // parser has dispatched every buffered frame. diff --git a/packages/support/acp-snapshot/tests/harness.spec.ts b/packages/support/acp-snapshot/tests/harness.spec.ts index 051db6a0ba..902f48d2cf 100644 --- a/packages/support/acp-snapshot/tests/harness.spec.ts +++ b/packages/support/acp-snapshot/tests/harness.spec.ts @@ -1,4 +1,5 @@ import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { once } from 'node:events' import { tmpdir } from 'node:os' import { delimiter, join } from 'node:path' import { fileURLToPath } from 'node:url' @@ -112,6 +113,41 @@ describe('runScenario', () => { expect(launched.stderr()).toContain('late inherited stderr') }) + it('waits for in-flight client callbacks after the ACP stream closes', { timeout: 20_000 }, async () => { + const { dir, fixtureFile } = await scenario({ permissionProbe: true }) + let releasePermission: (() => void) | undefined + const permissionReleased = new Promise((resolve) => { releasePermission = resolve }) + let markPermissionStarted: (() => void) | undefined + const permissionStarted = new Promise((resolve) => { markPermissionStarted = resolve }) + let permissionFinished = false + const launched = launchAcpTestAgent({ + agent: AGENT, + cwd: dir, + env: { DSH_SNAPSHOT_FILE: fixtureFile }, + async requestPermission() { + markPermissionStarted?.() + await permissionReleased + permissionFinished = true + return { outcome: { outcome: 'cancelled' } } + }, + }) + await launched.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} }) + const { sessionId } = await launched.client.newSession({ cwd: dir, mcpServers: [] }) + void launched.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }).catch(() => undefined) + await permissionStarted + + const childClosed = once(launched.child, 'close') + let closeSettled = false + const closing = launched.close('SIGKILL').then(() => { closeSettled = true }) + await childClosed + await launched.client.closed + expect(closeSettled).toBe(false) + + releasePermission?.() + await closing + expect(permissionFinished).toBe(true) + }) + it('drives a full turn: initialize (terminal caps), session, prompt, permission stub, harvest', { timeout: 20_000 }, async () => { const { fixtureFile } = await scenario({ permissionProbe: true, From d4c96deac3802164c430335309f936f3e5ab4f1d Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Tue, 14 Jul 2026 09:19:02 +0800 Subject: [PATCH 2/3] refactor: share ACP callback cleanup --- packages/support/acp-snapshot/src/launcher.ts | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/packages/support/acp-snapshot/src/launcher.ts b/packages/support/acp-snapshot/src/launcher.ts index 30702a5888..f1271986c0 100644 --- a/packages/support/acp-snapshot/src/launcher.ts +++ b/packages/support/acp-snapshot/src/launcher.ts @@ -141,10 +141,8 @@ export function launchAcpTestAgent(options: AcpTestLaunchOptions): LaunchedAcpTe const trackClientCallback = (callback: () => T | PromiseLike): Promise => { const pending = Promise.resolve().then(callback) inFlightClientCallbacks.add(pending) - void pending.then( - () => { inFlightClientCallbacks.delete(pending) }, - () => { inFlightClientCallbacks.delete(pending) }, - ) + const untrack = (): void => { inFlightClientCallbacks.delete(pending) } + void pending.then(untrack, untrack) return pending } const requestPermission = options.requestPermission From 2dea4afac0d49976700139a766d40734e3204f4e Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Tue, 14 Jul 2026 09:23:02 +0800 Subject: [PATCH 3/3] fix: bind stdio to its configured agent --- packages/ui/stdio-agent/README.md | 2 +- packages/ui/stdio-agent/src/index.ts | 5 ++- packages/ui/stdio-agent/src/stdio-chat.ts | 30 +++++++++++----- .../ui/stdio-agent/tests/stdio-chat.spec.ts | 35 ++++++++++++++----- 4 files changed, 53 insertions(+), 19 deletions(-) diff --git a/packages/ui/stdio-agent/README.md b/packages/ui/stdio-agent/README.md index 82da64cd81..d7d30ee49d 100644 --- a/packages/ui/stdio-agent/README.md +++ b/packages/ui/stdio-agent/README.md @@ -32,7 +32,7 @@ The leaf `cordis.yml` supplies only the **swappable backends** — an LLM adapte | `welcome` | `ready.` | the stdin-chat banner | | `resumeSessionId` | — | resume a persisted session id instead of starting fresh (sourced from an env var in the leaf) | -Fresh stdio sessions use the process launch directory as `session.header.cwd` and mint one combined `main-session-` agent/session id, so durable restarts cannot collide. The UI's `main` text is a display label, not a second routing id. Resumed sessions register under the exact `resumeSessionId` and keep the cwd stored in the persisted session header. +Fresh stdio sessions use the process launch directory as `session.header.cwd` and mint one combined `main-session-` agent/session id, so durable restarts cannot collide. The UI's `main` text is a display label, not a second routing id; the UI binds to that fresh-id namespace, or to the exact `resumeSessionId` for a resumed run, and never selects unrelated registry roots. Resumed sessions keep the cwd stored in the persisted session header. ## The bin diff --git a/packages/ui/stdio-agent/src/index.ts b/packages/ui/stdio-agent/src/index.ts index f01c48a615..2eca0cceb3 100644 --- a/packages/ui/stdio-agent/src/index.ts +++ b/packages/ui/stdio-agent/src/index.ts @@ -124,5 +124,8 @@ export function apply(ctx: Context, config: Config): void { ctx.plugin(SessionPersistenceJsonl, { root: config.persistenceRoot ?? './.sessions' }) ctx.plugin(UserInteractionService) ctx.plugin(toolAskUser) - ctx.plugin(uiStdio, { welcome: config.welcome ?? 'ready.' }) + ctx.plugin(uiStdio, { + welcome: config.welcome ?? 'ready.', + ...config.resumeSessionId !== undefined ? { resumeSessionId: config.resumeSessionId } : {}, + }) } diff --git a/packages/ui/stdio-agent/src/stdio-chat.ts b/packages/ui/stdio-agent/src/stdio-chat.ts index a295274786..76347c91e4 100644 --- a/packages/ui/stdio-agent/src/stdio-chat.ts +++ b/packages/ui/stdio-agent/src/stdio-chat.ts @@ -36,10 +36,13 @@ export const inject = ['agents', 'userInteraction'] export interface Config { /** Banner printed once on start, before the first `> ` prompt. */ welcome?: string + /** Exact persisted session id the app configured for resume; absent selects the app's fresh `main-session-*` identity. */ + resumeSessionId?: string } export const Config: z = z.object({ welcome: z.string().default('ready.'), + resumeSessionId: z.string(), }) /** @@ -95,16 +98,25 @@ export function createStdioChat(ctx: Context, config: Config, runtime: StdioRunt const welcome = config.welcome ?? 'ready.' const { input, output, exit } = runtime - // This app owns one configured top-level agent. Hold the live object - // directly: its per-run id is intentionally fresh, while `main` remains only - // the terminal's fixed display label. Runtime creator ownership distinguishes - // that root from its subagents even if a child is registered after an HMR - // replacement. Persisted parentSession lineage is deliberately irrelevant: - // a resumed child session can itself be this process's configured root. - let target: Agent | undefined = ctx.agents.roots()[0] - ctx.on('agent/created', () => { target ??= ctx.agents.roots()[0] }) + // Bind only to this app's configured top-level agent. Fresh runs own the + // `main-session-*` namespace; resumed runs own the exact persisted id. The + // registry's runtime-root relation excludes subagents without confusing it + // with durable parentSession lineage. Keeping the matching candidates also + // covers HMR's publish-new-before-dispose-old ordering without ever falling + // through to an unrelated root owned by another app or test fixture. + const matchesConfiguredIdentity = (agent: Agent): boolean => config.resumeSessionId === undefined + ? agent.id.startsWith('main-session-') + : agent.id === config.resumeSessionId + const configuredRoots = new Set(ctx.agents.roots().filter(matchesConfiguredIdentity)) + let target: Agent | undefined = [...configuredRoots].at(-1) + ctx.on('agent/created', (agent) => { + if (!matchesConfiguredIdentity(agent) || !ctx.agents.roots().includes(agent)) return + configuredRoots.add(agent) + target ??= agent + }) ctx.on('agent/disposed', (agent) => { - if (target === agent) target = ctx.agents.roots().at(-1) + configuredRoots.delete(agent) + if (target === agent) target = [...configuredRoots].at(-1) }) // Transcript rendering off the durable `session/event` feed — the assistant diff --git a/packages/ui/stdio-agent/tests/stdio-chat.spec.ts b/packages/ui/stdio-agent/tests/stdio-chat.spec.ts index 2600fdb24f..668068a72d 100644 --- a/packages/ui/stdio-agent/tests/stdio-chat.spec.ts +++ b/packages/ui/stdio-agent/tests/stdio-chat.spec.ts @@ -74,7 +74,7 @@ function chunkEvent(chunk: StreamChunk): SessionEvent { return { type: 'assistant/chunk', seq: 0, time: 0, data: { turn: 1, step: 0, chunk } } } -const CONFIG: Config = { welcome: 'hi there' } +const CONFIG: Config = { welcome: 'hi there', resumeSessionId: 'main' } async function setup(config: Config = CONFIG, runtimeOver: Partial = {}) { const ctx = new Context() @@ -199,7 +199,9 @@ describe('createStdioChat rendering', () => { }) it('accepts a lineage-bearing configured agent created after the UI installs', async () => { - const { ctx, input } = await setup() + const { ctx, input } = await setup({ welcome: 'hi there', resumeSessionId: 'resumed' }) + const unrelated = makeAgent('unrelated') + ctx.agents.register(unrelated) const resumed = makeAgent('resumed') ;(resumed.session.header as { parentSession?: string }).parentSession = 'persisted-parent' ctx.agents.register(resumed) @@ -207,6 +209,7 @@ describe('createStdioChat rendering', () => { input.feed('continue') await new Promise(resolve => setImmediate(resolve)) + expect(unrelated.sent).toEqual([]) expect(resumed.sent).toEqual([[{ type: 'text', text: 'continue' }]]) }) @@ -235,7 +238,7 @@ describe('createStdioChat rendering', () => { it('keeps the target when a different agent is disposed', async () => { const { ctx, out } = await setup() - const target = makeAgent('target') + const target = makeAgent('main') ctx.agents.register(target) ctx.emit('agent/disposed', makeAgent('other')) ctx.emit('session/event', target.session, { @@ -245,11 +248,11 @@ describe('createStdioChat rendering', () => { }) it('retargets a surviving root when HMR publishes it before disposing the old root', async () => { - const { ctx, input } = await setup() - const oldRoot = makeAgent('old-root') + const { ctx, input } = await setup({ welcome: 'hi there' }) + const oldRoot = makeAgent('main-session-old') const child = makeAgent('child') ;(child.session.header as { parentSession?: string }).parentSession = oldRoot.id - const replacement = makeAgent('replacement') + const replacement = makeAgent('main-session-replacement') const lateChild = makeAgent('late-child') const disposeOld = ctx.agents.register(oldRoot) const disposeChild = ctx.agents.enter(child, oldRoot) @@ -273,6 +276,22 @@ describe('createStdioChat rendering', () => { disposeChild() }) + it('does not retarget stdin to an unrelated root after the configured agent is disposed', async () => { + const { ctx, input } = await setup() + const unrelated = makeAgent('unrelated') + ctx.agents.register(unrelated) + const configured = makeAgent('main') + const disposeConfigured = ctx.agents.register(configured) + const error = vi.spyOn(ctx.logger, 'error').mockImplementation(() => {}) + + disposeConfigured() + input.feed('must not leak') + await new Promise(resolve => setImmediate(resolve)) + + expect(unrelated.sent).toEqual([]) + expect(error).toHaveBeenCalledWith('ui-stdio: main agent is not running') + }) + it('renders tool/call and tool/result session events', async () => { const { ctx, out } = await setup() const session = {} as Session @@ -720,8 +739,8 @@ describe('createStdioChat input', () => { expect(spy).toHaveBeenCalledWith('ui-stdio: main agent is not running') }) - it('drives the app-owned agent without a duplicate id config', async () => { - const { ctx, input } = await setup({ welcome: 'w' }) + it('drives the exact app-configured resumed session', async () => { + const { ctx, input } = await setup({ welcome: 'w', resumeSessionId: 'worker' }) const agent = makeAgent('worker') ctx.agents.register(agent) input.feed('hi')