import { Store } from '../store/index.js'; import { readOutboxSince, type OutboxRecord } from '../executor/protocol.js'; /** * Outbox → DB 摄取(daemon 是唯一 DB 写者)。 * * worker 全程不碰 DB,只把进度/结果追加进 runs//outbox.ndjson。daemon 在每轮 tick 里把 * 这些记录读出来,按 seq 升序逐条映射成 DB 写(finishRun / setResult / transition / failTaskAttempt …)。 * * 幂等靠 run.lastSeq 游标:只处理 seq>lastSeq 的记录,每处理一条就 setLastSeq(seq)。崩溃重启续读时 * 已落库的 seq 不会重做——故 result/failed/done 这类终态记录即便文件还在,也只会被消费一次。 */ export interface IngestLogger { info(msg: string): void; error(msg: string): void; } /** * 摄取单个 run 的新 outbox 记录(seq>run.lastSeq),按序落库。 * run 不存在则忽略(已被删/异常)。每条处理后推进 lastSeq 游标。 */ export function ingestRun(store: Store, log: IngestLogger, runId: string): void { const run = store.getRun(runId); if (!run) return; const taskId = run.taskId; const recs = readOutboxSince(runId, run.lastSeq); for (const rec of recs) { try { applyRecord(store, log, taskId, runId, run.kind, rec); } catch (e) { // 单条映射失败不阻断后续 run 的摄取;但本条不推进游标,下轮重试。 log.error(`ingest run=${runId} seq=${rec.seq} 失败:${(e as Error).message}`); return; } store.setLastSeq(runId, rec.seq); } } /** 把一条 OutboxRecord 映射为 DB 写。runKind 用于区分 failed 的收尾策略(executor vs planner)。 */ function applyRecord(store: Store, log: IngestLogger, taskId: string, runId: string, runKind: string, rec: OutboxRecord): void { switch (rec.type) { case 'started': // worker 启动自报;pid 已由 daemon spawn 时 setWorkerPid 写过,这里仅记日志。 log.info(`worker 启动 task=${taskId} run=${runId} pid=${rec.pid} model=${rec.model}`); return; case 'phase': // 进度阶段,仅日志(看板事件可后续接)。 log.info(`task=${taskId} run=${runId} 阶段=${rec.phase}`); return; case 'failed': { // 先把 run 收尾为 failed(带转录/会话),再走失败/重试策略(按 kind 分流)。 // failTaskAttempt/failPlanAttempt 见 run 已 ended(非 started)不会重复收尾,只做重试决策。 try { store.finishRun(runId, 'failed', { error: rec.error, transcriptRef: rec.transcriptRef ?? undefined, claudeSessionId: rec.sessionId ?? undefined, }); } catch (e) { log.error(`task=${taskId} run=${runId} 收尾 failed 出错(继续重试决策):${(e as Error).message}`); } if (runKind === 'planner') store.failPlanAttempt(taskId, runId, rec.error); // planner:退避留 analyzing/speccing else store.failTaskAttempt(taskId, runId, rec.error); // executor:failed→重试/needs_attention log.info(`task=${taskId} run=${runId}(${runKind}) 失败:${rec.error}`); return; } case 'spec-result': { // planner-spec 成功:写方案 → 转 spec_review(待你审)→ 收尾 planner run。 store.setSpec(taskId, rec.spec); store.transition(taskId, 'spec_review', { by: 'ingest', runId }); store.finishRun(runId, 'succeeded', { transcriptRef: rec.transcriptRef ?? undefined, claudeSessionId: rec.sessionId ?? undefined, }); log.info(`task=${taskId} run=${runId} 方案完成 → spec_review`); return; } case 'decompose-result': { // planner-decompose 成功:写分析 + 建子任务 + 转 plan_review(待你审)→ 收尾 planner run。 const task = store.getTask(taskId); if (!task) { log.error(`decompose: 任务不存在 ${taskId}`); return; } store.setPlan(taskId, rec.plan || '(无分析说明)'); let created = 0; const createdIds: string[] = []; for (const sub of rec.subtasks) { try { const priority = typeof (sub as { priority?: unknown }).priority === 'number' ? (sub as { priority: number }).priority : 1; const newTask = store.createTask({ projectId: task.projectId, parentId: taskId, title: sub.title, complexity: sub.complexity, priority }); createdIds.push(newTask.id); created++; } catch (e) { log.error(`decompose: 建子任务「${sub.title}」失败:${(e as Error).message}`); createdIds.push(''); // 占位,保持序号对齐 } } // 序号 → taskId 映射,写入 deps for (let i = 0; i < rec.subtasks.length; i++) { const sub = rec.subtasks[i] as { deps?: number[] }; if (!sub.deps?.length || !createdIds[i]) continue; const depIds = sub.deps .filter((idx) => idx >= 0 && idx < createdIds.length && createdIds[idx]) .map((idx) => createdIds[idx]); if (depIds.length > 0) { try { store.patchTask(createdIds[i], { deps: depIds }); } catch (e) { log.error(`decompose: 写子任务「${rec.subtasks[i].title}」deps 失败:${(e as Error).message}`); } } } store.transition(taskId, 'plan_review', { by: 'ingest', runId, subtasks: created }); store.finishRun(runId, 'succeeded', { transcriptRef: rec.transcriptRef ?? undefined, claudeSessionId: rec.sessionId ?? undefined, }); log.info(`task=${taskId} run=${runId} 拆解完成 → plan_review(${created} 个子任务)`); return; } case 'result': { // 成功终态:建 reviewer/security 两条复审 run + 收尾,写 result 四字段,转 exec_review,收尾 executor run。 const cr = store.startRun(taskId, 'reviewer', { worktree: rec.worktree, branch: rec.branch }); store.finishRun(cr.id, 'succeeded', { transcriptRef: rec.code.transcriptRef ?? undefined }); const sr = store.startRun(taskId, 'security', { worktree: rec.worktree, branch: rec.branch }); store.finishRun(sr.id, 'succeeded', { transcriptRef: rec.security.transcriptRef ?? undefined }); store.setResult(taskId, { branch: rec.branch, worktree: rec.worktree, diffSummary: rec.diffSummary, commits: rec.commits, prUrl: null, summary: rec.code.summary, verdict: rec.code.verdict, securitySummary: rec.security.summary, securityVerdict: rec.security.verdict, mergeTaskId: null, }); // verdict 硬闸:任一 reject → 退回重执行,不进 exec_review const codeReject = rec.code.verdict === 'reject'; const secReject = rec.security.verdict === 'reject'; if (codeReject || secReject) { const reason = [ codeReject ? `代码复审拒绝:${rec.code.summary?.slice(0, 200)}` : '', secReject ? `安全审计拒绝:${rec.security.summary?.slice(0, 200)}` : '', ].filter(Boolean).join(';'); store.finishRun(runId, 'succeeded', { transcriptRef: rec.executor.transcriptRef ?? undefined, claudeSessionId: rec.executor.sessionId ?? undefined, }); // 退回重执行(带复审意见) store.failTaskAttempt(taskId, null, `自动复审拒绝,退回重执行:${reason}`); log.info(`task=${taskId} run=${runId} 复审拒绝 → 退回重执行`); return; } // 双 approve → 正常流程 store.transition(taskId, 'exec_review', { by: 'ingest', runId }); store.finishRun(runId, 'succeeded', { transcriptRef: rec.executor.transcriptRef ?? undefined, claudeSessionId: rec.executor.sessionId ?? undefined, }); log.info(`task=${taskId} run=${runId} 执行完成 → exec_review(${rec.commits.length} commits)`); return; } case 'done': // worker 即将退出:终态标记。run 已由 result/failed 收尾,这里无须额外 DB 写。 log.info(`task=${taskId} run=${runId} worker 退出`); return; } } /** * 摄取所有在途 run(executor + planner):遍历所有 started run 调 ingestRun。 * daemon 每轮 tick 调一次,把 worker 期间累积的 outbox 落库。 */ export function ingestAll(store: Store, log: IngestLogger): void { for (const { run } of store.liveRunsWithTask()) { ingestRun(store, log, run.id); } }