Files
maestro/src/daemon/orchestrator.ts
T
wangjia 8392944877 feat(planner): 编排器把 analyzing/speccing 也当可执行——派 planner run 自动拆解/写方案
修复"hard/medium 任务建后永远停在 analyzing/speccing"的缺口:编排器现在对 analyzing(hard)
派 planner-decompose、对 speccing(medium) 派 planner-spec(auto-approved 下;auto-easy 仍只执行 easy),
planner 只读跑 CC(无 worktree/无写)产出 → daemon 落库 + 自动提交到 plan_review/spec_review 闸(仍人审)。

- protocol: JobSpec.runKind + spec-result/decompose-result 两型 outbox + DecomposeResult
- runner: runPlanner(只读 Read/Glob/Grep 跑在 repo,opus 档)+ buildPlannerPrompt
- pipeline: runKind 分支 + parseDecompose(取末尾 fenced JSON,非法→failed)
- orchestrator: claimable 纳入 analyzing/speccing + 排除在途;并发/in-flight 从 countExecuting 泛化为
  inflightTaskIds(有 started run 的任务,executor+planner 通用);claimOne 按状态分 executor/planner
  (planner 不转状态);reap 泛化(死 planner→failPlanAttempt)
- ingest: ingestAll 覆盖所有 started run;spec-result→setSpec+spec_review;decompose-result→setPlan+建子任务+plan_review;
  failed 按 kind 分流(planner→failPlanAttempt 退避留态、超限→needs_attention)
- store: inflightTaskIds / liveRunsWithTask / failPlanAttempt;reconcile 泛化到所有 started run
- status: analyzing/speccing 加 →needs_attention(planner 失败超限升级)
- models: planner 角色(opus);status: RunKind 已含 planner
- 测试 +12(pipeline 5 / ingest 3 / orchestrator 4),194 全绿

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-13 16:59:30 +08:00

276 lines
13 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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;
}
/**
* 默认 spawnWorkerunref + 脱离 stdio + 收敛 env,返回子进程 pid。
* 关键:【不 detached(不 setsid)】——保持与 daemon 同会话,worker 内起的 `claude` 才能沿用
* daemon 的鉴权上下文(本机无 API_KEY/凭证文件,鉴权走 macOS keychaindetached 后新会话读不到 → Not logged in)。
* 存活性不依赖 detachedunref 后 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 {
/** 跑一轮:先 ingestoutbox→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<OrchestratorDeps> = {}): Orchestrator {
const d: OrchestratorDeps = {
spawnWorker: defaultSpawnWorker,
isWorkerAlive,
ingestAll,
writeJobSpec,
nowMs: Date.now,
...deps,
};
/**
* 本项目可领取的任务(三类"可执行"状态):
* ready/queued → executor(执行改动);analyzing → planner 拆解 Hardspeccing → planner 写方案 Medium。
* 排除:已在途(inflight=有 started run 的任务)、非叶子容器、退避冷却中、deps 未全 done;
* auto-easy 只挑 easy(故不做规划——analyzing/speccing 都是 hard/medium 被排除)。按调度分降序返回。
*/
function claimable(project: Project, inflight: ReadonlySet<string>): 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 CLAIMABLE = new Set<string>(['ready', 'queued', 'analyzing', 'speccing']);
const candidates = tasks.filter((t) => {
if (!CLAIMABLE.has(t.status)) return false;
if (inflight.has(t.id)) return false; // 已有在途 runexecutor/planner
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);
}
/**
* 领取一个任务并起 worker:
* - ready/queued → executorready→queued→executingstartRun(executor, worktree)job.runKind=executor。
* - analyzing → planner-decomposespeccing → planner-spec:【不转状态】(留 analyzing/speccing 作"规划中"),
* startRun(planner)planner 只读跑在 repoworktreeDir 仅记录用)。失败按 kind 走 failTaskAttempt/failPlanAttempt。
*/
function claimOne(project: Project, task: Task, score: number): void {
const isPlanner = task.status === 'analyzing' || task.status === 'speccing';
const runKind = task.status === 'analyzing' ? 'planner-decompose'
: task.status === 'speccing' ? 'planner-spec' : 'executor';
const dbKind = isPlanner ? 'planner' : 'executor';
const role = isPlanner ? 'planner' : 'executor';
let runId: string | null = null;
try {
const branch = branchFor(task.id);
const dir = worktreeDirFor(project.repoPath, task.id);
if (!isPlanner) {
if (task.status === 'ready') store.transition(task.id, 'queued', { by: 'orchestrator', score });
store.transition(task.id, 'executing', { by: 'orchestrator' });
}
// planner run 只读跑在主仓,worktree 记 repoPathexecutor 记隔离 worktree 路径
const run = store.startRun(task.id, dbKind, { worktree: isPlanner ? project.repoPath : dir, branch });
runId = run.id;
d.writeJobSpec({
runId: run.id,
task: { ...task, status: isPlanner ? task.status : 'executing' },
project, worktreeDir: dir, branch, runKind,
});
log.info(`领取 ${runKind} ${task.id}${task.title}」score=${score} run=${run.id} model=${pickModel(task, project, role)}`);
const pid = d.spawnWorker(run.id);
store.setWorkerPid(run.id, pid);
log.info(`任务 ${task.id} 已起 worker pid=${pid}${runKind}`);
} catch (e) {
const msg = (e as Error).message;
log.error(`任务 ${task.id} 领取/spawn 失败:${msg}`);
try {
if (isPlanner) store.failPlanAttempt(task.id, runId, msg);
else store.failTaskAttempt(task.id, runId, msg);
} catch (e2) {
log.error(`任务 ${task.id} 失败收尾出错:${(e2 as Error).message}`);
}
}
}
/**
* 领取轮:对每个 active 且 autonomy≠manual 的项目,并发闸 = inflightTaskIds(有 started run 的任务,
* executor + planner 通用)。在 active<concurrency 时按 claimable 领新任务并 spawn worker。
*/
function claimTick(): void {
try {
for (const p of store.listProjects()) {
if (p.status !== 'active' || p.autonomy === 'manual') continue;
const inflight = store.inflightTaskIds(p.id);
let active = inflight.size;
if (active >= p.concurrency) continue;
for (const { task, score } of claimable(p, inflight)) {
if (active >= p.concurrency) break;
claimOne(p, task, score);
active++;
}
}
} catch (e) {
log.error(`编排器领取轮失败:${(e as Error).message}`);
}
}
/** ingest 所有在途 run 的 outboxdaemon 唯一 DB 写者把 worker 进度/结果落库)。 */
function ingest(): void {
try {
d.ingestAll(store, log);
} catch (e) {
log.error(`ingest 轮失败:${(e as Error).message}`);
}
}
/** 回收:对每个在途 runexecutor + planner)判活,死 worker → 按 kind 收尾(executor=failTaskAttempt / planner=failPlanAttempt)。 */
function reap(): void {
try {
for (const { task, run } of store.liveRunsWithTask()) {
const alive = 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} ${run.kind})→ 回收`);
try {
if (run.kind === 'executor') store.failTaskAttempt(task.id, run.id, 'worker 异常退出');
else if (run.kind === 'planner') store.failPlanAttempt(task.id, run.id, 'planner worker 异常退出');
else store.finishRun(run.id, 'failed', { error: '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<OrchestratorDeps> = {},
): 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;
}