import { spawn } from 'node:child_process'; import { fileURLToPath } from 'node:url'; import { dirname, join } from 'node:path'; import { Store } from '../store/index.js'; import type { Project, Task } from '../model/types.js'; import { rankByScore } from '../model/scoring.js'; import { branchFor, worktreeDirFor } from '../executor/worktree.js'; import { writeJobSpec, isWorkerAlive, heartbeatAgeMs, type JobSpec } from '../executor/protocol.js'; import { pickModel } from '../executor/models.js'; import { prepareSandboxedSpawn, describeSandbox } from '../executor/sandbox.js'; import { ingestAll, type IngestLogger } from './ingest.js'; /** 失败后默认最多自动重试次数(项目级 maxRetries 未设置时回退此值)。失败/退避策略本体在 store.failTaskAttempt。 */ export const DEFAULT_MAX_RETRIES = 2; export interface OrchestratorLogger { info(msg: string): void; error(msg: string): void; } /** * worker 入口解析:返回 spawn 用的 [cmd, ...args](不含 runId,由调用方追加)。 * - 默认:node <本文件同级 ../executor/worker.js>(编译后 dist 布局:daemon/ 与 executor/ 同级)。 * - env MAESTRO_WORKER_CMD 覆盖整条命令(空格分隔),供测试 / tsx 跑 .ts 入口用, * 例:MAESTRO_WORKER_CMD="npx tsx src/executor/worker.ts"。 */ export function workerEntry(): string[] { const override = process.env.MAESTRO_WORKER_CMD?.trim(); if (override) return override.split(/\s+/); const here = dirname(fileURLToPath(import.meta.url)); return ['node', join(here, '..', 'executor', 'worker.js')]; } /** * spawn worker 时传给子进程的 env:透传完整会话/系统环境,仅【剥离敏感密钥类变量】。 * 为什么不用严格 allowlist:本机鉴权走 macOS keychain(无 API_KEY/凭证文件),而 Security/keychain * 访问依赖一批会话变量(USER/LOGNAME/__CF_USER_TEXT_ENCODING/TMPDIR/XPC_* 等,均非密钥);只放 * PATH/HOME/LANG 会让 keychain 读不到 → "Not logged in"。故改为黑名单:会话变量照常透传、保 keychain 可用, * 同时剥掉 *_SECRET/*_TOKEN/*_KEY/AWS_/CF_ 等真正的密钥(鉴权用的 ANTHROPIC_/CLAUDE_ 始终保留)。 */ const KEEP_AUTH = (k: string): boolean => k.startsWith('ANTHROPIC_') || k.startsWith('CLAUDE_') || k === 'CLAUDECODE'; const SECRET_KEY_RE = /SECRET|PASSWORD|PASSWD|PRIVATE_KEY|CREDENTIAL|_TOKEN|API[_-]?KEY|ACCESS_KEY|^AWS_|^CF_|^GITHUB|^GH_|^NPM_/i; export function workerEnv(src: NodeJS.ProcessEnv = process.env): NodeJS.ProcessEnv { const out: NodeJS.ProcessEnv = {}; for (const [k, v] of Object.entries(src)) { if (v === undefined) continue; if (KEEP_AUTH(k)) { out[k] = v; continue; } // 鉴权变量始终保留(即便名含 TOKEN/KEY) if (SECRET_KEY_RE.test(k)) continue; // 真正的密钥类:剥离 out[k] = v; // 其余(会话/系统变量)透传 → keychain 可用 } return out; } /** * 默认 spawnWorker:unref + 脱离 stdio + 收敛 env,返回子进程 pid。 * 关键:【不 detached(不 setsid)】——保持与 daemon 同会话,worker 内起的 `claude` 才能沿用 * daemon 的鉴权上下文(本机无 API_KEY/凭证文件,鉴权走 macOS keychain;detached 后新会话读不到 → Not logged in)。 * 存活性不依赖 detached:unref 后 daemon 退出,worker 作为孤儿被 init/launchd 收养、继续运行(daemon 无控制终端, * 不会有进程组 SIGHUP;重启用单 pid kill 不波及 worker)。 * * OS 级沙箱(可选,默认关):MAESTRO_SANDBOX 开启时,prepareSandboxedSpawn 把命令包成 * `sh -c 'ulimit…; exec sandbox-exec -f profile <原命令>'`——文件写围栏 + 资源上限。 * 末尾 exec 链保证 pid 不变(仍是真 worker 的 pid),unref/同会话/keychain 等约束不受影响。 */ function defaultSpawnWorker(runId: string): number { const base = [...workerEntry(), runId]; const wrapped = prepareSandboxedSpawn(base, runId); // 沙箱关闭 → null,按原命令直接 spawn const [cmd, ...args] = wrapped ? [wrapped.cmd, ...wrapped.args] : base; const child = spawn(cmd, args, { detached: false, stdio: 'ignore', env: workerEnv(), }); child.unref(); if (child.pid === undefined) throw new Error('spawn worker 未返回 pid'); return child.pid; } /** 依赖注入点:生产用默认实现,测试传 mock(不真 spawn / 不真判活 / ingest no-op)。 */ export interface OrchestratorDeps { /** spawn 一个独立 worker 进程跑该 run,返回其 pid。 */ spawnWorker: (runId: string) => number; /** worker 是否存活(reaper 判死用),默认 protocol.isWorkerAlive。 */ isWorkerAlive: typeof isWorkerAlive; /** 把所有在途 run 的 outbox 摄取入库(daemon 唯一写者),默认 ingest.ingestAll。 */ ingestAll: (store: Store, log: IngestLogger) => void; /** 写 worker 输入 job.json,默认 protocol.writeJobSpec(测试可 mock,避免落盘)。 */ writeJobSpec: (job: JobSpec) => void; /** 当前时间戳(ms);测试注入可控时钟(默认 Date.now)。 */ nowMs: () => number; } export interface Orchestrator { /** 跑一轮:先 ingest(outbox→DB)+ reaper(回收死 worker),再领新任务 spawn worker。错误只记日志、不抛。 */ tick(): void; /** 仅领取/spawn 那一步(测试细粒度断言用,不含 ingest/reaper)。 */ claimTick(): void; /** ingest 所有在途 run(测试用)。 */ ingest(): void; /** 回收死 worker(测试用)。 */ reap(): void; } /** * 编排器(监工):自身不跑 agent,只领任务 + spawn 独立 worker + 把 worker 的 outbox 摄取入库 + 回收死 worker。 * 全部 DB 写都在 daemon 进程内(含 ingest 落的 run/result/transition),WS 推送照常经 store.subscribe 触发。 * 无内存在途态——并发闸读 store.countExecuting,退避读 task.nextEligibleAt,故跨 daemon 重启天然持久。 */ export function createOrchestrator(store: Store, log: OrchestratorLogger, deps: Partial = {}): Orchestrator { const d: OrchestratorDeps = { spawnWorker: defaultSpawnWorker, isWorkerAlive, ingestAll, writeJobSpec, nowMs: Date.now, ...deps, }; /** * 本项目可领取的任务:queued(重试/孤儿)+ ready 叶子且 deps 全 done;auto-easy 只挑 easy。 * 退避读持久化的 task.nextEligibleAt(早于它不领)。按调度分降序返回(见 model/scoring.ts)。 */ function claimable(project: Project): Array<{ task: Task; score: number }> { const nowMs = d.nowMs(); const tasks = store.listTasks(project.id); const byId = new Map(tasks.map((t) => [t.id, t])); const parents = new Set(tasks.filter((t) => t.parentId).map((t) => t.parentId as string)); const easyOnly = project.autonomy === 'auto-easy'; const candidates = tasks.filter((t) => { if (t.status !== 'ready' && t.status !== 'queued') return false; if (parents.has(t.id)) return false; // 非叶子(容器)跳过 if (easyOnly && t.complexity !== 'easy') return false; if (t.nextEligibleAt && nowMs < Date.parse(t.nextEligibleAt)) return false; // 退避冷却中 return t.deps.every((dep) => byId.get(dep)?.status === 'done'); }); return rankByScore(candidates, tasks); } /** 领取一个任务:建分支/worktree 路径 → executing → startRun → 写 job.json → spawn worker → 记 pid。 */ function claimOne(project: Project, task: Task, score: number): void { let runId: string | null = null; try { if (task.status === 'ready') store.transition(task.id, 'queued', { by: 'orchestrator', score }); const branch = branchFor(task.id); const dir = worktreeDirFor(project.repoPath, task.id); store.transition(task.id, 'executing', { by: 'orchestrator' }); const run = store.startRun(task.id, 'executor', { worktree: dir, branch }); runId = run.id; d.writeJobSpec({ runId: run.id, task: { ...task, status: 'executing' }, project, worktreeDir: dir, branch, }); log.info(`领取任务 ${task.id}「${task.title}」score=${score} run=${run.id} model=${pickModel(task, project, 'executor')}`); const pid = d.spawnWorker(run.id); store.setWorkerPid(run.id, pid); log.info(`任务 ${task.id} 已起 worker pid=${pid} worktree=${dir}`); } catch (e) { const msg = (e as Error).message; log.error(`任务 ${task.id} 领取/spawn 失败:${msg}`); try { store.failTaskAttempt(task.id, runId, msg); } catch (e2) { log.error(`任务 ${task.id} 失败收尾出错:${(e2 as Error).message}`); } } } /** 领取轮:对每个 active 且 autonomy≠manual 的项目,在 countExecuting= p.concurrency) continue; for (const { task, score } of claimable(p)) { if (active >= p.concurrency) break; claimOne(p, task, score); active++; } } } catch (e) { log.error(`编排器领取轮失败:${(e as Error).message}`); } } /** ingest 所有在途 run 的 outbox(daemon 唯一 DB 写者把 worker 进度/结果落库)。 */ function ingest(): void { try { d.ingestAll(store, log); } catch (e) { log.error(`ingest 轮失败:${(e as Error).message}`); } } /** 回收:对每个 executing 任务判活,死 worker → failTaskAttempt(收尾 + 重试/needs_attention)。 */ function reap(): void { try { for (const { task, run } of store.executingWithLatestExecutorRun()) { const alive = run !== null && d.isWorkerAlive({ pid: run.workerPid, heartbeatAgeMs: heartbeatAgeMs(run.id, d.nowMs()), startedAgeMs: d.nowMs() - Date.parse(run.startedAt), }); if (alive) continue; log.error(`任务 ${task.id} worker 异常退出(run=${run?.id ?? '无'})→ 回收`); try { store.failTaskAttempt(task.id, run?.id ?? null, 'worker 异常退出'); } catch (e) { log.error(`任务 ${task.id} 回收收尾出错:${(e as Error).message}`); } } } catch (e) { log.error(`编排器回收轮失败:${(e as Error).message}`); } } function tick(): void { ingest(); // 先把 worker 进度/结果落库(可能把 executing→exec_review,腾出并发槽) reap(); // 再回收死 worker(可能把 executing→queued/needs_attention,腾出并发槽) claimTick(); // 最后领新任务 } return { tick, claimTick, ingest, reap }; } /** * 接线入口:MAESTRO_ORCH_INTERVAL(秒)控制轮询间隔,默认 15,0=关闭。 * 每轮依次 ingest → reaper → 领取(见 tick)。返回 timer 供 shutdown 时 clearInterval。 */ export function startOrchestrator( store: Store, app: { log: OrchestratorLogger }, deps: Partial = {}, ): NodeJS.Timeout | null { const intervalSec = Number(process.env.MAESTRO_ORCH_INTERVAL ?? 15); if (!Number.isFinite(intervalSec) || intervalSec <= 0) { app.log.info('编排器已关闭(MAESTRO_ORCH_INTERVAL=0)'); return null; } const orch = createOrchestrator(store, app.log, deps); const timer = setInterval(() => orch.tick(), intervalSec * 1000); timer.unref(); app.log.info(`编排器已启用:每 ${intervalSec}s 一轮(ingest→回收→领取,autonomy≠manual 的 active 项目)`); app.log.info(describeSandbox()); // 记录当前生效的沙箱/资源上限策略 return timer; }