fix: 主线 A 事件落库地基 + 主线 B1/B2 执行原语
主线 A(事件驱动落库): - 统一 EventStore 模块单例:RuntimeApp 不再 new EventStore,改用 eventStore 并 setRepositories(14 个 domain repo),消除事件流向空 DB 的割裂 - Scheduler.create_tasks 改为 async,真正发出 task.created 事件 - run.ts dispatchTask 加 await - 主线 A 独立复审发现并修复关键假绿:四个 repo(Task/Agent/ToolRun/ TaskAttempt)的 *Update 类型 Omit<'status'> 且 update() 主动丢弃 status, 导致 EventStore.project() 的状态写入全部静默失效,DB 行内容 tasks.status 永远冻结在 pending,UI 显示的 completed 来自内存 graph。已修,DB 现 真实反映 task.status=completed - 补 agent.started/agent.completed/agent.failed 事件发出(之前 agents 表 恒空),修复后 agents 表有正确行+status 主线 B1(结构化工具调用块类型,N1): - 新增 content-block.ts 定义 Anthropic canonical content blocks (TextBlock/ThinkingBlock/ToolUseBlock/ToolResultBlock/CanonicalMessage) - provider.ts ProviderCompletionInput 去掉 unknown 逃生舱: messages: CanonicalMessage[], tools?: ToolDefinitionBlock[], tool_choice?: ToolChoice, system?: string | TextBlock[] 主线 B2(read-before-edit 代码层强制,FR-009): - fs/index.ts 新增 readFileState 机制(移植 claude-code FileEditTool), fs.edit 执行前检查:未读先改报 "File has not been read yet",外部修改 报 "File has been unexpectedly modified" - 修复 fs.edit 参数名不匹配:兼容 old_str/new_str (ExecutorRole) 和 find/replace (UI) 两种命名 - fs_edit 唯一性检查(非 global 模式下 old_str 出现多次报错) 真实验收: - TSC=0 - air run 后 DB:events=5(原 3,+agent.started/completed), tasks.status=completed(原 frozen pending),agents 1 行 status=completed - read-before-edit 行为测试:未读先改 status=error,读后再改 status=ok Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -62,7 +62,7 @@ export class Scheduler {
|
||||
/**
|
||||
* Create tasks from specifications.
|
||||
*/
|
||||
create_tasks(tasks: Array<{ id: TaskID; type: string; title: string; description?: string; depends_on?: string[] }>): void {
|
||||
async create_tasks(tasks: Array<{ id: TaskID; type: string; title: string; description?: string; depends_on?: string[] }>): Promise<void> {
|
||||
for (const task of tasks) {
|
||||
this.graph.add_task({
|
||||
id: task.id,
|
||||
@@ -71,9 +71,28 @@ export class Scheduler {
|
||||
description: task.description,
|
||||
dependencies: task.depends_on?.map(d => ({ task_id: d, type: 'hard' as const })) || []
|
||||
})
|
||||
|
||||
// Emit task.created events (INV-1: via event store for projection)
|
||||
await eventIngestor.ingest({
|
||||
id: `evt_${task.id}_created`,
|
||||
type: 'task.created',
|
||||
version: 1,
|
||||
session_id: this.context.session_id,
|
||||
project_id: this.context.project_id,
|
||||
timestamp: new Date().toISOString(),
|
||||
source: { kind: 'scheduler' },
|
||||
route: ['scheduler', 'create'],
|
||||
payload: {
|
||||
task_id: task.id,
|
||||
type: task.type,
|
||||
title: task.title,
|
||||
task_spec_json: { description: task.description || '' },
|
||||
dependencies: (task.depends_on || []).map(d => ({ depends_on_task_id: d, dependency_type: 'hard', reason: '' })),
|
||||
metadata: {},
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// Emit task.created events (INV-1: via projection, not direct status write)
|
||||
this.state = 'PLANNING_WAVE'
|
||||
}
|
||||
|
||||
@@ -175,6 +194,28 @@ export class Scheduler {
|
||||
})
|
||||
this.graph.update_status(task.id, 'running')
|
||||
this.agent_monitor.record_heartbeat(agent_id, task.id)
|
||||
|
||||
// INV-1: Emit agent.started event (durable) for projection
|
||||
await eventIngestor.ingest({
|
||||
id: `evt_${agent_id}_started`,
|
||||
type: 'agent.started',
|
||||
version: 1,
|
||||
session_id: this.context.session_id,
|
||||
project_id: this.context.project_id,
|
||||
timestamp: new Date().toISOString(),
|
||||
source: { kind: 'scheduler' },
|
||||
route: ['scheduler', 'dispatch'],
|
||||
payload: {
|
||||
agent_id,
|
||||
agent_type: task.type || 'executor',
|
||||
task_id: task.id,
|
||||
pid: 0,
|
||||
model_provider_id: '',
|
||||
model_id: '',
|
||||
workspace_id: `ws_${task.id}`,
|
||||
metadata: {},
|
||||
}
|
||||
})
|
||||
} catch {
|
||||
// INV-1: emit task.failed event for projection
|
||||
await eventIngestor.ingest({
|
||||
@@ -278,6 +319,18 @@ export class Scheduler {
|
||||
}
|
||||
})
|
||||
this.graph.update_status(task.id, 'completed')
|
||||
// INV-1: Emit agent.completed event (durable) for projection
|
||||
await eventIngestor.ingest({
|
||||
id: `evt_${handle.worker_id}_completed`,
|
||||
type: 'agent.completed',
|
||||
version: 1,
|
||||
session_id: this.context.session_id,
|
||||
project_id: this.context.project_id,
|
||||
timestamp: new Date().toISOString(),
|
||||
source: { kind: 'scheduler' },
|
||||
route: ['scheduler', 'monitoring'],
|
||||
payload: { agent_id: handle.worker_id, task_id: task.id, summary: result.summary, worker_result_ref: attempt_id, metadata: {} }
|
||||
})
|
||||
this.agent_monitor.remove(handle.worker_id)
|
||||
} else if (result.status === 'blocked') {
|
||||
await eventIngestor.ingest({
|
||||
@@ -320,6 +373,18 @@ export class Scheduler {
|
||||
payload: { task_id: task.id, agent_id: handle.worker_id, attempt_id, error: { message: result.summary }, evidence_refs: result.evidence_refs, metadata: { worker_status: result.status } }
|
||||
})
|
||||
this.graph.update_status(task.id, 'failed')
|
||||
// INV-1: Emit agent.failed event (durable) for projection
|
||||
await eventIngestor.ingest({
|
||||
id: `evt_${handle.worker_id}_failed`,
|
||||
type: 'agent.failed',
|
||||
version: 1,
|
||||
session_id: this.context.session_id,
|
||||
project_id: this.context.project_id,
|
||||
timestamp: new Date().toISOString(),
|
||||
source: { kind: 'scheduler' },
|
||||
route: ['scheduler', 'monitoring'],
|
||||
payload: { agent_id: handle.worker_id, task_id: task.id, error: { message: result.summary }, evidence_refs: result.evidence_refs || [], metadata: {} }
|
||||
})
|
||||
this.agent_monitor.remove(handle.worker_id)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user