Merge remote-tracking branch 'origin/worktree/context-source-cards' into worktree/context-forms-remaining
# Conflicts: # apps/web/tests/snapshots/queue-actions/layout.expected.md # docs/cordis-catalog/events.md # docs/cordis-catalog/services.md # docs/core-data-structures/core.i18n.yaml # docs/core-data-structures/goal.i18n.yaml # docs/core-data-structures/goal.md # docs/core-data-structures/goal.zh.md # docs/event-producer-consumer.md # examples/acp-agent/tests/goal-snapshots/goal-session/session.expected.jsonl # examples/acp-agent/tests/goal-snapshots/goal-wrapup/session.expected.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.1.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.2.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.jsonl # examples/acp-agent/tests/snapshots/bash-spill/session.jsonl # examples/acp-agent/tests/snapshots/bash-tool-turn/session.jsonl # examples/acp-agent/tests/snapshots/both-mode-turn/session.jsonl # examples/acp-agent/tests/snapshots/cancel-tool-calls/session.jsonl # examples/acp-agent/tests/snapshots/cancel/session.jsonl # examples/acp-agent/tests/snapshots/code-mode-turn/session.jsonl # examples/acp-agent/tests/snapshots/code-mode-workspace-context/session.jsonl # examples/acp-agent/tests/snapshots/cordis-inspect-jsdoc/session.jsonl # examples/acp-agent/tests/snapshots/empty-response-retry/session.jsonl # examples/acp-agent/tests/snapshots/error-finish/session.jsonl # examples/acp-agent/tests/snapshots/escalation-approved/session.jsonl # examples/acp-agent/tests/snapshots/escalation-rejected/session.jsonl # examples/acp-agent/tests/snapshots/fs-edit/session.jsonl # examples/acp-agent/tests/snapshots/fs-escalation-approved/session.jsonl # examples/acp-agent/tests/snapshots/fs-policy-reject/session.jsonl # examples/acp-agent/tests/snapshots/fs-read-window/session.jsonl # examples/acp-agent/tests/snapshots/fs-read/session.jsonl # examples/acp-agent/tests/snapshots/fs-write-overwrite/session.jsonl # examples/acp-agent/tests/snapshots/fs-write/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-invalid-matcher/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-posttool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-posttool-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-pretool-ask/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-pretool-deny/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-promptsubmit-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-stop-continue/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-invalid-matcher/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-posttool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-posttool-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-pretool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-promptsubmit-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-stop-continue/session.jsonl # examples/acp-agent/tests/snapshots/lsp-definition/session.jsonl # examples/acp-agent/tests/snapshots/missing-sandbox-runner/session.jsonl # examples/acp-agent/tests/snapshots/multi-turn/session.jsonl # examples/acp-agent/tests/snapshots/packed-chunks/session.jsonl # examples/acp-agent/tests/snapshots/parallel-tool-calls/session.jsonl # examples/acp-agent/tests/snapshots/partial-landlock-child-failure/session.jsonl # examples/acp-agent/tests/snapshots/pty-tools/session.jsonl # examples/acp-agent/tests/snapshots/repeat-tool-guard/session.jsonl # examples/acp-agent/tests/snapshots/session-query-spill/session.jsonl # examples/acp-agent/tests/snapshots/session-sandbox-root/session.jsonl # examples/acp-agent/tests/snapshots/session-title-after-turn/session.jsonl # examples/acp-agent/tests/snapshots/skill-load/session.jsonl # examples/acp-agent/tests/snapshots/subagent-continuable/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-continuable/session.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.jsonl # examples/acp-agent/tests/snapshots/subagent-fork/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-fork/session.jsonl # examples/acp-agent/tests/snapshots/subagent-list-agents/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-list-agents/session.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.jsonl # examples/acp-agent/tests/snapshots/subagent-published-run-failure/session.jsonl # examples/acp-agent/tests/snapshots/subagent-report/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-report/session.jsonl # examples/acp-agent/tests/snapshots/subagent-spawn/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-spawn/session.jsonl # examples/acp-agent/tests/snapshots/text-turn/session.jsonl # examples/acp-agent/tests/snapshots/todo-write/session.jsonl # examples/acp-agent/tests/snapshots/tool-call-turn/session.jsonl # examples/acp-agent/tests/snapshots/web-fetch/session.jsonl # examples/acp-agent/tests/snapshots/workflow-run/session.1.jsonl # examples/acp-agent/tests/snapshots/workflow-run/session.jsonl # examples/acp-agent/tests/snapshots/workspace-context/session.jsonl # examples/acp-agent/tests/snapshots/workspace-edit/session.jsonl # examples/headless-agent/tests/snapshots/goal-tools/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/pty-tools/session.jsonl # examples/headless-agent/tests/snapshots/pty-tools/stream-json.expected.jsonl # examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/child.expected.jsonl # examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/parent.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/persistent-tools/notifications.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/persistent-tools/session.jsonl # packages/bash/tool-bash/tests/integration.spec.ts # packages/context/time-context/src/index.ts # packages/context/tmux-context/src/index.ts # packages/core/agent-loop/src/agent.ts # packages/core/system-prompt/src/index.ts # packages/goal/goal/src/domain.ts # packages/goal/goal/src/index.ts # packages/goal/goal/src/render.ts # packages/plan/plan-mode/src/index.ts
This commit is contained in:
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/subagent/subagent/README.md
|
||||
README.md: e54f0b98ec3649cec428a47026e6657a9749608b
|
||||
README.zh.md: 074117cdbad52d754a449a0b6adc10daffa0c2f1
|
||||
README.md: a21aa6ae2822d68d513fd9409d77b3f3bf74a7a3
|
||||
README.zh.md: 3caa612aefcdac4f1dcdbcf4a3c1b81adc52c3d3
|
||||
|
||||
@@ -4,21 +4,7 @@ English | [中文](README.zh.md)
|
||||
|
||||
The subagent seam lets one agent delegate work to a child through a named provider. Callers use one service API (`ctx.subagents`); providers decide whether the child runs in this process, in another process, or through a future transport.
|
||||
|
||||
## Package roles
|
||||
|
||||
The family separates the stable interface from implementations and model-facing tools:
|
||||
|
||||
| Package | Role |
|
||||
|---|---|
|
||||
| `@deepseek-ai/dsh-subagent` | Provider registry, request/result/descriptor types, lifecycle events, and continuable-child orchestration. |
|
||||
| `@deepseek-ai/dsh-subagent-spawn` | Fresh in-process child; supports continuable children. |
|
||||
| `@deepseek-ai/dsh-subagent-fork` | In-process child seeded with completed parent turns; supports continuable children. |
|
||||
| `@deepseek-ai/dsh-subagent-acp` | Fresh out-of-process ACP child (one-shot). |
|
||||
| `@deepseek-ai/dsh-tool-subagent` | Model-facing delegation tool over one configured provider. |
|
||||
| `@deepseek-ai/dsh-tool-subagent-control` | The globally named `send_message` follow-up tool. |
|
||||
| `@deepseek-ai/dsh-tool-subagent-report` | Child-scoped return channel to the direct parent. |
|
||||
|
||||
Multiple providers may coexist under different names. This lets a deployment expose, for example, a cheap in-process child and an isolated ACP child without changing the service contract.
|
||||
The [subagent family overview](../README.md) maps implementations and model-facing consumers. This package owns the provider registry, shared request and result contracts, durable descriptors, and continuable-child orchestration. Multiple named providers may coexist behind that contract.
|
||||
|
||||
## Service API
|
||||
|
||||
|
||||
@@ -4,21 +4,7 @@
|
||||
|
||||
subagent seam 允许一个 agent(智能体)通过具名提供方把工作委派给子 agent。调用方使用统一的服务 API(`ctx.subagents`);提供方决定子 agent 在当前进程、另一进程还是未来的传输之上运行。
|
||||
|
||||
## 包角色
|
||||
|
||||
该系列包把稳定接口与实现、面向模型的工具分开:
|
||||
|
||||
| 包 | 角色 |
|
||||
|---|---|
|
||||
| `@deepseek-ai/dsh-subagent` | 提供方注册表、请求/结果/描述符类型、生命周期事件和可继续子 agent 编排。 |
|
||||
| `@deepseek-ai/dsh-subagent-spawn` | 全新的进程内子 agent;支持可继续子 agent。 |
|
||||
| `@deepseek-ai/dsh-subagent-fork` | 以父 agent 已完成轮次作为初始内容的进程内子 agent;支持可继续子 agent。 |
|
||||
| `@deepseek-ai/dsh-subagent-acp` | 全新的进程外 ACP(Agent Client Protocol)子 agent(一次性)。 |
|
||||
| `@deepseek-ai/dsh-tool-subagent` | 基于一个已配置提供方、面向模型的委派工具。 |
|
||||
| `@deepseek-ai/dsh-tool-subagent-control` | 全局具名 `send_message` 后续操作工具。 |
|
||||
| `@deepseek-ai/dsh-tool-subagent-report` | 子级作用域的返回通道,指向直接父级。 |
|
||||
|
||||
多个提供方可以使用不同名称共存。因此,部署可以同时公开低成本的进程内子 agent 和隔离的 ACP 子 agent,而无需改变服务契约。
|
||||
[subagent 家族概述](../README.md)列出了实现和面向模型的消费方。本包负责提供方注册表、共享请求和结果契约、持久描述符以及可继续子级编排。多个具名提供方可以在该契约背后共存。
|
||||
|
||||
## 服务 API
|
||||
|
||||
@@ -62,7 +48,7 @@ subagent seam 允许一个 agent(智能体)通过具名提供方把工作委
|
||||
|
||||
该 seam 拥有实现和消费方共享的深度词汇:`AgentOptions.subagentDepth` 声明、`assertSubagentMaxDepth` 和 `delegationDepthOf(agent)`。持久化的 `SessionHeader.delegationDepth` 具有权威性且单调:运行时选项可以加深计数,但绝不能降低它,因此恢复后的子 agent 不会被重新计为顶层。
|
||||
|
||||
`inheritsParentContext` 只用于描述,不能强制执行。它仅说明子 agent 是否能看到父级已完成的对话历史(`fork` 可以;`spawn` 和 ACP 不可以),不表示是否继承工具、服务或权限。
|
||||
`inheritsParentContext` 只用于描述,不能强制执行。它仅说明子 agent 是否能看到父级已完成的对话历史(`fork` 可以;`spawn` 和 ACP(Agent Client Protocol)不可以),不表示是否继承工具、服务或权限。
|
||||
|
||||
## 一次性所有权与生命周期
|
||||
|
||||
|
||||
@@ -693,7 +693,7 @@ export class SubagentContinuationManager {
|
||||
}
|
||||
|
||||
/**
|
||||
* Cold-resume a persisted child: load and authorize its Session, fold the
|
||||
* Cold-resume a persisted child: inspect and authorize its Session, fold the
|
||||
* generic descriptor, create the Activation through `ctx.agents.resume()`,
|
||||
* and submit the waiting turn. This never dispatches through a subagent
|
||||
* provider — the persisted Session already holds the initial prefix and the
|
||||
@@ -706,13 +706,13 @@ export class SubagentContinuationManager {
|
||||
options: SubagentFollowupOptions,
|
||||
): Promise<MessageId> {
|
||||
const persistence = this.requirePersistence()
|
||||
let loaded: Awaited<ReturnType<typeof persistence.load>>
|
||||
let loaded: Awaited<ReturnType<typeof persistence.inspect>>
|
||||
try {
|
||||
loaded = await persistence.load(childId)
|
||||
loaded = await persistence.inspect(childId, options.signal)
|
||||
} catch (error: unknown) {
|
||||
options.signal.throwIfAborted()
|
||||
throw new SubagentError(`subagent "${childId}" is unavailable`, 'NOT_RESUMABLE', { cause: error })
|
||||
}
|
||||
// The persistence seam takes no signal; recheck before any child work.
|
||||
options.signal.throwIfAborted()
|
||||
this.assertAdmitting(parent)
|
||||
// Authorize the persisted header before folding: only the durable child's
|
||||
@@ -729,17 +729,24 @@ export class SubagentContinuationManager {
|
||||
'NOT_RESUMABLE',
|
||||
)
|
||||
}
|
||||
const activation = await this.materialize({
|
||||
childId,
|
||||
provider: descriptor.provider,
|
||||
parent,
|
||||
agentOptions: {
|
||||
...descriptor.agentProvider !== undefined ? { provider: descriptor.agentProvider } : {},
|
||||
...descriptor.agentModel !== undefined ? { model: descriptor.agentModel } : {},
|
||||
},
|
||||
composition: { persona: descriptor.persona, toolFilter: descriptor.toolFilter },
|
||||
signal: options.signal,
|
||||
})
|
||||
let activation: Activation
|
||||
try {
|
||||
activation = await this.materialize({
|
||||
childId,
|
||||
provider: descriptor.provider,
|
||||
parent,
|
||||
agentOptions: {
|
||||
...descriptor.agentProvider !== undefined ? { provider: descriptor.agentProvider } : {},
|
||||
...descriptor.agentModel !== undefined ? { model: descriptor.agentModel } : {},
|
||||
},
|
||||
composition: { persona: descriptor.persona, toolFilter: descriptor.toolFilter },
|
||||
signal: options.signal,
|
||||
})
|
||||
} catch (error: unknown) {
|
||||
options.signal.throwIfAborted()
|
||||
if (error instanceof SubagentError) throw error
|
||||
throw new SubagentError(`subagent "${childId}" is unavailable`, 'NOT_RESUMABLE', { cause: error })
|
||||
}
|
||||
return this.submitMaterialized(activation, content, options.source, parent, options.signal)
|
||||
}
|
||||
|
||||
@@ -852,16 +859,13 @@ export class SubagentContinuationManager {
|
||||
// quiet Agent from one whose accepted turn has not been admitted yet.
|
||||
// Registered through the child's own scoped context, so scope filtering
|
||||
// already restricts both listeners to this exact agent.
|
||||
handle.agent.ctx.on('agent/inbox/dequeue', (_agent, item) => {
|
||||
/* v8 ignore next -- a dequeue of an id this manager never admitted needs
|
||||
handle.agent.ctx.on('agent/inbox/claimed', (_agent, { message }) => {
|
||||
/* v8 ignore next -- a claim of an id this manager never admitted needs
|
||||
* another sender on the same child, which no current path allows. */
|
||||
if (activation.accepted.delete(item.message.id)) this.wake(activation)
|
||||
if (activation.accepted.delete(message.id)) this.wake(activation)
|
||||
})
|
||||
handle.agent.ctx.on('agent/inbox/discard', (_agent, items) => {
|
||||
// Deleting every id in the batch is unconditional; waking once afterwards
|
||||
// costs nothing and avoids branching on which ids this manager admitted.
|
||||
for (const item of items) activation.accepted.delete(item.message.id)
|
||||
this.wake(activation)
|
||||
handle.agent.ctx.on('agent/inbox/discarded', (_agent, { message }) => {
|
||||
if (activation.accepted.delete(message.id)) this.wake(activation)
|
||||
})
|
||||
// Agent creation committed setup at its publication boundary;
|
||||
// revocations from here on are immediate live revocation.
|
||||
|
||||
@@ -207,7 +207,6 @@ function epochStopReason(events: readonly SessionEvent[]): SubagentResult['stopR
|
||||
return 'max-tokens'
|
||||
case 'aborted':
|
||||
case 'interrupted':
|
||||
case 'disposed':
|
||||
return 'aborted'
|
||||
case 'error':
|
||||
return 'error'
|
||||
|
||||
@@ -1,8 +1,9 @@
|
||||
/**
|
||||
* Read-only interpretation of session-query lineage as durable subagent
|
||||
* children. The module owns no catalog state and does not consult Activation,
|
||||
* Agent-registry, continuation-manager, or provider state. A child's
|
||||
* descriptor distinguishes one-shot work from a continuable conversation.
|
||||
* children. Only descendants with durable `origin: 'subagent'` enter per-child
|
||||
* inspection. The module owns no catalog state and does not consult Activation,
|
||||
* Agent-registry, continuation-manager, or provider state. A child's descriptor
|
||||
* distinguishes one-shot work from a continuable conversation.
|
||||
*
|
||||
* @module @deepseek-ai/dsh-subagent
|
||||
*/
|
||||
@@ -20,12 +21,13 @@ type SessionQueryRuntime = Pick<
|
||||
>
|
||||
|
||||
/**
|
||||
* One entry of a {@link listChildren} result in trace candidate order. A valid
|
||||
* descriptor produces a `child`, a per-child inspection failure produces a
|
||||
* `diagnostic`, and a descriptor-less ordinary child is omitted. Healthy rows
|
||||
* include a one-level, origin-classified descendant hint. Diagnostics are
|
||||
* transient query results, never session events or catalog state, and never
|
||||
* expose model-hidden descriptor content.
|
||||
* One entry of a {@link listChildren} result in trace candidate order. Only a
|
||||
* candidate whose durable header has `origin: 'subagent'` is inspected. A
|
||||
* valid descriptor produces a `child`, a per-child inspection failure produces
|
||||
* a `diagnostic`, and a candidate without its own descriptor is omitted.
|
||||
* Healthy rows include a one-level, origin-classified descendant hint.
|
||||
* Diagnostics are transient query results, never session events or catalog
|
||||
* state, and never expose model-hidden descriptor content.
|
||||
*/
|
||||
export type SubagentListEntry =
|
||||
| {
|
||||
@@ -69,8 +71,9 @@ export type SubagentListEntry =
|
||||
}
|
||||
|
||||
/**
|
||||
* Interpret one parent's direct session descendants as session-backed subagents
|
||||
* without loading or resuming an Agent.
|
||||
* Interpret one parent's origin-classified direct descendants as session-backed
|
||||
* subagents without loading or resuming an Agent. Ordinary forks are skipped
|
||||
* before per-child event inspection.
|
||||
* @see {@link SubagentService.listChildren} for the public cancellation and
|
||||
* failure contract.
|
||||
* @param ctx - context carrying the optional session-query service.
|
||||
@@ -103,6 +106,7 @@ export async function listChildren(
|
||||
)
|
||||
const entries: SubagentListEntry[] = []
|
||||
for (const node of trace.descendants) {
|
||||
if (node.session.header.origin !== 'subagent') continue
|
||||
const hasChildren = node.descendants.some(
|
||||
descendant => descendant.session.header.origin === 'subagent',
|
||||
)
|
||||
@@ -218,6 +222,8 @@ function perChildDiagnosticReason(
|
||||
): 'corrupt' | 'unavailable' | undefined {
|
||||
if (!(error instanceof SessionQueryError)) return undefined
|
||||
switch (error.code) {
|
||||
case 'SESSION_QUERY_CORRUPT_SESSION':
|
||||
return 'corrupt'
|
||||
case 'SESSION_QUERY_SESSION_NOT_FOUND':
|
||||
case 'SESSION_QUERY_EVENT_NOT_FOUND':
|
||||
case 'SESSION_QUERY_PERSISTENCE_FAILED':
|
||||
|
||||
@@ -142,11 +142,24 @@ async function waitNoActivation(ctx: Context, childId: SessionId): Promise<void>
|
||||
}, { timeout: 5_000 })
|
||||
}
|
||||
|
||||
/** Observe calls at the Agent cancellation boundary without a production event. */
|
||||
function observeCancel(agent: Agent, callback: () => void): void {
|
||||
const cancel = agent.cancel.bind(agent)
|
||||
let observed = false
|
||||
vi.spyOn(agent, 'cancel').mockImplementation((cause, options) => {
|
||||
if (!observed) {
|
||||
observed = true
|
||||
callback()
|
||||
}
|
||||
cancel(cause, options)
|
||||
})
|
||||
}
|
||||
|
||||
describe('SubagentService.startContinuable', () => {
|
||||
it('returns both identities at inbox acceptance, without waiting for the turn or the log', async () => {
|
||||
const { ctx, parent, adapter } = await setup([textResponse('first answer')])
|
||||
const enqueued: { id: MessageId; loggedYet: boolean }[] = []
|
||||
ctx.on('agent/inbox/enqueue', (agent, accepted) => {
|
||||
ctx.on('agent/inbox/inserted', (agent, accepted) => {
|
||||
// Acceptance is the boundary `startContinuable` resolves at, so observe
|
||||
// the log state exactly there rather than after later microtasks.
|
||||
enqueued.push({ id: accepted.message.id, loggedYet: hasUserText(agent.session.events, 'child task') })
|
||||
@@ -555,6 +568,47 @@ describe('SubagentService.followup residency routing', () => {
|
||||
.rejects.toMatchObject({ code: 'NOT_RESUMABLE' })
|
||||
})
|
||||
|
||||
it('propagates cancellation while inspecting a cold child', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('first')])
|
||||
const started = await ctx.subagents.startContinuable(startSpec(parent))
|
||||
await waitNoActivation(ctx, started.childId)
|
||||
const inspectStarted = Promise.withResolvers<undefined>()
|
||||
const inspect = vi.spyOn(ctx.sessionPersistence, 'inspect').mockImplementation((_id, signal) => {
|
||||
return new Promise<never>((_resolve, reject) => {
|
||||
if (signal === undefined) {
|
||||
reject(new Error('cold inspection must receive the followup signal'))
|
||||
return
|
||||
}
|
||||
inspectStarted.resolve(undefined)
|
||||
signal.addEventListener('abort', () => {
|
||||
reject(reason)
|
||||
}, { once: true })
|
||||
})
|
||||
})
|
||||
const controller = new AbortController()
|
||||
const reason = new Error('cold inspection cancelled')
|
||||
|
||||
try {
|
||||
const delivery = followup(ctx, parent, started.childId, message('cancel me'), controller.signal)
|
||||
await inspectStarted.promise
|
||||
controller.abort(reason)
|
||||
await expect(delivery).rejects.toBe(reason)
|
||||
} finally {
|
||||
inspect.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('preserves a SubagentError raised while cold-materializing a child', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('first')])
|
||||
const started = await ctx.subagents.startContinuable(startSpec(parent))
|
||||
await waitNoActivation(ctx, started.childId)
|
||||
const failure = new SubagentError('materialization denied', 'UNAUTHORIZED')
|
||||
ctx.agents.resume = () => Promise.reject(failure)
|
||||
|
||||
await expect(followup(ctx, parent, started.childId, message('continue')))
|
||||
.rejects.toBe(failure)
|
||||
})
|
||||
|
||||
it('cold-resumes a delivery that lost the race with final disposal', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('first'), textResponse('after the race')])
|
||||
const started = await ctx.subagents.startContinuable(startSpec(parent))
|
||||
@@ -737,7 +791,9 @@ describe('continuable durability and teardown', () => {
|
||||
const grandchild = await ctx.subagents.startContinuable(startSpec(targetChild))
|
||||
await vi.waitFor(() => { expect(adapter.requests).toHaveLength(3) })
|
||||
const cancellations: SessionId[] = []
|
||||
ctx.on('agent/cancel-requested', (agent) => { cancellations.push(agent.id) })
|
||||
observeCancel(targetChild, () => { cancellations.push(targetChild.id) })
|
||||
const grandchildAgent = ctx.agents.get(grandchild.childId)!
|
||||
observeCancel(grandchildAgent, () => { cancellations.push(grandchildAgent.id) })
|
||||
|
||||
const drained = ctx.subagents.drainContinuableDescendants([parent])
|
||||
const convergedDrain = ctx.subagents.drainContinuableDescendants([parent])
|
||||
@@ -784,7 +840,8 @@ describe('continuable durability and teardown', () => {
|
||||
const grandchild = await ctx.subagents.startContinuable(startSpec(child))
|
||||
await vi.waitFor(() => { expect(adapter.requests).toHaveLength(2) })
|
||||
const cancellations: SessionId[] = []
|
||||
ctx.on('agent/cancel-requested', (agent) => { cancellations.push(agent.id) })
|
||||
const grandchildAgent = ctx.agents.get(grandchild.childId)!
|
||||
observeCancel(grandchildAgent, () => { cancellations.push(grandchildAgent.id) })
|
||||
|
||||
const drained = ctx.subagents.drainContinuableDescendants([child])
|
||||
|
||||
@@ -828,7 +885,8 @@ describe('continuable durability and teardown', () => {
|
||||
expect(ctx.agents.get(intermediateId)).toBeUndefined()
|
||||
expect(ctx.agents.get(descendant.childId)).toBeDefined()
|
||||
const cancellations: SessionId[] = []
|
||||
ctx.on('agent/cancel-requested', (agent) => { cancellations.push(agent.id) })
|
||||
const descendantAgent = ctx.agents.get(descendant.childId)!
|
||||
observeCancel(descendantAgent, () => { cancellations.push(descendantAgent.id) })
|
||||
|
||||
const drained = ctx.subagents.drainContinuableDescendants([parent])
|
||||
|
||||
@@ -926,7 +984,7 @@ describe('continuable durability and teardown', () => {
|
||||
const drains: Promise<void>[] = []
|
||||
const accepted: MessageId[] = []
|
||||
ctx.on('subagent/start', () => { drains.push(drainManager(ctx)) })
|
||||
ctx.on('agent/inbox/enqueue', (_agent, item) => { accepted.push(item.message.id) })
|
||||
ctx.on('agent/inbox/inserted', (_agent, item) => { accepted.push(item.message.id) })
|
||||
|
||||
await expect(ctx.subagents.startContinuable(startSpec(parent)))
|
||||
.rejects.toMatchObject({ code: 'DRAINING' })
|
||||
@@ -967,12 +1025,12 @@ describe('continuable durability and teardown', () => {
|
||||
await vi.waitFor(() => { expect(adapter.requests).toHaveLength(1) })
|
||||
const child = ctx.agents.get(started.childId)!
|
||||
const order: string[] = []
|
||||
child.ctx.on('agent/inbox/enqueue', (_agent, accepted) => {
|
||||
child.ctx.on('agent/inbox/inserted', (_agent, accepted) => {
|
||||
if (accepted.message.content.some(block => block.type === 'text' && block.text === 'before drain')) {
|
||||
order.push('enqueue')
|
||||
}
|
||||
})
|
||||
child.ctx.on('agent/cancel-requested', () => { order.push('cancel') })
|
||||
observeCancel(child, () => { order.push('cancel') })
|
||||
|
||||
const delivery = followup(ctx, parent, started.childId, message('before drain'))
|
||||
// Let the child-lock operation reach the live admission cutoff. Admission
|
||||
@@ -1150,9 +1208,9 @@ describe('continuable review regressions', () => {
|
||||
const ends: SubagentRunEndInfo[] = []
|
||||
ctx.on('subagent/end', (info) => { ends.push(info) })
|
||||
// Block the resumed prompt so this epoch produces nothing of its own.
|
||||
ctx.on('agent/prompt-submit', async (subject, _message, _signal, next) => {
|
||||
ctx.on('agent/pre-step', async (subject, _messages, _context, next) => {
|
||||
if (subject === parent) return next()
|
||||
return { kind: 'block', reason: 'blocked by policy' }
|
||||
return { kind: 'reject' }
|
||||
})
|
||||
await followup(ctx, parent, started.childId, message('again'))
|
||||
await waitNoActivation(ctx, started.childId)
|
||||
@@ -1258,7 +1316,7 @@ describe('continuable review regressions', () => {
|
||||
expect(found).toBeDefined()
|
||||
return found!
|
||||
})
|
||||
child.ctx.on('agent/cancel-requested', () => { order.push('cancel') })
|
||||
observeCancel(child, () => { order.push('cancel') })
|
||||
|
||||
const drained = drainManager(ctx)
|
||||
hold.resolve(undefined)
|
||||
@@ -1298,7 +1356,7 @@ describe('continuable review regressions', () => {
|
||||
|
||||
// Cancel from the synchronous enqueue observer: the discard fires after the
|
||||
// id is recorded but before `followup()` returns.
|
||||
const off = child.ctx.on('agent/inbox/enqueue', (_agent, accepted) => {
|
||||
const off = child.ctx.on('agent/inbox/inserted', (_agent, accepted) => {
|
||||
if (accepted.message.content.some(block => block.type === 'text' && block.text === 'doomed')) {
|
||||
child.cancel({ kind: 'user' })
|
||||
}
|
||||
@@ -1330,7 +1388,7 @@ describe('continuable review regressions', () => {
|
||||
|
||||
await followup(ctx, parent, started.childId, message('queued'))
|
||||
expect(activation.accepted.size).toBe(1)
|
||||
const off = child.ctx.on('agent/inbox/enqueue', (_agent, accepted) => {
|
||||
const off = child.ctx.on('agent/inbox/inserted', (_agent, accepted) => {
|
||||
if (accepted.message.content.some(block => block.type === 'text' && block.text === 'doomed')) {
|
||||
child.cancel({ kind: 'user' })
|
||||
}
|
||||
@@ -1348,9 +1406,9 @@ describe('continuable review regressions', () => {
|
||||
const ends: SubagentRunEndInfo[] = []
|
||||
ctx.on('subagent/end', (info) => { ends.push(info) })
|
||||
// Block admission so the child's only turn never opens.
|
||||
ctx.on('agent/prompt-submit', async (subject, _message, _signal, next) => {
|
||||
ctx.on('agent/pre-step', async (subject, _messages, _context, next) => {
|
||||
if (subject === parent) return next()
|
||||
return { kind: 'block', reason: 'blocked by policy' }
|
||||
return { kind: 'reject' }
|
||||
})
|
||||
|
||||
const started = await ctx.subagents.startContinuable(startSpec(parent))
|
||||
@@ -1370,7 +1428,7 @@ describe('continuable review regressions', () => {
|
||||
const registeredAtEnqueue: boolean[] = []
|
||||
// A synchronous inbox observer runs before the admitting microtask, the
|
||||
// exact window where `Agent.status` is still idle.
|
||||
ctx.on('agent/inbox/enqueue', (agent) => {
|
||||
ctx.on('agent/inbox/inserted', (agent) => {
|
||||
if (agent.session.header.parentSession !== undefined) {
|
||||
registeredAtEnqueue.push(ctx.agents.get(agent.id) === agent)
|
||||
}
|
||||
|
||||
@@ -113,10 +113,11 @@ describe('SubagentService.listChildren', () => {
|
||||
const parentId = SessionId('query-only-parent')
|
||||
ctx.sessions.create(parentId)
|
||||
const childId = SessionId('query-only-child')
|
||||
const child = ctx.sessions.create(childId, { meta: { parentSession: parentId } })
|
||||
const child = ctx.sessions.create(childId, {
|
||||
meta: { parentSession: parentId, origin: 'subagent' },
|
||||
})
|
||||
child.append('turn/start', {
|
||||
turn: 1,
|
||||
trigger: { kind: 'message', source: { kind: 'user' } },
|
||||
})
|
||||
child.append('subagent/descriptor', descriptorPayload('query-only child'))
|
||||
|
||||
@@ -193,6 +194,7 @@ describe('SubagentService.listChildren', () => {
|
||||
] as SessionEvent[])
|
||||
const childId = await authorChild(ctx, '00000000-0000-4000-8000-00000000cdcd', {
|
||||
parentSession: coldParent,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('persisted parent case')))
|
||||
const entries = await ctx.subagents.listChildren(coldParent)
|
||||
expect(entries).toEqual([
|
||||
@@ -203,28 +205,33 @@ describe('SubagentService.listChildren', () => {
|
||||
])
|
||||
})
|
||||
|
||||
it('orders children by createdAt then id and omits ordinary forks without a diagnostic', async () => {
|
||||
it('orders children by createdAt then id without inspecting ordinary forks', async () => {
|
||||
const { ctx, parent } = await setup([])
|
||||
// Authored headers pin the ordering key deterministically: same createdAt
|
||||
// ties break on id, different createdAt orders ascending.
|
||||
const late = await authorChild(ctx, '00000000-0000-4000-8000-000000000003', {
|
||||
parentSession: parent.id,
|
||||
createdAt: 9,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('late child')))
|
||||
const tieB = await authorChild(ctx, '00000000-0000-4000-8000-000000000002', {
|
||||
parentSession: parent.id,
|
||||
createdAt: 5,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('tie b')))
|
||||
const tieA = await authorChild(ctx, '00000000-0000-4000-8000-000000000001', {
|
||||
parentSession: parent.id,
|
||||
createdAt: 5,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('tie a')))
|
||||
// An ordinary session fork shares parentSession but has no descriptor.
|
||||
// An ordinary session fork shares parentSession but has no subagent origin.
|
||||
const fork = ctx.sessions.fork(parent.session, undefined, SessionId('plain-fork'))
|
||||
await ctx.sessions.flush(fork)
|
||||
const listEvents = vi.spyOn(ctx.sessionQuery, 'listEvents')
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries.map(entry => entry.id)).toEqual([tieA, tieB, late])
|
||||
expect(entries.every(entry => entry.kind === 'child')).toBe(true)
|
||||
expect(listEvents).not.toHaveBeenCalledWith(fork.id)
|
||||
})
|
||||
|
||||
it('reports a live child as running while keeping settled siblings complete', async () => {
|
||||
@@ -233,8 +240,10 @@ describe('SubagentService.listChildren', () => {
|
||||
// A live child session outside persistence: publish a live session with a
|
||||
// descriptor and the parent lineage, without starting an Activation.
|
||||
const liveId = SessionId('live-child')
|
||||
const live = ctx.sessions.create(liveId, { meta: { parentSession: parent.id } })
|
||||
live.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
const live = ctx.sessions.create(liveId, {
|
||||
meta: { parentSession: parent.id, origin: 'subagent' },
|
||||
})
|
||||
live.append('turn/start', { turn: 1 })
|
||||
live.append('subagent/descriptor', descriptorPayload('live child'))
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toContainEqual({
|
||||
@@ -260,6 +269,7 @@ describe('SubagentService.listChildren', () => {
|
||||
events[4] = { ...events[4]!, seq: 4 }
|
||||
const corrupt = await authorChild(ctx, '00000000-0000-4000-8000-00000000dupe', {
|
||||
parentSession: parent.id,
|
||||
origin: 'subagent',
|
||||
}, events)
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toContainEqual({ kind: 'diagnostic', id: corrupt, reason: 'corrupt' })
|
||||
@@ -269,12 +279,13 @@ describe('SubagentService.listChildren', () => {
|
||||
})
|
||||
})
|
||||
|
||||
it('diagnoses an invalid child event surface as corrupt', async () => {
|
||||
it('diagnoses a child rejected by persisted Session preparation as corrupt', async () => {
|
||||
const { ctx, parent } = await setup([])
|
||||
// The surface-eligible user/message lacks its required surfaceOp, so the
|
||||
// per-child listEvents fold fails with SESSION_QUERY_INVALID_SURFACE.
|
||||
// The surface-eligible user/message lacks its required surfaceOp. The
|
||||
// first-party persistence inspection rejects before session-query can fold it.
|
||||
const invalid = await authorChild(ctx, '00000000-0000-4000-8000-0000000000ee', {
|
||||
parentSession: parent.id,
|
||||
origin: 'subagent',
|
||||
}, [
|
||||
{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
|
||||
{
|
||||
@@ -293,6 +304,7 @@ describe('SubagentService.listChildren', () => {
|
||||
const { ctx, parent } = await setup([])
|
||||
const malformed = await authorChild(ctx, '00000000-0000-4000-8000-0000000000ff', {
|
||||
parentSession: parent.id,
|
||||
origin: 'subagent',
|
||||
}, childEvents({ version: SUBAGENT_DESCRIPTOR_VERSION, mode: 'continuable', provider: 7 }))
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([{ kind: 'diagnostic', id: malformed, reason: 'corrupt' }])
|
||||
@@ -302,6 +314,7 @@ describe('SubagentService.listChildren', () => {
|
||||
const { ctx, parent } = await setup([])
|
||||
const future = await authorChild(ctx, '00000000-0000-4000-8000-0000000000aa', {
|
||||
parentSession: parent.id,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('from the future', SUBAGENT_DESCRIPTOR_VERSION + 1)))
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([{ kind: 'diagnostic', id: future, reason: 'unsupported' }])
|
||||
@@ -315,6 +328,7 @@ describe('SubagentService.listChildren', () => {
|
||||
await authorChild(ctx, '00000000-0000-4000-8000-0000000000f0', {
|
||||
parentSession: parent.id,
|
||||
seedLength: seed.length,
|
||||
origin: 'subagent',
|
||||
}, seed)
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([])
|
||||
@@ -324,6 +338,7 @@ describe('SubagentService.listChildren', () => {
|
||||
const { ctx, parent } = await setup([])
|
||||
const foreign = await authorChild(ctx, '00000000-0000-4000-8000-0000000000bb', {
|
||||
parentSession: parent.id,
|
||||
origin: 'subagent',
|
||||
}, childEvents({
|
||||
version: SUBAGENT_DESCRIPTOR_VERSION,
|
||||
mode: 'continuable',
|
||||
@@ -354,16 +369,30 @@ describe('SubagentService.listChildren', () => {
|
||||
expect(entries).toEqual([{ kind: 'diagnostic', id: childId, reason: 'unavailable' }])
|
||||
})
|
||||
|
||||
it('maps a mid-scan disappearance to unavailable', async () => {
|
||||
it.each([
|
||||
['session', 'SESSION_QUERY_SESSION_NOT_FOUND'],
|
||||
['descriptor event', 'SESSION_QUERY_EVENT_NOT_FOUND'],
|
||||
] as const)('maps a missing child %s to unavailable', async (_target, code) => {
|
||||
const { ctx, parent } = await setup([textResponse('done')])
|
||||
const childId = await startChild(ctx, parent, 'vanishing child')
|
||||
const query = ctx.get('sessionQuery')!
|
||||
query.listEvents = () =>
|
||||
Promise.reject(new SessionQueryError('gone', 'SESSION_QUERY_SESSION_NOT_FOUND'))
|
||||
Promise.reject(new SessionQueryError('gone', code))
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([{ kind: 'diagnostic', id: childId, reason: 'unavailable' }])
|
||||
})
|
||||
|
||||
it('maps an invalid child surface to corrupt', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('done')])
|
||||
const childId = await startChild(ctx, parent, 'invalid surface')
|
||||
const query = ctx.get('sessionQuery')!
|
||||
query.listEvents = () =>
|
||||
Promise.reject(new SessionQueryError('invalid surface', 'SESSION_QUERY_INVALID_SURFACE'))
|
||||
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([{ kind: 'diagnostic', id: childId, reason: 'corrupt' }])
|
||||
})
|
||||
|
||||
it('diagnoses a read whose header no longer names this parent as corrupt', async () => {
|
||||
const { ctx, parent } = await setup([textResponse('done')])
|
||||
const childId = await startChild(ctx, parent, 'reparented child')
|
||||
@@ -429,6 +458,7 @@ describe('SubagentService.listChildren', () => {
|
||||
const plain = await authorChild(ctx, '00000000-0000-4000-8000-00000000c0de', {
|
||||
parentSession: parent.id,
|
||||
createdAt: 1,
|
||||
origin: 'subagent',
|
||||
}, childEvents(descriptorPayload('twin child')))
|
||||
// The compacted twin: a compaction checkpoint replaces the whole surface,
|
||||
// while the append-only log retains the model-hidden descriptor event.
|
||||
@@ -447,6 +477,7 @@ describe('SubagentService.listChildren', () => {
|
||||
const compacted = await authorChild(ctx, '00000000-0000-4000-8000-00000000c1de', {
|
||||
parentSession: parent.id,
|
||||
createdAt: 2,
|
||||
origin: 'subagent',
|
||||
}, compactedEvents)
|
||||
const entries = await ctx.subagents.listChildren(parent.id)
|
||||
expect(entries).toEqual([
|
||||
|
||||
Reference in New Issue
Block a user