diff --git a/src/daemon/index.ts b/src/daemon/index.ts index 82dcc60..3648623 100644 --- a/src/daemon/index.ts +++ b/src/daemon/index.ts @@ -62,7 +62,7 @@ async function main(): Promise { await app.listen({ host: cfg.host, port: cfg.port }); app.log.info(`maestrod 就绪 · db=${cfg.dbFile} · http://${cfg.host}:${cfg.port} · ws ${cfg.host}:${cfg.port}/ws`); if (rec.blocked || rec.released) app.log.info(`依赖对账:转入等依赖 ${rec.blocked} · 放行可执行 ${rec.released}`); - if (itr.tasks) app.log.info(`中断恢复:${itr.tasks} 个执行中任务重新入队(${itr.runs} 个 run 标记中断)`); + if (itr.readopted || itr.reclaimed) app.log.info(`中断恢复:re-adopt ${itr.readopted} · 回收 ${itr.reclaimed}`); const syncTimer = startSyncLoop(store, app); const orchTimer = startOrchestrator(store, app); // 编排器:自动领取可执行任务(MAESTRO_ORCH_INTERVAL 秒,0=关闭) diff --git a/src/executor/protocol.ts b/src/executor/protocol.ts new file mode 100644 index 0000000..7878297 --- /dev/null +++ b/src/executor/protocol.ts @@ -0,0 +1,162 @@ +// Worker ↔ daemon 契约(Phase 3 多进程执行)。 +// +// 铁律:worker 进程【完全不碰 DB】(连读都不碰)。daemon 是唯一 DB 读写者。 +// 二者只经【文件 + 进程信号】通讯,天然跨 daemon 重启持久: +// daemon → worker:runs//job.json(唯一输入);SIGTERM(取消/超时) +// worker → daemon:runs//outbox.ndjson(追加、带 seq);runs//heartbeat(mtime 心跳) +// +// 本模块同时被 worker 侧(写 outbox/读 job/刷心跳)与 daemon 侧(写 job/读 outbox/判活)import, +// 是唯一的共享面——双方都不应另起自己的协议常量/路径计算。 + +import { + existsSync, mkdirSync, readFileSync, writeFileSync, appendFileSync, + statSync, utimesSync, closeSync, openSync, +} from 'node:fs'; +import { homedir } from 'node:os'; +import { join } from 'node:path'; +import type { Project, Task, ReviewVerdict } from '../model/types.js'; + +/** 数据根目录(与 worktree.ts 同约定): */ +function dataDir(): string { + return process.env.MAESTRO_DATA_DIR ?? join(homedir(), '.maestro'); +} + +/** 全部 run 工作目录根:/runs */ +export function runsBase(): string { + return join(dataDir(), 'runs'); +} +/** 单个 run 的工作目录:/(放 job.json / outbox.ndjson / heartbeat) */ +export function runDir(runId: string): string { + return join(runsBase(), runId); +} +export function jobPath(runId: string): string { return join(runDir(runId), 'job.json'); } +export function outboxPath(runId: string): string { return join(runDir(runId), 'outbox.ndjson'); } +export function heartbeatPath(runId: string): string { return join(runDir(runId), 'heartbeat'); } + +// ───────────────────────── daemon → worker:JobSpec ───────────────────────── + +/** + * worker 的唯一输入。daemon 在 spawn 前写 runs//job.json,worker 读它即可执行, + * 【无需访问 DB】。携带完整 Task / Project(均 JSON 可序列化),worker 自行 pickModel/buildPrompt。 + */ +export interface JobSpec { + runId: string; + task: Task; + project: Project; + /** 确定性 worktree 路径与分支(daemon 预填;worker 也能由 worktree.ts 算出,传入避免重复) */ + worktreeDir: string; + branch: string; +} + +export function writeJobSpec(job: JobSpec): void { + mkdirSync(runDir(job.runId), { recursive: true }); + writeFileSync(jobPath(job.runId), JSON.stringify(job)); +} +export function readJobSpec(runId: string): JobSpec { + return JSON.parse(readFileSync(jobPath(runId), 'utf8')) as JobSpec; +} + +// ───────────────────────── worker → daemon:OutboxRecord ──────────────────── + +/** 一次复审(code review / 安全审计)的产出,worker 报给 daemon,由 daemon 落成 run 行 + result 字段 */ +export interface ReviewReport { + summary: string | null; + verdict: ReviewVerdict | null; + transcriptRef: string | null; +} + +/** + * worker 追加进 outbox 的记录。`seq` 单调递增(每 run 内),`at` ISO 时间,均由 appendOutbox 填。 + * - started:worker 启动,报自己的 pid / worktree / branch / model + * - phase :进度阶段(executing|verifying|reviewing…),daemon 落成 event 推看板(可选) + * - failed :本次尝试终态失败(executor 报错或 verify 失败)。daemon 据此 finishRun(failed)+重试决策 + * - result :成功终态。daemon 据此建 reviewer/security run + setResult(四字段) + 转 exec_review + * - done :worker 即将退出(成功或失败都发,daemon 据此停止 tail) + */ +export type OutboxRecord = + | { seq: number; at: string; type: 'started'; pid: number; worktree: string; branch: string; model: string } + | { seq: number; at: string; type: 'phase'; phase: string } + | { seq: number; at: string; type: 'failed'; error: string; transcriptRef: string | null; sessionId: string | null } + | { + seq: number; at: string; type: 'result'; + branch: string; worktree: string; diffSummary: string; commits: string[]; + executor: { transcriptRef: string | null; sessionId: string | null }; + code: ReviewReport; + security: ReviewReport; + } + | { seq: number; at: string; type: 'done' }; + +/** OutboxRecord 去掉 seq/at(由 appendOutbox 填) */ +export type OutboxPayload = + | Omit, 'seq' | 'at'> + | Omit, 'seq' | 'at'> + | Omit, 'seq' | 'at'> + | Omit, 'seq' | 'at'> + | Omit, 'seq' | 'at'>; + +/** 追加一条 outbox 记录(worker 侧调用)。seq = 现有行数+1(worker 单线程,无并发写)。返回写入的完整记录。 */ +export function appendOutbox(runId: string, payload: OutboxPayload): OutboxRecord { + mkdirSync(runDir(runId), { recursive: true }); + const prev = readOutboxAll(runId); + const seq = prev.length + 1; + const rec = { seq, at: new Date().toISOString(), ...payload } as OutboxRecord; + appendFileSync(outboxPath(runId), JSON.stringify(rec) + '\n'); + return rec; +} + +/** 读全部 outbox 记录(坏行跳过,容忍写一半的尾行)。 */ +export function readOutboxAll(runId: string): OutboxRecord[] { + const p = outboxPath(runId); + if (!existsSync(p)) return []; + const out: OutboxRecord[] = []; + for (const line of readFileSync(p, 'utf8').split('\n')) { + const s = line.trim(); + if (!s) continue; + try { out.push(JSON.parse(s) as OutboxRecord); } catch { /* 半截尾行:跳过,下轮再读 */ } + } + return out; +} + +/** 读 seq > lastSeq 的新记录(daemon ingest 侧调用,幂等去重靠 lastSeq)。 */ +export function readOutboxSince(runId: string, lastSeq: number): OutboxRecord[] { + return readOutboxAll(runId).filter((r) => r.seq > lastSeq); +} + +// ───────────────────────── worker → daemon:心跳 + 存活判定 ────────────────── + +/** worker 刷心跳:touch runs//heartbeat(文件不存在则建)。daemon 用其 mtime 判活。 */ +export function touchHeartbeat(runId: string): void { + mkdirSync(runDir(runId), { recursive: true }); + const p = heartbeatPath(runId); + if (!existsSync(p)) { closeSync(openSync(p, 'w')); return; } + const now = new Date(); + utimesSync(p, now, now); +} + +/** 心跳文件的「年龄」毫秒(now - mtime);无心跳文件返回 null。 */ +export function heartbeatAgeMs(runId: string, now = Date.now()): number | null { + try { return now - statSync(heartbeatPath(runId)).mtimeMs; } + catch { return null; } +} + +/** pid 是否存活(signal 0 探测;EPERM 视为存活——进程在但无权) */ +export function pidAlive(pid: number): boolean { + if (!pid || pid <= 0) return false; + try { process.kill(pid, 0); return true; } + catch (e) { return (e as NodeJS.ErrnoException).code === 'EPERM'; } +} + +export const HEARTBEAT_INTERVAL_MS = 10_000; // worker 刷心跳间隔 +export const HEARTBEAT_GRACE_MS = 60_000; // 心跳超过此年龄视为可疑(须 > interval,容忍 GC/慢盘) +export const BOOT_GRACE_MS = 30_000; // started 后此窗口内额外认 pid 存活(覆盖启动到首次心跳空窗) + +/** + * worker 是否存活。主信号=心跳新鲜;boot 窗口内额外接受 pid 存活(覆盖启动空窗); + * 窗口后仅认心跳——规避 pid 复用误判(被复用的 pid 不会更新本 run 的心跳)。 + */ +export function isWorkerAlive(args: { pid: number | null; heartbeatAgeMs: number | null; startedAgeMs: number }): boolean { + const { pid, heartbeatAgeMs: hb, startedAgeMs } = args; + if (hb !== null && hb < HEARTBEAT_GRACE_MS) return true; // 心跳新鲜 → 活 + if (startedAgeMs < BOOT_GRACE_MS && pid !== null && pidAlive(pid)) return true; // 启动空窗 → pid 兜底 + return false; +} diff --git a/src/model/types.ts b/src/model/types.ts index 884c1ed..20a694f 100644 --- a/src/model/types.ts +++ b/src/model/types.ts @@ -70,6 +70,7 @@ export interface Task { result: TaskResult | null; assignee: 'agent' | 'human' | null; retryBaseline: number; // 上次手动重投时已有的失败 run 数(重置重试计数用) + nextEligibleAt: string | null; // 持久化退避:早于此时间不被领取(重试退避,重启不丢);null=即刻可领 createdAt: string; updatedAt: string; } @@ -89,6 +90,8 @@ export interface Run { transcriptRef: string | null; // agent 转录日志文件路径 claudeSessionId: string | null; error: string | null; + workerPid: number | null; // 执行该 run 的 worker 进程 pid(多进程执行;daemon 写,判活用) + lastSeq: number; // 已 ingest 的 outbox 最大 seq(daemon 写,幂等游标;重启续读) } export type EventType = diff --git a/src/store/db.ts b/src/store/db.ts index 37b5aeb..3d8bf9a 100644 --- a/src/store/db.ts +++ b/src/store/db.ts @@ -29,6 +29,9 @@ export function openDb(file: string): Database.Database { ensureColumn(db, 'projects', 'max_retries', 'max_retries INTEGER NOT NULL DEFAULT 2'); // 失败后最大自动重试次数 ensureColumn(db, 'projects', 'timeout_ms', 'timeout_ms INTEGER NOT NULL DEFAULT 1800000'); // 单次执行超时毫秒 ensureColumn(db, 'tasks', 'retry_baseline', 'retry_baseline INTEGER NOT NULL DEFAULT 0'); // 手动重投时的失败基线 + ensureColumn(db, 'tasks', 'next_eligible_at', 'next_eligible_at TEXT'); // 持久化退避(多进程执行) + ensureColumn(db, 'runs', 'worker_pid', 'worker_pid INTEGER'); // worker 进程 pid + ensureColumn(db, 'runs', 'last_seq', 'last_seq INTEGER NOT NULL DEFAULT 0'); // outbox ingest 游标 const schema = readFileSync(join(HERE, 'schema.sql'), 'utf8'); db.exec(schema); return db; diff --git a/src/store/mappers.ts b/src/store/mappers.ts index 13c10f5..81d5a9f 100644 --- a/src/store/mappers.ts +++ b/src/store/mappers.ts @@ -16,6 +16,7 @@ export interface TaskRow { title: string; complexity: string; status: string; priority: number; deps: string; plan: string | null; spec: string | null; operations: string | null; result: string | null; assignee: string | null; retry_baseline: number; + next_eligible_at: string | null; created_at: string; updated_at: string; source_ref: string | null; } @@ -27,6 +28,7 @@ export interface RunRow { id: string; task_id: string; kind: string; worktree: string | null; branch: string | null; status: string; started_at: string; ended_at: string | null; transcript_ref: string | null; claude_session_id: string | null; error: string | null; + worker_pid: number | null; last_seq: number; } export interface EventRow { id: string; project_id: string; task_id: string | null; type: string; payload: string; at: string; @@ -68,6 +70,7 @@ export function rowToTask(r: TaskRow, approvals: ApprovalRecord[] = []): Task { plan: r.plan, spec: r.spec, operations: r.operations, approvals, result: r.result ? parseResult(r.result) : null, assignee: r.assignee as Task['assignee'], retryBaseline: r.retry_baseline ?? 0, + nextEligibleAt: r.next_eligible_at ?? null, createdAt: r.created_at, updatedAt: r.updated_at, }; } @@ -85,6 +88,7 @@ export function rowToRun(r: RunRow): Run { worktree: r.worktree, branch: r.branch, status: r.status as RunStatus, startedAt: r.started_at, endedAt: r.ended_at, transcriptRef: r.transcript_ref, claudeSessionId: r.claude_session_id, error: r.error, + workerPid: r.worker_pid ?? null, lastSeq: r.last_seq ?? 0, }; } diff --git a/src/store/schema.sql b/src/store/schema.sql index 0965981..8fff1c7 100644 --- a/src/store/schema.sql +++ b/src/store/schema.sql @@ -36,6 +36,7 @@ CREATE TABLE IF NOT EXISTS tasks ( result TEXT, -- JSON TaskResult assignee TEXT, -- agent | human retry_baseline INTEGER NOT NULL DEFAULT 0, -- 上次手动重投时已有的失败 run 数(重置重试计数用) + next_eligible_at TEXT, -- 持久化退避:早于此时间不被领取(null=即刻可领) created_at TEXT NOT NULL, updated_at TEXT NOT NULL, source_ref TEXT -- 旧 todo 来源标识(todo:17 / todo:17/1A),项目内唯一 @@ -66,7 +67,9 @@ CREATE TABLE IF NOT EXISTS runs ( ended_at TEXT, transcript_ref TEXT, claude_session_id TEXT, - error TEXT + error TEXT, + worker_pid INTEGER, -- 多进程执行:worker 进程 pid(daemon 写,判活用) + last_seq INTEGER NOT NULL DEFAULT 0 -- 已 ingest 的 outbox 最大 seq(幂等游标) ); CREATE INDEX IF NOT EXISTS idx_runs_task ON runs(task_id, started_at); diff --git a/src/store/store.ts b/src/store/store.ts index ed9d1b5..eca79b7 100644 --- a/src/store/store.ts +++ b/src/store/store.ts @@ -231,8 +231,8 @@ export class Store { id: id('tsk'), project_id: input.projectId, parent_id: input.parentId ?? null, depth, title: input.title, complexity: input.complexity, status, priority: input.priority ?? 1, deps: JSON.stringify(input.deps ?? []), plan: null, spec: null, operations: null, - result: null, assignee: null, retry_baseline: 0, created_at: now(), updated_at: now(), - source_ref: null, + result: null, assignee: null, retry_baseline: 0, next_eligible_at: null, + created_at: now(), updated_at: now(), source_ref: null, }; this.db.prepare( `INSERT INTO tasks (id,project_id,parent_id,depth,title,complexity,status,priority,deps,plan,spec,operations,result,assignee,retry_baseline,created_at,updated_at,source_ref) @@ -562,15 +562,86 @@ export class Store { * 卡在 executing 的任务转 failed→queued 等编排器重新领取。 * 注:中断产生的 failed run 会计入该任务的失败次数(多次中断+真失败可能提前转 needs_attention,可接受)。 */ - reconcileInterrupted(): { runs: number; tasks: number } { - const runs = this.db.prepare(`SELECT * FROM runs WHERE status = 'started'`).all() as RunRow[]; - for (const r of runs) this.finishRun(r.id, 'failed', { error: 'daemon 重启,执行中断' }); - const rows = this.db.prepare(`SELECT * FROM tasks WHERE status = 'executing'`).all() as TaskRow[]; - for (const t of rows) { - this.transition(t.id, 'failed', { auto: 'daemon-restart' }); - this.transition(t.id, 'queued', { auto: 'daemon-restart' }); + /** + * 失败/重试策略的唯一落点(ingest 的 failed 记录、reaper 判死、reconcile 判死 都调它)。 + * 把(可选)执行 run 收尾为 failed,按「净失败次数 vs project.maxRetries」决定: + * 未到上限 → queued + 持久化退避(next_eligible_at = now + 指数退避) + * 到上限 → needs_attention(清退避) + * netFailed 相对 retry_baseline(手动重投时设的基线)。退避:min(30s·2^(n-1), 10min)。 + */ + failTaskAttempt(taskId: string, runId: string | null, error: string): Task { + const row = this.getTaskRow(taskId); + if (!row) throw new StoreError(`任务不存在: ${taskId}`); + const project = this.getProject(row.project_id); + if (!project) throw new StoreError(`项目不存在: ${row.project_id}`); + + // 收尾执行 run:给了 runId 且仍 started → 标 failed;没给 → 补记一条 failed(保证重试计数不漏) + if (runId) { + const rr = this.db.prepare(`SELECT status FROM runs WHERE id = ?`).get(runId) as { status: string } | undefined; + if (rr && rr.status === 'started') this.finishRun(runId, 'failed', { error }); + } else { + const r = this.startRun(taskId, 'executor'); + this.finishRun(r.id, 'failed', { error }); } - return { runs: runs.length, tasks: rows.length }; + + if ((this.getTaskRow(taskId)!.status as TaskStatus) !== 'failed') { + this.transition(taskId, 'failed', { by: 'failTaskAttempt', error }); + } + + const allFailed = (this.db.prepare( + `SELECT COUNT(*) AS n FROM runs WHERE task_id = ? AND kind = 'executor' AND status = 'failed'`, + ).get(taskId) as { n: number }).n; + const netFailed = Math.max(0, allFailed - row.retry_baseline); + const priorNetFailed = Math.max(0, netFailed - 1); // 不含本次 + if (priorNetFailed < project.maxRetries) { + const attempt = priorNetFailed + 1; + const backoffMs = Math.min(30_000 * 2 ** (attempt - 1), 10 * 60_000); + this.setNextEligibleAt(taskId, new Date(Date.now() + backoffMs).toISOString()); + this.transition(taskId, 'queued', { by: 'failTaskAttempt', retry: attempt, backoffMs }); + } else { + this.setNextEligibleAt(taskId, null); + this.transition(taskId, 'needs_attention', { by: 'failTaskAttempt', allFailed, netFailed }); + } + return this.getTask(taskId)!; + } + + /** 持久化退避:早于 next_eligible_at 不被领取(null=即刻可领)。 */ + setNextEligibleAt(taskId: string, iso: string | null): void { + this.db.prepare(`UPDATE tasks SET next_eligible_at = ?, updated_at = ? WHERE id = ?`).run(iso, now(), taskId); + } + + /** 项目内 executing 任务数(多进程执行的并发闸:每个 executing 对应一个活 worker)。 */ + countExecuting(projectId: string): number { + return (this.db.prepare( + `SELECT COUNT(*) AS n FROM tasks WHERE project_id = ? AND status = 'executing'`, + ).get(projectId) as { n: number }).n; + } + + /** 所有 executing 任务 + 其最近一条 executor run(reaper / reconcile 判活用)。run 可能为 null(异常)。 */ + executingWithLatestExecutorRun(): Array<{ task: Task; run: Run | null }> { + const rows = this.db.prepare(`SELECT * FROM tasks WHERE status = 'executing'`).all() as TaskRow[]; + return rows.map((t) => { + const rr = this.db.prepare( + `SELECT * FROM runs WHERE task_id = ? AND kind = 'executor' ORDER BY started_at DESC LIMIT 1`, + ).get(t.id) as RunRow | undefined; + return { task: rowToTask(t), run: rr ? rowToRun(rr) : null }; + }); + } + + /** + * 启动中断恢复(多进程执行):逐个 executing 任务按注入的 isAlive 判活—— + * 活 → re-adopt(保留 executing,daemon 续 ingest 其 outbox); + * 死 → failTaskAttempt(收尾 + 重试/needs_attention)。 + * isAlive 由 daemon 注入(protocol.isWorkerAlive,基于 worker_pid + 心跳);默认保守判死(无判活信息时回收)。 + */ + reconcileInterrupted(isAlive: (run: Run | null) => boolean = () => false): { readopted: number; reclaimed: number } { + let readopted = 0, reclaimed = 0; + for (const { task, run } of this.executingWithLatestExecutorRun()) { + if (run && isAlive(run)) { readopted++; continue; } + this.failTaskAttempt(task.id, run?.id ?? null, 'daemon 重启时发现 worker 已退出'); + reclaimed++; + } + return { readopted, reclaimed }; } /** @@ -687,6 +758,7 @@ export class Store { const rr: RunRow = { id: id('run'), task_id: taskId, kind, worktree: fields.worktree ?? null, branch: fields.branch ?? null, status: 'started', started_at: now(), ended_at: null, transcript_ref: null, claude_session_id: null, error: null, + worker_pid: null, last_seq: 0, }; this.db.prepare( `INSERT INTO runs (id,task_id,kind,worktree,branch,status,started_at,ended_at,transcript_ref,claude_session_id,error) @@ -707,6 +779,16 @@ export class Store { return rowToRun(this.db.prepare(`SELECT * FROM runs WHERE id = ?`).get(runId) as RunRow); } + /** 记录执行该 run 的 worker 进程 pid(多进程执行,daemon spawn 后写;判活用)。 */ + setWorkerPid(runId: string, pid: number): void { + this.db.prepare(`UPDATE runs SET worker_pid = ? WHERE id = ?`).run(pid, runId); + } + + /** 更新已 ingest 的 outbox 最大 seq(幂等游标;daemon ingest 后写,重启续读)。 */ + setLastSeq(runId: string, seq: number): void { + this.db.prepare(`UPDATE runs SET last_seq = ? WHERE id = ?`).run(seq, runId); + } + /** 所有进行中的 run(status='started'),联 tasks 取任务标题与项目。 */ activeRuns(): ActiveRun[] { const rows = this.db.prepare(