diff --git a/.agents/notes/implemented/feature/2026-08-05-agent-teams.i18n.yaml b/.agents/notes/implemented/feature/2026-08-05-agent-teams.i18n.yaml index 50b5887ede..7b01f5858a 100644 --- a/.agents/notes/implemented/feature/2026-08-05-agent-teams.i18n.yaml +++ b/.agents/notes/implemented/feature/2026-08-05-agent-teams.i18n.yaml @@ -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 .agents/notes/implemented/feature/2026-08-05-agent-teams.md -2026-08-05-agent-teams.md: 28828f6d0672dfc8e9e44b1763a7b9eb69a61139 -2026-08-05-agent-teams.zh.md: 2e2296126c54046905e89babdcf1c1df1bc84f54 +2026-08-05-agent-teams.md: aedd31bb4ec2b1833390350fbb0b0740f08ec4cf +2026-08-05-agent-teams.zh.md: 4b2268e07f87249d0c0bc9eb336baf60629b0c0c diff --git a/.agents/notes/implemented/feature/2026-08-05-agent-teams.md b/.agents/notes/implemented/feature/2026-08-05-agent-teams.md index 28828f6d06..aedd31bb4e 100644 --- a/.agents/notes/implemented/feature/2026-08-05-agent-teams.md +++ b/.agents/notes/implemented/feature/2026-08-05-agent-teams.md @@ -34,7 +34,7 @@ Peer communication is a Lead-log mailbox. `team/message/queued` is appended and Quiet `send_message` injects, flushes, and acknowledges immediately for a live target without waking it; an inactive target remains queued until another event materializes that teammate. Waking `followup_task` becomes the target's next FIFO turn and may cold-resume it. Success means the message is already durable even when immediate delivery is deferred. The mechanism provides process-local retry and target-Session de-duplication, not a cross-process exactly-once claim. -Shared tasks are complete snapshots with Team-local ids and monotonic revisions. Every mutation carries `expectedRevision`. Any member creates, reads, or claims a ready unowned task; the owner or Lead edits and transitions it, while only the Lead assigns another member. Dependencies must name non-deleted tasks and form a complete DAG. Deleted tasks are retained tombstones. `writeScopes` are normalized path prefixes that produce overlap diagnostics but never block claim or authorize a write. +Shared tasks are complete snapshots with Team-local ids and monotonic revisions. Every mutation carries `expectedRevision`. Any member creates, reads, or claims a ready unowned task; the owner or Lead edits and transitions it, while only the Lead assigns another member. Numeric task ids remain within the safe-integer allocation range, and exhaustion fails without reusing an id. Dependencies must name non-deleted tasks and form a complete DAG. Deleted tasks are retained tombstones. `writeScopes` are normalized path prefixes that produce overlap diagnostics but never block claim or authorize a write. `wait_agent` blocks on one roster, mailbox, task, or live-status edge registered after the call starts instead of encouraging model polling. It does not replay an earlier edge, so callers re-read authoritative state after wakeup or timeout. Lead-only interruption cancels the current turn with inbox preservation and does not alter mailbox or task ownership. @@ -62,7 +62,7 @@ Worktree isolation is not a harness runtime behavior. A deployment or prompt may ## Testing -Package tests cover identity, name and authority checks, provider selection, reserved-id persistence collisions, child-before-Lead flush ordering, durable provisioning failure and pending-inbox JSONL/SQLite reconciliation, concurrent target-local ordering, pending/history de-duplication, mailbox limits, post-flush notification, bounded disposal with in-flight creation and dispatch cancellation, failed-member cleanup, task CAS and DAG validation, write-scope warnings, wait cancellation/timeout, inbox-preserving interruption, ordinary-fork isolation, legacy-control shadowing, compact declared-schema result rendering, and scoped registration HMR at per-file 100% coverage. +Package tests cover identity, name and authority checks, provider selection, reserved-id persistence collisions, child-before-Lead flush ordering, durable provisioning failure and pending-inbox JSONL/SQLite reconciliation, concurrent target-local ordering, pending/history de-duplication, mailbox limits, post-flush notification, bounded disposal with in-flight creation and dispatch cancellation, failed-member cleanup, task CAS and DAG validation, write-scope warnings, wait cancellation/timeout, inbox-preserving interruption, ordinary-fork isolation, legacy-control shadowing, compact declared-schema result rendering, and scoped registration HMR at per-file 100% coverage. A keyless headless Loader snapshot assembles the real Team plugins and records teammate creation, peer mail, dependent tasks, waiting, and Lead aggregation. ## Consequences diff --git a/.agents/notes/implemented/feature/2026-08-05-agent-teams.zh.md b/.agents/notes/implemented/feature/2026-08-05-agent-teams.zh.md index 2e2296126c..4b2268e07f 100644 --- a/.agents/notes/implemented/feature/2026-08-05-agent-teams.zh.md +++ b/.agents/notes/implemented/feature/2026-08-05-agent-teams.zh.md @@ -34,7 +34,7 @@ Peer 通讯使用 Lead 日志 mailbox。投递前先追加并 flush `team/messag 对于 live target,quiet `send_message` 会立即注入、flush 并确认,但不会唤醒它;inactive target 会保持 queued,直到其他事件 materialize 该 teammate。waking `followup_task` 成为 target 的下一个 FIFO turn,并可冷恢复。即使即时投递被推迟,成功也表示消息已经持久化。该机制提供进程内重试与 target Session 去重,不宣称跨进程 exactly-once。 -共享 task 是带 Team-local id 与单调 revision 的完整快照。每次变更都携带 `expectedRevision`。任意 member 可以创建、读取或 claim ready 且无 owner 的任务;Owner 或 Lead 可以编辑和转换;只有 Lead 可以分配给另一个 member。依赖必须指向未删除任务,并形成完整 DAG。删除任务保留为 tombstone。`writeScopes` 是规范化路径前缀,只产生重叠诊断,绝不会阻止 claim 或授予写权限。 +共享 task 是带 Team-local id 与单调 revision 的完整快照。每次变更都携带 `expectedRevision`。任意 member 可以创建、读取或 claim ready 且无 owner 的任务;Owner 或 Lead 可以编辑和转换;只有 Lead 可以分配给另一个 member。数字 task id 保持在安全整数分配范围内;该范围耗尽时会失败,不会复用 id。依赖必须指向未删除任务,并形成完整 DAG。删除任务保留为 tombstone。`writeScopes` 是规范化路径前缀,只产生重叠诊断,绝不会阻止 claim 或授予写权限。 `wait_agent` 等待调用注册后发生的下一条 roster、mailbox、task 或实时 status 边,避免模型轮询。它不会回放更早的边,因此调用方需要在唤醒或超时后重新读取权威状态。仅限 Lead 的 interrupt 使用 inbox preservation 取消当前 turn,不改变 mailbox 或 task owner。 @@ -62,7 +62,7 @@ Worktree isolation 不是 harness runtime 行为。deployment 或 prompt 可以 ## Testing -Package test 以逐文件 100% coverage 覆盖身份、名字与权限检查、provider 选择、预留 id 持久化冲突、child-before-Lead flush 顺序、持久 provisioning 失败与 pending-inbox JSONL/SQLite 对账、target-local 并发顺序、pending/history 去重、mailbox 限额、flush 后 notification、取消在途创建与 dispatch 的有界 dispose、failed member cleanup、task CAS 与 DAG 校验、write-scope warning、wait cancel/timeout、保留 inbox 的 interrupt、普通 fork 隔离、旧 control shadowing、声明 schema 的紧凑结果渲染与 scoped registration HMR。 +Package test 以逐文件 100% coverage 覆盖身份、名字与权限检查、provider 选择、预留 id 持久化冲突、child-before-Lead flush 顺序、持久 provisioning 失败与 pending-inbox JSONL/SQLite 对账、target-local 并发顺序、pending/history 去重、mailbox 限额、flush 后 notification、取消在途创建与 dispatch 的有界 dispose、failed member cleanup、task CAS 与 DAG 校验、write-scope warning、wait cancel/timeout、保留 inbox 的 interrupt、普通 fork 隔离、旧 control shadowing、声明 schema 的紧凑结果渲染与 scoped registration HMR。一条 keyless headless Loader 快照会组合真实 Team 插件,并记录 teammate 创建、peer mail、依赖任务、等待与 Lead 汇总。 ## Consequences diff --git a/examples/headless-agent/team.cordis.snapshot.yml b/examples/headless-agent/team.cordis.snapshot.yml new file mode 100644 index 0000000000..1d0ab574f4 --- /dev/null +++ b/examples/headless-agent/team.cordis.snapshot.yml @@ -0,0 +1,36 @@ +# Keyless Agent Teams composition over the real headless app and deterministic fixture adapter. +- id: base + name: '@deepseek-ai/cordis-plugin-include' + config: + path: ./cordis.yml + patches: + - id: llm-deepseek + name: '@deepseek-ai/dsh-llm-deepseek' + disabled: true + - id: tool-subagent-control + name: '@deepseek-ai/dsh-tool-subagent-control' + disabled: true + - id: tool-subagent-report + name: '@deepseek-ai/dsh-tool-subagent-report' + disabled: true + - id: tool-subagent + name: '@deepseek-ai/dsh-tool-subagent' + config: + provider: spawn + toolName: subagent + backgroundMode: one-shot + maxDepth: 1 + - id: tool-subagent-fork + name: '@deepseek-ai/dsh-tool-subagent' + config: + provider: fork + toolName: subagent_fork + backgroundMode: one-shot + maxDepth: 1 + - insert: + - id: team + name: '@deepseek-ai/dsh-team' + - id: tool-team + name: '@deepseek-ai/dsh-tool-team' + - id: team-fixture-llm + name: './tests/fixtures/team-llm.mjs' diff --git a/examples/headless-agent/tests/fixtures/team-llm.mjs b/examples/headless-agent/tests/fixtures/team-llm.mjs new file mode 100644 index 0000000000..c1c973c771 --- /dev/null +++ b/examples/headless-agent/tests/fixtures/team-llm.mjs @@ -0,0 +1,194 @@ +/** Deterministic keyless Agent Teams adapter for the real headless Loader snapshot. */ + +import { CallId, LlmAdapter } from '@deepseek-ai/dsh-llm' + +let nextCall = 0 + +function calls(messages) { + return messages.flatMap(message => message.role === 'assistant' + ? message.content.filter(block => block.type === 'tool-call').map(block => block.name) + : []) +} + +function latestAssistantCalls(messages) { + const assistant = messages.findLast(message => message.role === 'assistant') + return assistant?.content.filter(block => block.type === 'tool-call').map(block => block.name) ?? [] +} + +function hasTaskAction(messages, action) { + return messages.some(message => message.role === 'assistant' + && message.content.some((block) => { + if (block.type !== 'tool-call' || block.name !== 'team_task_update') return false + try { + return JSON.parse(block.arguments).action === action + } catch { + return false + } + })) +} + +function latestToolText(messages) { + const message = messages.findLast(candidate => candidate.content.some(block => block.type === 'tool-result')) + if (message === undefined) return '' + return message.content.flatMap(block => block.type === 'tool-result' + ? block.content.filter(item => item.type === 'text').map(item => item.text) + : []).join('\n') +} + +function toolChunks(specs) { + const chunks = [] + for (const [index, spec] of specs.entries()) { + const id = CallId(`team-fixture-${++nextCall}`) + const args = JSON.stringify(spec.args) + chunks.push( + { type: 'block-start', index, blockType: 'tool-call' }, + { type: 'tool-call-delta', index, id, name: spec.name, argumentsDelta: args }, + { type: 'block-end', index, block: { type: 'tool-call', id, name: spec.name, arguments: args } }, + ) + } + chunks.push( + { type: 'usage', usage: { inputTokens: 10, outputTokens: 5 } }, + { type: 'finish', reason: { kind: 'tool-calls' } }, + ) + return chunks +} + +function textChunks(text) { + return [ + { type: 'block-start', index: 0, blockType: 'text' }, + { type: 'text-delta', index: 0, text }, + { type: 'block-end', index: 0, block: { type: 'text', text } }, + { type: 'usage', usage: { inputTokens: 10, outputTokens: 3 } }, + { type: 'finish', reason: { kind: 'stop' } }, + ] +} + +function researcher(messages) { + const names = calls(messages) + if (!names.includes('team_task_create')) { + return toolChunks([{ name: 'team_task_create', args: { + subject: 'Research', description: 'Collect the deterministic finding.', write_scopes: ['research'], + } }]) + } + if (!names.includes('team_task_update')) { + return toolChunks([{ name: 'team_task_update', args: { + task_id: 'task-1', expected_revision: 1, action: 'claim', + } }]) + } + if (!names.includes('send_message')) { + return toolChunks([ + { name: 'team_task_update', args: { task_id: 'task-1', expected_revision: 2, action: 'complete' } }, + { name: 'send_message', args: { target: 'implementer', message: 'Research complete: use the deterministic finding.' } }, + ]) + } + return textChunks('Research teammate complete.') +} + +function implementer(messages) { + const names = calls(messages) + const last = latestAssistantCalls(messages) + const text = latestToolText(messages) + if (!names.includes('team_task_create')) { + if (last.includes('team_task_get') && text.includes('"subject":"Research"')) { + return toolChunks([{ name: 'team_task_create', args: { + subject: 'Implementation', + description: 'Apply the deterministic finding.', + blocked_by: ['task-1'], + write_scopes: ['implementation'], + } }]) + } + if (last.includes('wait_agent')) { + return toolChunks([{ name: 'team_task_get', args: { task_id: 'task-1' } }]) + } + return toolChunks([{ name: 'wait_agent', args: { timeout_ms: 10000 } }]) + } + if (!hasTaskAction(messages, 'claim')) { + if (last.includes('team_task_get') && text.includes('"status":"completed"')) { + return toolChunks([{ name: 'team_task_update', args: { + task_id: 'task-2', expected_revision: 1, action: 'claim', + } }]) + } + if (last.includes('team_task_get')) { + return toolChunks([{ name: 'team_task_get', args: { task_id: 'task-1' } }]) + } + if (last.includes('wait_agent')) { + return toolChunks([{ name: 'team_task_get', args: { task_id: 'task-1' } }]) + } + return toolChunks([{ name: 'team_task_get', args: { task_id: 'task-1' } }]) + } + if (!names.includes('send_message')) { + return toolChunks([ + { name: 'team_task_update', args: { task_id: 'task-2', expected_revision: 2, action: 'complete' } }, + { name: 'send_message', args: { target: 'lead', message: 'Implementation complete and verified.' } }, + ]) + } + return textChunks('Implementation teammate complete.') +} + +function lead(messages) { + const names = calls(messages) + const last = latestAssistantCalls(messages) + const spawned = names.filter(name => name === 'spawn_teammate').length + if (spawned === 0) { + return toolChunks([{ + name: 'spawn_teammate', + args: { + name: 'implementer', + description: 'Own deterministic implementation.', + prompt: 'IMPLEMENTER_MARK: wait for research, complete dependent task 2, report to lead.', + context: 'fresh', + }, + }]) + } + if (spawned === 1) { + return toolChunks([{ + name: 'spawn_teammate', + args: { + name: 'researcher', + description: 'Own deterministic research.', + prompt: 'RESEARCHER_MARK: complete research task 1, message implementer, then finish.', + context: 'fresh', + }, + }]) + } + const result = latestToolText(messages) + if (last.includes('team_task_list')) { + const completed = result.match(/"status":"completed"/gu)?.length ?? 0 + if (completed >= 2) return toolChunks([{ name: 'list_agents', args: {} }]) + return toolChunks([{ name: 'team_task_list', args: {} }]) + } + if (last.includes('list_agents')) { + const inactive = result.match(/"status":"inactive"/gu)?.length ?? 0 + if (inactive >= 2) return textChunks('TEAM_WORKFLOW_OK: both teammates and dependent tasks completed.') + return toolChunks([{ name: 'list_agents', args: {} }]) + } + if (last.includes('wait_agent')) return toolChunks([{ name: 'team_task_list', args: {} }]) + return toolChunks([{ name: 'wait_agent', args: { timeout_ms: 10000 } }]) +} + +class TeamFixtureAdapter extends LlmAdapter { + async * stream(options) { + const userText = options.messages.flatMap(message => message.role === 'user' + ? message.content.filter(block => block.type === 'text').map(block => block.text) + : []).join('\n') + const chunks = userText.includes('RESEARCHER_MARK') + ? researcher(options.messages) + : userText.includes('IMPLEMENTER_MARK') + ? implementer(options.messages) + : lead(options.messages) + for (const chunk of chunks) { + options.signal?.throwIfAborted() + yield chunk + } + } +} + +/** Cordis plugin name. */ +export const name = 'team-fixture-llm' +/** LLM registry dependency. */ +export const inject = ['llm'] + +/** Register the keyless adapter on the shipped default provider route. */ +export function apply(ctx) { + ctx.llm.registerAdapter(['deepseek-official'], new TeamFixtureAdapter()) +} diff --git a/examples/headless-agent/tests/headless.snapshot.ts b/examples/headless-agent/tests/headless.snapshot.ts index d22a48a7ab..19247c3c6b 100644 --- a/examples/headless-agent/tests/headless.snapshot.ts +++ b/examples/headless-agent/tests/headless.snapshot.ts @@ -47,6 +47,7 @@ const ralphScenarioDir = join(snapshotsDir, 'ralph-loop') const ralphConfigPath = fileURLToPath(new URL('../ralph.cordis.snapshot.yml', import.meta.url)) const settlementScenarioDir = join(snapshotsDir, 'subagent-settlement') const settlementConfigPath = fileURLToPath(new URL('../subagent-settlement.cordis.snapshot.yml', import.meta.url)) +const teamConfigPath = fileURLToPath(new URL('../team.cordis.snapshot.yml', import.meta.url)) const startupFailureConfigPath = fileURLToPath(new URL('./fixtures/startup-activation-error/cordis.yml', import.meta.url)) const startupFailureExpected = join(snapshotsDir, 'startup-activation-error', 'stderr.expected.txt') const binScript = fileURLToPath(new URL('./fixtures/headless-driver.ts', import.meta.url)) @@ -645,6 +646,85 @@ describe('headless stream-json snapshots', () => { expect(normalized).toBe(await readFile(advancedStreamExpected, 'utf8')) }, LOADER_SMOKE_TEST_TIMEOUT_MS) + it('runs a keyless Agent Team with peer mail, dependent tasks, waiting, and Lead aggregation', async () => { + let projection: unknown + const result = await runLoaderSmoke({ + label: 'Agent Teams headless snapshot', + tempDirPrefix: 'headless-snapshot-agent-team-', + binScript, + libBinScript: binScript, + configPath: teamConfigPath, + binArgs: [ + teamConfigPath, + '请明确使用 Agent Teams,把调研和实现拆给两个 teammate,等待完成后汇总。', + ], + tsconfigPath, + processTimeoutMs: 60_000, + env: { + DSH_SNAPSHOT: 'team', + NODE_OPTIONS: [process.env.NODE_OPTIONS, '--disable-warning=ExperimentalWarning'].filter(Boolean).join(' '), + }, + inspect: async (cwd) => { + const logs = await persistedLogs(cwd) + const parent = logs.find(log => typeof log.header.parentSession !== 'string') + if (parent === undefined) throw new Error('Agent Teams snapshot did not persist its Lead') + const rows = parseJsonl(parent.content) + const members = rows.filter(row => row.type === 'team/member') + .map(row => ((row.data as JsonObject).member as JsonObject)) + const tasks = rows.filter(row => row.type === 'team/task') + .map(row => ((row.data as JsonObject).task as JsonObject)) + const latestTasks = Object.values(Object.fromEntries(tasks.map(task => [String(task.subject), task]))) + projection = { + sessions: logs.length, + memberEdges: members.length, + activeMembers: members.filter(member => member.phase === 'active').map(member => member.name).sort(), + tasks: latestTasks.map(task => ({ + subject: task.subject, + revision: task.revision, + status: task.status, + })).sort((left, right) => String(left.subject).localeCompare(String(right.subject))), + queuedMessages: rows.filter(row => row.type === 'team/message/queued').length, + deliveredMessages: rows.filter(row => row.type === 'team/message/delivered').length, + waited: rows.some(row => row.type === 'tool/call' + && (row.data as JsonObject).name === 'wait_agent'), + checkedRoster: rows.some(row => row.type === 'tool/call' + && (row.data as JsonObject).name === 'list_agents'), + } + }, + }) + expect(result.stderr).toBe('') + expect(parseJsonl(result.stdout).at(-1)).toMatchObject({ + type: 'result', + output: 'TEAM_WORKFLOW_OK: both teammates and dependent tasks completed.', + }) + expect(projection).toMatchInlineSnapshot(` + { + "activeMembers": [ + "implementer", + "researcher", + ], + "checkedRoster": true, + "deliveredMessages": 2, + "memberEdges": 4, + "queuedMessages": 2, + "sessions": 3, + "tasks": [ + { + "revision": 3, + "status": "completed", + "subject": "Implementation", + }, + { + "revision": 3, + "status": "completed", + "subject": "Research", + }, + ], + "waited": true, + } + `) + }, 75_000) + it('replays persisted goal tools through the one-shot app', async () => { const prompt = await scenarioPrompt(goalScenarioDir, 'goal-tools') const streamExpected = join(goalScenarioDir, 'stream-json.expected.jsonl') diff --git a/examples/package.json b/examples/package.json index cd348d5fe2..021ed0dc32 100644 --- a/examples/package.json +++ b/examples/package.json @@ -82,6 +82,7 @@ "@deepseek-ai/dsh-subprocess-local": "workspace:*", "@deepseek-ai/dsh-system-prompt": "workspace:*", "@deepseek-ai/dsh-jobs-local": "workspace:*", + "@deepseek-ai/dsh-team": "workspace:*", "@deepseek-ai/dsh-time-context": "workspace:*", "@deepseek-ai/dsh-tool-call-timeout-policy": "workspace:*", "@deepseek-ai/dsh-token-meter": "workspace:*", @@ -103,6 +104,7 @@ "@deepseek-ai/dsh-tool-subagent-control": "workspace:*", "@deepseek-ai/dsh-tool-subagent-report": "workspace:*", "@deepseek-ai/dsh-tool-jobs": "workspace:*", + "@deepseek-ai/dsh-tool-team": "workspace:*", "@deepseek-ai/dsh-tool-todo": "workspace:*", "@deepseek-ai/dsh-tool-web": "workspace:*", "@deepseek-ai/dsh-tool-workflow": "workspace:*", diff --git a/packages/team/team/README.i18n.yaml b/packages/team/team/README.i18n.yaml index 8ebba28885..cdef92faf0 100644 --- a/packages/team/team/README.i18n.yaml +++ b/packages/team/team/README.i18n.yaml @@ -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/team/team/README.md -README.md: ba8446ec1655abb273fdc0ddbc95a4cf0f938f12 -README.zh.md: 93a7006ef350c8596a18dc4d87d96742e60e9eb7 +README.md: d24b1a8f313db0329f9e048e8e1371896484a59e +README.zh.md: d76852f189dc99eabe1c7f99bf7b1f797eb334b6 diff --git a/packages/team/team/README.md b/packages/team/team/README.md index ba8446ec16..d24b1a8f31 100644 --- a/packages/team/team/README.md +++ b/packages/team/team/README.md @@ -41,7 +41,7 @@ The guarantee is process-local retry plus target-Session de-duplication, not cro ## Shared task board -Tasks are complete versioned snapshots. Every mutation carries `expectedRevision`; stale callers receive `TEAM_TASK_STALE_REVISION` instead of overwriting a newer value. Any member can create, read, or claim a ready unowned task. The owner or Lead can edit, release, complete, reopen, or delete it; only the Lead can assign another member. +Tasks are complete versioned snapshots. Every mutation carries `expectedRevision`; stale callers receive `TEAM_TASK_STALE_REVISION` instead of overwriting a newer value. Any member can create, read, or claim a ready unowned task. The owner or Lead can edit, release, complete, reopen, or delete it; only the Lead can assign another member. Numeric `task-` ids require a safe-integer suffix; creation reports `TEAM_TASK_LIMIT` instead of reusing the final safe id. Dependencies must name current non-deleted tasks and form a complete DAG with no self or duplicate edge. A pending task is ready only after every blocker completes. Deleting a task that still has a non-deleted dependent is rejected. Deleted tasks remain tombstones for replay and id stability but do not consume `maxTasks` or appear in `listTasks()`. @@ -49,7 +49,7 @@ Dependencies must name current non-deleted tasks and form a complete DAG with no `waitForChange()` waits for one roster, task, mailbox, or live-status edge that occurs after registration, for 10 seconds through one hour; it reports only whether the wait timed out and does not replay a change that already happened. Runtime disposal releases current waits and makes later waits return immediately without a timeout. Callers re-read authoritative state after wakeup or timeout. Cancellation preserves an Error reason or reports a non-Error reason through `TEAM_WAIT_ABORTED` with structural inspection instead of object coercion. `interrupt()` is Lead-only and delegates to the continuable-subagent interrupt path, which cancels only a live teammate's current turn with `keepInbox`; it neither releases task ownership nor deletes durable mail. -The separately published `./invariant` companion replays each candidate Team event against its committed Session prefix. Replay validates every current-version Team payload before it enters folded state, then rejects invalid member transitions, reused names, discontinuous task revisions, invalid task dependencies, duplicate queue/ack records, and acknowledgements with the wrong target before append. Session event `seq` and `time` own ordering and timing instead of duplicated snapshot timestamps. +The separately published `./invariant` companion replays each candidate Team event against its committed Session prefix. Replay validates every current-version Team payload before it enters folded state, then rejects invalid member transitions, reused names, out-of-range numeric task ids, discontinuous task revisions, invalid task dependencies, duplicate queue/ack records, and acknowledgements with the wrong target before append. Session event `seq` and `time` own ordering and timing instead of duplicated snapshot timestamps. ## Model Experience diff --git a/packages/team/team/README.zh.md b/packages/team/team/README.zh.md index 93a7006ef3..d76852f189 100644 --- a/packages/team/team/README.zh.md +++ b/packages/team/team/README.zh.md @@ -41,7 +41,7 @@ roster 同时报告持久 provisioning/failed phase 与实时 `running`/`idl ## 共享任务板 -任务是完整的版本化快照。每次变更都携带 `expectedRevision`;陈旧调用方会收到 `TEAM_TASK_STALE_REVISION`,不会覆盖更新值。任意成员都可以创建、读取或 claim ready 且无 owner 的任务。Owner 或 Lead 可以编辑、释放、完成、重开或删除任务;只有 Lead 可以分配给其他成员。 +任务是完整的版本化快照。每次变更都携带 `expectedRevision`;陈旧调用方会收到 `TEAM_TASK_STALE_REVISION`,不会覆盖更新值。任意成员都可以创建、读取或 claim ready 且无 owner 的任务。Owner 或 Lead 可以编辑、释放、完成、重开或删除任务;只有 Lead 可以分配给其他成员。数字 `task-` id 的后缀必须是安全整数;最后一个安全 id 已被占用时,创建会报告 `TEAM_TASK_LIMIT`,而不会复用该 id。 依赖必须指向当前未删除任务,并组成完整 DAG,不允许 self edge 或重复 edge。只有所有 blocker 都 completed,pending 任务才 ready。仍被未删除任务依赖的任务不能删除。删除任务作为 tombstone 保留以供回放和维持 id 稳定,但不占用 `maxTasks`,也不出现在 `listTasks()` 中。 @@ -49,7 +49,7 @@ roster 同时报告持久 provisioning/failed phase 与实时 `running`/`idl `waitForChange()` 可以等待注册后发生的下一条 roster、task、mailbox 或实时 status 边,时长范围为 10 秒到 1 小时;它只报告等待是否超时,也不会回放调用前已经发生的变化。运行时 dispose 会释放当前等待,并使后续等待不经超时立即返回。调用方需要在唤醒或超时后重新读取权威状态。取消会保留 Error reason;非 Error reason 则通过 `TEAM_WAIT_ABORTED` 以结构化检查结果报告,不再强制转成 object 字符串。`interrupt()` 仅限 Lead,并委托 continuable-subagent 的 interrupt 路径以 `keepInbox` 只取消 live teammate 的当前 turn;它既不释放任务 owner,也不删除持久 mail。 -单独发布的 `./invariant` 配套模块会把每条候选 Team event 对照已提交 Session 前缀回放。回放会先验证每个当前版本 Team payload,再将其纳入折叠状态;随后会在 append 前拒绝非法 member 转换、名字复用、不连续任务 revision、非法任务依赖、重复 queue/ack,以及 target 不匹配的 acknowledgement。顺序与时间由 Session event 的 `seq` 和 `time` 负责,不在 snapshot 中重复保存。 +单独发布的 `./invariant` 配套模块会把每条候选 Team event 对照已提交 Session 前缀回放。回放会先验证每个当前版本 Team payload,再将其纳入折叠状态;随后会在 append 前拒绝非法 member 转换、名字复用、超出范围的数字 task id、不连续任务 revision、非法任务依赖、重复 queue/ack,以及 target 不匹配的 acknowledgement。顺序与时间由 Session event 的 `seq` 和 `time` 负责,不在 snapshot 中重复保存。 ## 模型体验 diff --git a/packages/team/team/src/fold.ts b/packages/team/team/src/fold.ts index f6e248198b..cf439a5322 100644 --- a/packages/team/team/src/fold.ts +++ b/packages/team/team/src/fold.ts @@ -23,7 +23,11 @@ const nonNegativeSafeInteger = z.number().int().nonnegative().max(Number.MAX_SAF const positiveSafeInteger = nonNegativeSafeInteger.min(1) const sessionIdSchema = z.string().min(1).transform(value => SessionId(value)) const teamIdSchema = z.string().min(1).transform(value => toTeamId(value)) -const teamTaskIdSchema = z.string().min(1).transform(value => toTeamTaskId(value)) +const numericTaskIdPattern = /^task-(\d+)$/u +const teamTaskIdSchema = z.string().min(1).refine((value) => { + const match = numericTaskIdPattern.exec(value) + return match === null || Number.isSafeInteger(Number(match[1])) +}, { message: 'numeric task id suffix must be a safe integer' }).transform(value => toTeamTaskId(value)) const teamMessageIdSchema = z.string().min(1).transform(value => toTeamMessageId(value)) const coreContentBlockTypes = new Set(['text', 'reasoning', 'image', 'tool-call', 'tool-result']) @@ -243,8 +247,14 @@ export function applyTeamEvent(state: TeamFoldState, event: SessionEvent): void throw new Error(`team task "${task.id}" revision is not contiguous`) } assertTaskGraphCandidate(state.tasks, task) - const match = /^task-(\d+)$/u.exec(task.id) - if (match !== null) state.nextTaskNumber = Math.max(state.nextTaskNumber, Number(match[1]) + 1) + const match = numericTaskIdPattern.exec(task.id) + if (match !== null) { + const number = Number(match[1]) + state.nextTaskNumber = Math.max( + state.nextTaskNumber, + number === Number.MAX_SAFE_INTEGER ? number : number + 1, + ) + } state.tasks.set(task.id, task) break } diff --git a/packages/team/team/src/lifecycle.ts b/packages/team/team/src/lifecycle.ts index cf5f1e324f..55bdfadb9e 100644 --- a/packages/team/team/src/lifecycle.ts +++ b/packages/team/team/src/lifecycle.ts @@ -27,6 +27,20 @@ export class TeamRuntimeLifecycle { return reason } + /** Whether a rejection is the runtime cancellation, directly or through an Error cause chain. */ + private isCancellation(reason: unknown): boolean { + const seen = new Set() + let current = reason + while (!seen.has(current)) { + if (this.disposed && current === this.reason) return true + if (this.disposed && current instanceof TeamError && current.code === 'TEAM_DISPOSED') return true + if (!(current instanceof Error)) return false + seen.add(current) + current = current.cause + } + return false + } + /** Close Team runtime admission and cancel admitted interruptible work. */ close(): void { this.controller.abort(new TeamError('Agent Teams service disposed', 'TEAM_DISPOSED')) @@ -42,7 +56,7 @@ export class TeamRuntimeLifecycle { try { const outcomes = await this.withTimeout(Promise.allSettled(operations)) for (const outcome of outcomes) { - if (outcome.status === 'rejected' && outcome.reason !== this.reason) failures.push(outcome.reason) + if (outcome.status === 'rejected' && !this.isCancellation(outcome.reason)) failures.push(outcome.reason) } } catch (error: unknown) { failures.push(error) diff --git a/packages/team/team/src/mailbox.ts b/packages/team/team/src/mailbox.ts index 101a2cad5c..b234cef408 100644 --- a/packages/team/team/src/mailbox.ts +++ b/packages/team/team/src/mailbox.ts @@ -242,7 +242,7 @@ export class TeamMailbox { const input = createUserMessage({ content, source }) if (message.delivery === 'wakeup') { root.followup(input) - return true + return await this.checkpointDelivered(root, root.session, message.id) } root.inject(input) return await this.checkpointDelivered(root, root.session, message.id) diff --git a/packages/team/team/src/task-board.ts b/packages/team/team/src/task-board.ts index c4486968da..f987dd9217 100644 --- a/packages/team/team/src/task-board.ts +++ b/packages/team/team/src/task-board.ts @@ -53,8 +53,12 @@ export class TeamTaskBoard { if (active >= this.maxTasks) { throw new TeamError(`Team task limit ${this.maxTasks} reached`, 'TEAM_TASK_LIMIT') } + const id = TeamTaskId(`task-${state.nextTaskNumber}`) + if (state.tasks.has(id)) { + throw new TeamError('Team task id space exhausted', 'TEAM_TASK_LIMIT') + } const task: TeamTaskSnapshot = { - id: TeamTaskId(`task-${state.nextTaskNumber}`), + id, revision: 1, subject: requiredText(request.subject, 'subject', 200), description: requiredText(request.description, 'description', 16_384), diff --git a/packages/team/team/tests/fold.spec.ts b/packages/team/team/tests/fold.spec.ts index 41161a17cb..d1ed095a99 100644 --- a/packages/team/team/tests/fold.spec.ts +++ b/packages/team/team/tests/fold.spec.ts @@ -202,6 +202,14 @@ describe('Agent Teams fold', () => { expect(state.nextTaskNumber).toBe(1) }) + it('rejects a persisted numeric task id outside the safe integer range', () => { + expect(() => foldTeam(ROOT, [event('team/task', { + version: 1, + teamId: TEAM, + task: task({ id: TeamTaskId('task-9007199254740992') }), + }, 0)])).toThrow(/persisted Agent Teams team\/task payload is invalid/) + }) + it('enforces mailbox queue and acknowledgement relations', () => { const queued = event('team/message/queued', { version: 1, teamId: TEAM, message: message() }, 0) const delivered = event('team/message/delivered', { diff --git a/packages/team/team/tests/team.spec.ts b/packages/team/team/tests/team.spec.ts index 99a959b58f..dd03cd6438 100644 --- a/packages/team/team/tests/team.spec.ts +++ b/packages/team/team/tests/team.spec.ts @@ -14,6 +14,7 @@ import * as SubagentFork from '@deepseek-ai/dsh-subagent-fork-in-process' import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn-in-process' import { MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts' import TeamService, { foldTeam, TeamError, TeamId, TeamMessageId, TeamTaskId } from '../src/index.ts' +import { TeamRuntimeLifecycle } from '../src/lifecycle.ts' import type { TeamMemberSnapshot, TeamMessageSnapshot, TeamTaskSnapshot } from '../src/index.ts' const SIGNAL = new AbortController().signal @@ -509,6 +510,30 @@ describe('Team identity and provisioning', () => { }) describe('Team shared task DAG', () => { + it('fails loudly when the durable numeric task id space is exhausted', async () => { + const { ctx, lead } = await setup([]) + const id = TeamTaskId(`task-${Number.MAX_SAFE_INTEGER}`) + lead.session.append('team/task', { + version: 1, + teamId: TeamId(lead.id), + task: { + id, + revision: 1, + subject: 'last numeric task', + description: 'occupies the final safe numeric task id', + status: 'pending', + blockedBy: [], + writeScopes: [], + }, + }) + await ctx.sessions.flush(lead.session) + + await expect(ctx.teams.createTask(lead, { + subject: 'cannot allocate', + description: 'no safe numeric task id remains', + })).rejects.toMatchObject({ code: 'TEAM_TASK_LIMIT' }) + }) + it('bounds non-deleted tasks while retaining deleted task ids as tombstones', async () => { const { ctx, lead } = await setup([], { maxTasks: 1 }) const first = await ctx.teams.createTask(lead, { subject: 'first', description: 'first task' }) @@ -811,28 +836,54 @@ describe('Team shared task DAG', () => { }) describe('Team mailbox and waiting', () => { - it('delivers quiet and waking teammate messages to the Lead inbox', async () => { - const { ctx, lead } = await setup(['hang', textResponse('lead follow-up answer')]) + it('acknowledges waking messages persisted by a busy Lead before model claim', async () => { + const { ctx, lead, teamFiber } = await setup(['hang', 'hang'], { maxPendingMessagesPerMember: 1 }) const started = await spawn(ctx, lead, 'lead-reporter') const reporter = await waitRunning(ctx, started.member.id) + lead.followup(createUserMessage({ content: content('keep the Lead busy'), source: { kind: 'user' } })) + await waitRunning(ctx, lead.id) - const quiet = await ctx.teams.sendMessage(reporter, { - target: 'lead', content: content('quiet report'), delivery: 'quiet', signal: SIGNAL, + const first = await ctx.teams.sendMessage(reporter, { + target: 'lead', content: content('first wakeup report'), delivery: 'wakeup', signal: SIGNAL, }) - expect(quiet.status).toBe('accepted') - expect(lead.status).toBe('idle') + const second = await ctx.teams.sendMessage(reporter, { + target: 'lead', content: content('second wakeup report'), delivery: 'wakeup', signal: SIGNAL, + }) + expect([first.status, second.status]).toEqual(['accepted', 'accepted']) + expect(lead.status).toBe('running') expect(durable(lead).pendingMessages).toEqual([]) - expect(lead.inbox.nextStep.some(message => message.source.kind === 'team-message' - && message.source.messageId === quiet.messageId)).toBe(true) - const waking = await ctx.teams.sendMessage(reporter, { - target: 'lead', content: content('wake the lead'), delivery: 'wakeup', signal: SIGNAL, - }) - expect(waking.status).toBe('accepted') - await lead.whenIdle() - await vi.waitFor(() => { expect(durable(lead).pendingMessages).toEqual([]) }) - ctx.teams.interrupt(lead, 'lead-reporter') - await waitNoAgent(ctx, reporter.id) + const messageIds = new Set([first.messageId, second.messageId]) + const persisted = await ctx.sessionPersistence.inspect(lead.id) + const receiptOrder = persisted.events.flatMap((event) => { + if (event.type === 'agent/inbox/spliced' && event.data.inserted.some(message => + message.source.kind === 'team-message' && messageIds.has(message.source.messageId))) { + return ['agent/inbox/spliced'] + } + if (event.type === 'team/message/delivered' && messageIds.has(event.data.messageId)) { + return ['team/message/delivered'] + } + return [] + }) + expect(receiptOrder).toEqual([ + 'agent/inbox/spliced', + 'team/message/delivered', + 'agent/inbox/spliced', + 'team/message/delivered', + ]) + + const receiptCount = lead.session.events.filter(event => event.type === 'agent/inbox/spliced' + && event.data.inserted.some(message => message.source.kind === 'team-message' + && messageIds.has(message.source.messageId))).length + await teamFiber.dispose() + await ctx.plugin(TeamService, { maxPendingMessagesPerMember: 1 }) + await vi.waitFor(() => { expect(durable(lead).pendingMessages).toEqual([]) }) + expect(lead.session.events.filter(event => event.type === 'agent/inbox/spliced' + && event.data.inserted.some(message => message.source.kind === 'team-message' + && messageIds.has(message.source.messageId)))).toHaveLength(receiptCount) + + lead.cancel({ kind: 'parent' }) + await lead.whenIdle() }) it('flushes a live pending receipt before acknowledgement without inserting a duplicate', async () => { @@ -1326,6 +1377,28 @@ describe('Team mailbox and waiting', () => { await expect(internal.disposeRuntime()).rejects.toMatchObject({ errors: [cleanupFailure] }) }) + it('recognizes wrapped and coded runtime cancellation during disposal settlement', async () => { + const open = new TeamRuntimeLifecycle(100) + const ordinaryFailure = new Error('ordinary failure before disposal') + const openFailures: unknown[] = [] + await open.settle([Promise.reject(ordinaryFailure)], openFailures) + expect(openFailures).toEqual([ordinaryFailure]) + + const lifecycle = new TeamRuntimeLifecycle(100) + lifecycle.close() + const failures: unknown[] = [] + await lifecycle.settle([ + Promise.reject(new Error('wrapped cancellation', { cause: lifecycle.reason })), + Promise.reject(new TeamError('translated cancellation', 'TEAM_DISPOSED')), + ], failures) + expect(failures).toEqual([]) + + const cyclic = new Error('unrelated cyclic failure') + cyclic.cause = cyclic + await lifecycle.settle([Promise.reject(cyclic)], failures) + expect(failures).toEqual([cyclic]) + }) + it('disposes a live child even after its durable member edge becomes failed', async () => { const { ctx, lead } = await setup(['hang']) const childId = SessionId('failed-live-child') diff --git a/packages/team/tool-team/src/index.ts b/packages/team/tool-team/src/index.ts index 118297dee2..0e35b85337 100644 --- a/packages/team/tool-team/src/index.ts +++ b/packages/team/tool-team/src/index.ts @@ -250,6 +250,8 @@ function install(agent: Agent, ctx: Context, config: Required): () => vo if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 10_000 || timeoutMs > 3_600_000) { return await ctx.teams.waitForChange(caller, timeoutMs, exec.signal) } + // The active-peer read and waiter registration must remain one synchronous + // span; awaiting between them can lose the only peer-status edge. const hasActivePeer = ctx.teams.listMembers(caller).some(member => member.id !== caller.id && ACTIVE_WAIT_STATUSES.has(member.status)) if (!hasActivePeer) { diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index d81740da1a..fd2af2f6ed 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -640,6 +640,9 @@ importers: '@deepseek-ai/dsh-system-prompt': specifier: workspace:* version: link:../packages/core/system-prompt + '@deepseek-ai/dsh-team': + specifier: workspace:* + version: link:../packages/team/team '@deepseek-ai/dsh-terminal': specifier: workspace:* version: link:../packages/terminal/terminal @@ -706,6 +709,9 @@ importers: '@deepseek-ai/dsh-tool-subagent-report': specifier: workspace:* version: link:../packages/subagent/tool-subagent-report + '@deepseek-ai/dsh-tool-team': + specifier: workspace:* + version: link:../packages/team/tool-team '@deepseek-ai/dsh-tool-terminal': specifier: workspace:* version: link:../packages/terminal/tool-terminal