fix(team): address runtime review feedback
This commit is contained in:
parent
3546f595b9
commit
0283d62c76
18 changed files with 463 additions and 34 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
36
examples/headless-agent/team.cordis.snapshot.yml
Normal file
36
examples/headless-agent/team.cordis.snapshot.yml
Normal file
|
|
@ -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'
|
||||
194
examples/headless-agent/tests/fixtures/team-llm.mjs
vendored
Normal file
194
examples/headless-agent/tests/fixtures/team-llm.mjs
vendored
Normal file
|
|
@ -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())
|
||||
}
|
||||
|
|
@ -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')
|
||||
|
|
|
|||
|
|
@ -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:*",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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-<n>` 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
|
||||
|
||||
|
|
|
|||
|
|
@ -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-<n>` 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 中重复保存。
|
||||
|
||||
## 模型体验
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<unknown>()
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -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', {
|
||||
|
|
|
|||
|
|
@ -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')
|
||||
|
|
|
|||
|
|
@ -250,6 +250,8 @@ function install(agent: Agent, ctx: Context, config: Required<Config>): () => 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) {
|
||||
|
|
|
|||
6
pnpm-lock.yaml
generated
6
pnpm-lock.yaml
generated
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue