fix: correlate subagent completion by parent scope

This commit is contained in:
Tianyi Cui
2026-07-14 08:41:31 +08:00
parent d898d1faf0
commit 45e3997efc
4 changed files with 47 additions and 73 deletions

View File

@@ -7,8 +7,8 @@ This matrix shows which packages dispatch each harness-owned event and which pac
| Event | Mode | Declared in | Dispatchers | Listeners | | 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/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`) | [`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`) | [`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/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/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) | | `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) |

View File

@@ -4,7 +4,7 @@ The **SDK server plugin** (`jsonrpc`): mounting it serves a stdio JSON-RPC serve
## Wiring ## 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 ## Config

View File

@@ -15,8 +15,10 @@
import type { Context } from 'cordis' import type { Context } from 'cordis'
import { resolve } from 'node:path' import { resolve } from 'node:path'
import type { ContentBlock } from '@deepseek-ai/dsh-llm' 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 { 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 type { SubagentRunEndInfo, SubagentRunInfo } from '@deepseek-ai/dsh-subagent'
import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek' import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
import type { JsonRpcTransportPeer } from './transport.ts' import type { JsonRpcTransportPeer } from './transport.ts'
@@ -58,21 +60,14 @@ interface SessionRecord {
activePrompt: boolean activePrompt: boolean
} }
/** Runtime-local agent identity plus optional durable fork lineage. */ /** Recover the delegating parent carried by every service-owned subagent lifecycle event. */
interface LocalAgentRecord { function subagentParentOf(carrier: Scoped<SubagentService>): Agent {
parentSessionId?: SessionId return carrierKeyOf(carrier) as Agent
}
/** Pending local runs that share one provider/id correlation key. */
interface PendingLocalRuns {
count: number
parentSessionId?: SessionId
parentAmbiguous: boolean
} }
/** /**
* The SDK server over a booted harness context. Constructing it subscribes to * 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 * events and SDK-facing completion notifications while retaining local-run
* identity across child disposal. The subscriptions live until * identity across child disposal. The subscriptions live until
* {@link shutdown}. One instance serves one transport peer for the process * {@link shutdown}. One instance serves one transport peer for the process
@@ -84,8 +79,7 @@ export class HarnessSdkServer {
private llmFiber: { dispose(): Promise<void> } | undefined private llmFiber: { dispose(): Promise<void> } | undefined
private readonly sessions = new Map<string, SessionRecord>() private readonly sessions = new Map<string, SessionRecord>()
private readonly sessionCreations = new Map<string, Promise<SessionRecord>>() private readonly sessionCreations = new Map<string, Promise<SessionRecord>>()
private readonly localAgents = new Map<SessionId, LocalAgentRecord>() private readonly localRuns = new Map<string, Map<SessionId, Map<Agent, number>>>()
private readonly localRuns = new Map<string, Map<SessionId, PendingLocalRuns>>()
private readonly disposers: (() => void)[] = [] private readonly disposers: (() => void)[] = []
private shutdownTask: Promise<Record<string, never>> | undefined private shutdownTask: Promise<Record<string, never>> | undefined
private shuttingDown = false private shuttingDown = false
@@ -109,64 +103,40 @@ export class HarnessSdkServer {
childSessionId: String(session.id), childSessionId: String(session.id),
}) })
})) }))
// Cache runtime-local identity and optional lineage for each agent lifetime. // In-process providers publish the child before start. Count those starts by
// Parent lineage is not required by the provider contract, so an empty // the exact delegating-parent carrier so later completions remain local after
// record remains a load-bearing locality marker. // child disposal and reused ids need no settlement-order assumption.
this.disposers.push(ctx.on('agent/created', (agent) => { const localRuns = this.localRuns
const parentSessionId = agent.session.header.parentSession this.disposers.push(ctx.on('subagent/start', function (this: Scoped<SubagentService>, info: SubagentRunInfo) {
this.localAgents.set(agent.id, parentSessionId === undefined ? {} : { parentSessionId }) if (ctx.agents.get(info.id) === undefined) return
const parent = subagentParentOf(this)
const providerRuns = localRuns.get(info.provider) ?? new Map<SessionId, Map<Agent, number>>()
const parentRuns = providerRuns.get(info.id) ?? new Map<Agent, number>()
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.disposers.push(ctx.on('subagent/end', function (this: Scoped<SubagentService>, info: SubagentRunEndInfo) {
this.localAgents.delete(agent.id) const agent = ctx.agents.get(info.id)
})) const parent = subagentParentOf(this)
// Snapshot locality per provider/id run key. A provider may settle one run, const providerRuns = localRuns.get(info.provider)
// continue the same live child in another run, and dispose that child before const parentRuns = providerRuns?.get(info.id)
// the later result settles. Counts preserve every completion without const pendingCount = parentRuns?.get(parent)
// assuming settlement order. If id reuse produces disagreeing lineage, the if (pendingCount !== undefined) {
// optional parent is omitted until that pending group drains rather than if (pendingCount === 1) parentRuns?.delete(parent)
// attributed to the wrong completion. else parentRuns?.set(parent, pendingCount - 1)
this.disposers.push(ctx.on('subagent/start', (info: SubagentRunInfo) => { if (parentRuns?.size === 0) providerRuns?.delete(info.id)
const agent = this.ctx.agents.get(info.id) if (providerRuns?.size === 0) localRuns.delete(info.provider)
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<SessionId, PendingLocalRuns>()
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 protocol reports LOCAL child sessions. A lineage-bearing child // This protocol reports LOCAL child sessions. A lineage-bearing child
// has the session/created-driven start notification above; a parentless // has the session/created-driven start notification above; a parentless
// local provider still gets its terminal notification. A remote provider // local provider still gets its terminal notification. A remote provider
// has neither a cached creation nor a live local agent and is ignored. // has neither a pending local start nor a live local agent and is ignored.
if (pending === undefined && agent === undefined) return if (pendingCount === undefined && agent === undefined) return
const parentSessionId = pending === undefined transport.notify('subagent.finished', {
? agent?.session.header.parentSession
: pending.parentAmbiguous ? undefined : pending.parentSessionId
this.transport.notify('subagent.finished', {
provider: info.provider, provider: info.provider,
agentId: String(info.id), agentId: String(info.id),
...(parentSessionId === undefined ? {} : { parentSessionId: String(parentSessionId) }), parentSessionId: String(parent.session.id),
childSessionId: String(info.id), childSessionId: String(info.id),
status: info.stopReason === 'completed' ? 'ok' : 'error', status: info.stopReason === 'completed' ? 'ok' : 'error',
stopReason: info.stopReason, stopReason: info.stopReason,
@@ -240,7 +210,6 @@ export class HarnessSdkServer {
this.sessionCreations.clear() this.sessionCreations.clear()
const records = [...this.sessions.values()] const records = [...this.sessions.values()]
this.sessions.clear() this.sessions.clear()
this.localAgents.clear()
this.localRuns.clear() this.localRuns.clear()
const failures: unknown[] = [] const failures: unknown[] = []
while (this.disposers.length > 0) { while (this.disposers.length > 0) {

View File

@@ -316,6 +316,7 @@ describe('HarnessSdkServer', () => {
params: { params: {
provider: 'spawn', provider: 'spawn',
agentId: 'parentless-child-session', agentId: 'parentless-child-session',
parentSessionId: 'main',
childSessionId: 'parentless-child-session', childSessionId: 'parentless-child-session',
status: 'error', status: 'error',
stopReason: '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 storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-reuse-'))
const ctx = await makeHarness(storageDir) const ctx = await makeHarness(storageDir)
try { try {
@@ -450,8 +451,11 @@ describe('HarnessSdkServer', () => {
[{ type: 'text', text: 'new lifetime' }], [{ type: 'text', text: 'new lifetime' }],
[{ type: 'text', text: 'old lifetime' }], [{ type: 'text', text: 'old lifetime' }],
]) ])
expect(finished[0]?.params?.parentSessionId).toBe('old-parent') expect(finished.map(notification => notification.params?.parentSessionId)).toEqual([
expect(finished.slice(1).every(notification => !Object.hasOwn(notification.params ?? {}, 'parentSessionId'))).toBe(true) 'old-parent',
'new-parent',
'old-parent',
])
await firstRun.dispose() await firstRun.dispose()
await sameLifetimeRun.dispose() await sameLifetimeRun.dispose()
@@ -551,6 +555,7 @@ describe('HarnessSdkServer', () => {
params: { params: {
provider: 'fork', provider: 'fork',
agentId: 'failed-child-session', agentId: 'failed-child-session',
parentSessionId: 'fallback-parent',
childSessionId: 'failed-child-session', childSessionId: 'failed-child-session',
status: 'error', status: 'error',
stopReason: 'error', stopReason: 'error',
@@ -753,6 +758,6 @@ describe('HarnessSdkServer', () => {
const server = new HarnessSdkServer(ctx, new FakeTransport()) const server = new HarnessSdkServer(ctx, new FakeTransport())
await expect(server.shutdown()).rejects.toBe(listenerFailure) await expect(server.shutdown()).rejects.toBe(listenerFailure)
expect(on).toHaveBeenCalledTimes(6) expect(on).toHaveBeenCalledTimes(4)
}) })
}) })