From 45e3997efcecf48c57e171a7d77401b1c914ea32 Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Tue, 14 Jul 2026 08:41:31 +0800 Subject: [PATCH] fix: correlate subagent completion by parent scope --- docs/event-producer-consumer.md | 4 +- packages/ui/jsonrpc/README.md | 2 +- packages/ui/jsonrpc/src/server.ts | 101 ++++++++--------------- packages/ui/jsonrpc/tests/server.spec.ts | 13 ++- 4 files changed, 47 insertions(+), 73 deletions(-) diff --git a/docs/event-producer-consumer.md b/docs/event-producer-consumer.md index 04aff464a9..bc1c0d38d0 100644 --- a/docs/event-producer-consumer.md +++ b/docs/event-producer-consumer.md @@ -7,8 +7,8 @@ This matrix shows which packages dispatch each harness-owned event and which pac | Event | Mode | Declared in | Dispatchers | Listeners | | --- | --- | --- | --- | --- | -| `agent/created` | `emit` | [`packages/core/agent/src/types.ts:304`](../packages/core/agent/src/types.ts) | [`agent`](../packages/core/agent) (`events.dispatch`) | [`jsonrpc`](../packages/ui/jsonrpc), [`stdio-agent`](../packages/ui/stdio-agent) | -| `agent/disposed` | `emit` | [`packages/core/agent/src/types.ts:319`](../packages/core/agent/src/types.ts) | [`agent`](../packages/core/agent) (`events.dispatch`) | [`jsonrpc`](../packages/ui/jsonrpc), [`stdio-agent`](../packages/ui/stdio-agent) | +| `agent/created` | `emit` | [`packages/core/agent/src/types.ts:304`](../packages/core/agent/src/types.ts) | [`agent`](../packages/core/agent) (`events.dispatch`) | [`stdio-agent`](../packages/ui/stdio-agent) | +| `agent/disposed` | `emit` | [`packages/core/agent/src/types.ts:319`](../packages/core/agent/src/types.ts) | [`agent`](../packages/core/agent) (`events.dispatch`) | [`stdio-agent`](../packages/ui/stdio-agent) | | `agent/error` | `emit` | [`packages/core/agent/src/types.ts:593`](../packages/core/agent/src/types.ts) | [`agent-loop`](../packages/core/agent-loop) (`emit`) | - | | `agent/pre-step` | `serial` | [`packages/core/agent/src/types.ts:426`](../packages/core/agent/src/types.ts) | [`agent-loop`](../packages/core/agent-loop) (`serial`) | [`compact-basic`](../packages/compact/compact-basic), [`user-approval`](../packages/ui/user-approval) | | `agent/prompt-submit` | `waterfall` | [`packages/core/agent/src/types.ts:444`](../packages/core/agent/src/types.ts) | [`agent-loop`](../packages/core/agent-loop) (`waterfall`) | [`acp`](../packages/ui/acp), [`hooks-claude`](../packages/hooks/hooks-claude), [`hooks-codex`](../packages/hooks/hooks-codex), [`repeat-tool-guard`](../packages/guard/repeat-tool-guard) | diff --git a/packages/ui/jsonrpc/README.md b/packages/ui/jsonrpc/README.md index 5ff7eb7f8e..ffbd7a1288 100644 --- a/packages/ui/jsonrpc/README.md +++ b/packages/ui/jsonrpc/README.md @@ -4,7 +4,7 @@ The **SDK server plugin** (`jsonrpc`): mounting it serves a stdio JSON-RPC serve ## Wiring -`inject: ['agents']` — the server creates one agent per SDK `sessionId` (get-or-create on `session/prompt`). A local subagent's shared agent/session id supplies `subagent.finished.childSessionId` directly; the server caches runtime-local identity plus optional parent lineage for the agent lifetime and counts pending runs per provider/id, because a continuation may reuse one child and the child may be disposed before a later `subagent/end`. Settlement order is not assumed: if concurrent ID reuse makes parent lineage ambiguous, the completion remains local but omits the optional `parentSessionId` rather than attributing the wrong parent. Runs from remote providers are not reported because they create no local agent. The LLM seam is read opportunistically via `ctx.get('llm')` (not injected): when `initialize.model` has no registered adapter, the plugin mounts `dsh-llm-deepseek` for it (credentials from `$DEEPSEEK_API_KEY` / `$DEEPSEEK_BASE_URL`) — a config-registered adapter for the model wins. Everything else — persistence, the tool stacks, the adapter set — comes from the surrounding `cordis.yml`. +`inject: ['agents']` — the server creates one agent per SDK `sessionId` (get-or-create on `session/prompt`). A local subagent's shared agent/session id supplies `subagent.finished.childSessionId` directly; the server counts local starts by provider/id and the exact delegating-parent carrier, because a continuation may reuse one child and the child may be disposed before a later `subagent/end`. The paired event carrier preserves parent correlation even when reused ids settle out of order. Runs from remote providers are not reported because they create no local agent. The LLM seam is read opportunistically via `ctx.get('llm')` (not injected): when `initialize.model` has no registered adapter, the plugin mounts `dsh-llm-deepseek` for it (credentials from `$DEEPSEEK_API_KEY` / `$DEEPSEEK_BASE_URL`) — a config-registered adapter for the model wins. Everything else — persistence, the tool stacks, the adapter set — comes from the surrounding `cordis.yml`. ## Config diff --git a/packages/ui/jsonrpc/src/server.ts b/packages/ui/jsonrpc/src/server.ts index ce0ac9c75a..352563f7f2 100644 --- a/packages/ui/jsonrpc/src/server.ts +++ b/packages/ui/jsonrpc/src/server.ts @@ -15,8 +15,10 @@ import type { Context } from 'cordis' import { resolve } from 'node:path' import type { ContentBlock } from '@deepseek-ai/dsh-llm' -import type { AgentHandle } from '@deepseek-ai/dsh-agent' +import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent' +import { carrierKeyOf, type Scoped } from '@deepseek-ai/dsh-scope' import { SessionId, type TurnEndReason } from '@deepseek-ai/dsh-session' +import type SubagentService from '@deepseek-ai/dsh-subagent' import type { SubagentRunEndInfo, SubagentRunInfo } from '@deepseek-ai/dsh-subagent' import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek' import type { JsonRpcTransportPeer } from './transport.ts' @@ -58,21 +60,14 @@ interface SessionRecord { activePrompt: boolean } -/** Runtime-local agent identity plus optional durable fork lineage. */ -interface LocalAgentRecord { - parentSessionId?: SessionId -} - -/** Pending local runs that share one provider/id correlation key. */ -interface PendingLocalRuns { - count: number - parentSessionId?: SessionId - parentAmbiguous: boolean +/** Recover the delegating parent carried by every service-owned subagent lifecycle event. */ +function subagentParentOf(carrier: Scoped): Agent { + return carrierKeyOf(carrier) as Agent } /** * The SDK server over a booted harness context. Constructing it subscribes to - * session, agent, and subagent lifecycle events, forwarding durable session + * session and subagent lifecycle events, forwarding durable session * events and SDK-facing completion notifications while retaining local-run * identity across child disposal. The subscriptions live until * {@link shutdown}. One instance serves one transport peer for the process @@ -84,8 +79,7 @@ export class HarnessSdkServer { private llmFiber: { dispose(): Promise } | undefined private readonly sessions = new Map() private readonly sessionCreations = new Map>() - private readonly localAgents = new Map() - private readonly localRuns = new Map>() + private readonly localRuns = new Map>>() private readonly disposers: (() => void)[] = [] private shutdownTask: Promise> | undefined private shuttingDown = false @@ -109,64 +103,40 @@ export class HarnessSdkServer { childSessionId: String(session.id), }) })) - // Cache runtime-local identity and optional lineage for each agent lifetime. - // Parent lineage is not required by the provider contract, so an empty - // record remains a load-bearing locality marker. - this.disposers.push(ctx.on('agent/created', (agent) => { - const parentSessionId = agent.session.header.parentSession - this.localAgents.set(agent.id, parentSessionId === undefined ? {} : { parentSessionId }) + // In-process providers publish the child before start. Count those starts by + // the exact delegating-parent carrier so later completions remain local after + // child disposal and reused ids need no settlement-order assumption. + const localRuns = this.localRuns + this.disposers.push(ctx.on('subagent/start', function (this: Scoped, info: SubagentRunInfo) { + if (ctx.agents.get(info.id) === undefined) return + const parent = subagentParentOf(this) + const providerRuns = localRuns.get(info.provider) ?? new Map>() + const parentRuns = providerRuns.get(info.id) ?? new Map() + parentRuns.set(parent, (parentRuns.get(parent) ?? 0) + 1) + providerRuns.set(info.id, parentRuns) + localRuns.set(info.provider, providerRuns) })) - this.disposers.push(ctx.on('agent/disposed', (agent) => { - this.localAgents.delete(agent.id) - })) - // Snapshot locality per provider/id run key. A provider may settle one run, - // continue the same live child in another run, and dispose that child before - // the later result settles. Counts preserve every completion without - // assuming settlement order. If id reuse produces disagreeing lineage, the - // optional parent is omitted until that pending group drains rather than - // attributed to the wrong completion. - this.disposers.push(ctx.on('subagent/start', (info: SubagentRunInfo) => { - const agent = this.ctx.agents.get(info.id) - const cachedLocalAgent = this.localAgents.get(info.id) - const localAgent = cachedLocalAgent ?? (agent === undefined - ? undefined - : agent.session.header.parentSession === undefined - ? {} - : { parentSessionId: agent.session.header.parentSession }) - if (localAgent === undefined) return - const providerRuns = this.localRuns.get(info.provider) ?? new Map() - const pending = providerRuns.get(info.id) - if (pending === undefined) { - providerRuns.set(info.id, localAgent.parentSessionId === undefined - ? { count: 1, parentAmbiguous: false } - : { count: 1, parentSessionId: localAgent.parentSessionId, parentAmbiguous: false }) - } else { - pending.count += 1 - if (pending.parentSessionId !== localAgent.parentSessionId) pending.parentAmbiguous = true - } - this.localRuns.set(info.provider, providerRuns) - })) - this.disposers.push(ctx.on('subagent/end', (info: SubagentRunEndInfo) => { - const agent = this.ctx.agents.get(info.id) - const providerRuns = this.localRuns.get(info.provider) - const pending = providerRuns?.get(info.id) - if (pending !== undefined) { - pending.count -= 1 - if (pending.count === 0) providerRuns?.delete(info.id) - if (providerRuns?.size === 0) this.localRuns.delete(info.provider) + this.disposers.push(ctx.on('subagent/end', function (this: Scoped, info: SubagentRunEndInfo) { + const agent = ctx.agents.get(info.id) + const parent = subagentParentOf(this) + const providerRuns = localRuns.get(info.provider) + const parentRuns = providerRuns?.get(info.id) + const pendingCount = parentRuns?.get(parent) + if (pendingCount !== undefined) { + if (pendingCount === 1) parentRuns?.delete(parent) + else parentRuns?.set(parent, pendingCount - 1) + if (parentRuns?.size === 0) providerRuns?.delete(info.id) + if (providerRuns?.size === 0) localRuns.delete(info.provider) } // This protocol reports LOCAL child sessions. A lineage-bearing child // has the session/created-driven start notification above; a parentless // local provider still gets its terminal notification. A remote provider - // has neither a cached creation nor a live local agent and is ignored. - if (pending === undefined && agent === undefined) return - const parentSessionId = pending === undefined - ? agent?.session.header.parentSession - : pending.parentAmbiguous ? undefined : pending.parentSessionId - this.transport.notify('subagent.finished', { + // has neither a pending local start nor a live local agent and is ignored. + if (pendingCount === undefined && agent === undefined) return + transport.notify('subagent.finished', { provider: info.provider, agentId: String(info.id), - ...(parentSessionId === undefined ? {} : { parentSessionId: String(parentSessionId) }), + parentSessionId: String(parent.session.id), childSessionId: String(info.id), status: info.stopReason === 'completed' ? 'ok' : 'error', stopReason: info.stopReason, @@ -240,7 +210,6 @@ export class HarnessSdkServer { this.sessionCreations.clear() const records = [...this.sessions.values()] this.sessions.clear() - this.localAgents.clear() this.localRuns.clear() const failures: unknown[] = [] while (this.disposers.length > 0) { diff --git a/packages/ui/jsonrpc/tests/server.spec.ts b/packages/ui/jsonrpc/tests/server.spec.ts index 48b8c4b1be..427a190c36 100644 --- a/packages/ui/jsonrpc/tests/server.spec.ts +++ b/packages/ui/jsonrpc/tests/server.spec.ts @@ -316,6 +316,7 @@ describe('HarnessSdkServer', () => { params: { provider: 'spawn', agentId: 'parentless-child-session', + parentSessionId: 'main', childSessionId: 'parentless-child-session', status: 'error', stopReason: 'error', @@ -373,7 +374,7 @@ describe('HarnessSdkServer', () => { } }) - it('omits ambiguous lineage when one local id is reused and runs settle out of order', async () => { + it('correlates reused local ids by parent scope when runs settle out of order', async () => { const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-reuse-')) const ctx = await makeHarness(storageDir) try { @@ -450,8 +451,11 @@ describe('HarnessSdkServer', () => { [{ type: 'text', text: 'new lifetime' }], [{ type: 'text', text: 'old lifetime' }], ]) - expect(finished[0]?.params?.parentSessionId).toBe('old-parent') - expect(finished.slice(1).every(notification => !Object.hasOwn(notification.params ?? {}, 'parentSessionId'))).toBe(true) + expect(finished.map(notification => notification.params?.parentSessionId)).toEqual([ + 'old-parent', + 'new-parent', + 'old-parent', + ]) await firstRun.dispose() await sameLifetimeRun.dispose() @@ -551,6 +555,7 @@ describe('HarnessSdkServer', () => { params: { provider: 'fork', agentId: 'failed-child-session', + parentSessionId: 'fallback-parent', childSessionId: 'failed-child-session', status: 'error', stopReason: 'error', @@ -753,6 +758,6 @@ describe('HarnessSdkServer', () => { const server = new HarnessSdkServer(ctx, new FakeTransport()) await expect(server.shutdown()).rejects.toBe(listenerFailure) - expect(on).toHaveBeenCalledTimes(6) + expect(on).toHaveBeenCalledTimes(4) }) })