feat(executor): SDK 抖动 resume 续跑 + worker 池化评估暂缓 [tsk_74Acz1Nh4K1J]

可选优化任务,先评估收益再做:

- worker 池化:评估后暂缓。当前「一 run ↔ 一 pid」模型干净(判活/回收/
  re-adopt 全靠每 run 一 pid+心跳,worker 与 DB 隔离 + 重启 re-adopt 天然
  崩溃恢复),池化会打破该不变量且只在高并发量下兑现收益。现规模保持
  「每任务一进程 + 失败从头重跑」简单模型,仅记入 DESIGN.md。

- SDK resume:实现 worker 内部健壮性小优化。cc.ts 单次会话因流式异常中断
  (流断/没收到 result,非超时取消、非 max_turns 终态)且已拿到 sessionId 时,
  用 resume 接着原会话续跑一次(保留已干的活,不从头重来),剩余预算不足留 60s。
  与「重启 re-adopt / 整体失败 daemon 从头重跑」正交不替代;与模型回退互斥;
  MAESTRO_SDK_RESUME=0 可关闭。query 改为可注入便于单测。

测试:新增 test/cc.test.ts(6 例:续跑成功/终态不续/无 session 不续/
开关关闭/超时不续/模型回退)。npm test 全绿 163/163。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
wangjia
2026-06-13 15:03:25 +08:00
parent f020e15137
commit 53efe32a6e
4 changed files with 187 additions and 11 deletions
+128
View File
@@ -0,0 +1,128 @@
import { test, before, after } from 'node:test';
import assert from 'node:assert/strict';
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { runClaude, type CCQuery, type CCOptions } from '../src/executor/cc.js';
// cc.ts 的 SDK resume 续跑 / 模型回退逻辑:注入 fake query,不真起 Claude Code。
let dataDir: string;
let prevData: string | undefined;
let prevResume: string | undefined;
before(() => {
dataDir = mkdtempSync(join(tmpdir(), 'maestro-cc-'));
prevData = process.env.MAESTRO_DATA_DIR;
prevResume = process.env.MAESTRO_SDK_RESUME;
process.env.MAESTRO_DATA_DIR = dataDir; // transcript 落临时目录
delete process.env.MAESTRO_SDK_RESUME; // 默认开启 resume
});
after(() => {
if (prevData === undefined) delete process.env.MAESTRO_DATA_DIR; else process.env.MAESTRO_DATA_DIR = prevData;
if (prevResume === undefined) delete process.env.MAESTRO_SDK_RESUME; else process.env.MAESTRO_SDK_RESUME = prevResume;
rmSync(dataDir, { recursive: true, force: true });
});
type Script = (args: { prompt: string; options: Record<string, unknown> }) => AsyncGenerator<unknown>;
/** 串接多段脚本:第 i 次调用走第 i 段(越界用最后一段);记录每次的 options(看 resume)。 */
function makeQuery(scripts: Script[]): { fn: CCQuery; calls: Record<string, unknown>[] } {
const calls: Record<string, unknown>[] = [];
let i = 0;
const fn = ((args: { prompt: string; options: Record<string, unknown> }) => {
calls.push(args.options);
const script = scripts[Math.min(i, scripts.length - 1)];
i += 1;
return script(args);
}) as unknown as CCQuery;
return { fn, calls };
}
function baseOpts(runId: string, timeoutMs = 30_000): CCOptions {
return { prompt: 'do it', cwd: dataDir, model: 'claude-x', runId, maxTurns: 10, timeoutMs, allowedTools: [] };
}
// 段:拿到 sessionId 后流式中途抛错(SDK 抖动)
const transientThrow: Script = async function* (args) {
yield { type: 'system', session_id: 'sess-1' };
void args;
throw new Error('socket hang up');
};
// 段:一条成功 result
const successResult: Script = async function* () {
yield { type: 'result', subtype: 'success', is_error: false, result: '完工', session_id: 'sess-1' };
};
test('SDK 抖动(流中途抛错)+ 有 sessionId → resume 续跑一次并成功', async () => {
const { fn, calls } = makeQuery([transientThrow, successResult]);
const r = await runClaude(baseOpts('r-resume'), fn);
assert.equal(r.ok, true, '续跑后应成功');
assert.equal(r.finalText, '完工');
assert.equal(r.fellBack, false, '这是 resume,不是模型回退');
assert.equal(calls.length, 2, '应正好两次 attempt(首攻 + resume');
assert.equal(calls[0].resume, undefined, '首攻不带 resume');
assert.equal(calls[1].resume, 'sess-1', 'resume 应续上首攻的 sessionId');
});
test('终态错误(result=error_max_turns)→ 不 resume(重来只会再撞墙)', async () => {
const terminal: Script = async function* () {
yield { type: 'result', subtype: 'error_max_turns', is_error: true, session_id: 'sess-1' };
};
const { fn, calls } = makeQuery([terminal, successResult]);
const r = await runClaude(baseOpts('r-terminal'), fn);
assert.equal(r.ok, false);
assert.equal(calls.length, 1, 'max_turns 是终态,不应触发 resume');
assert.match(r.error ?? '', /error_max_turns/);
});
test('首攻无 sessionId(连 system 都没拿到就抛)→ 无可续之 session,不 resume', async () => {
const noSession: Script = async function* () {
throw new Error('connect ECONNREFUSED');
};
const { fn, calls } = makeQuery([noSession, successResult]);
const r = await runClaude(baseOpts('r-nosess'), fn);
assert.equal(r.ok, false);
assert.equal(calls.length, 1, '没有 sessionId 无法 resume');
});
test('MAESTRO_SDK_RESUME=0 → 即便可续也不 resume', async () => {
process.env.MAESTRO_SDK_RESUME = '0';
try {
const { fn, calls } = makeQuery([transientThrow, successResult]);
const r = await runClaude(baseOpts('r-disabled'), fn);
assert.equal(r.ok, false);
assert.equal(calls.length, 1, '开关关闭时不应 resume');
} finally {
delete process.env.MAESTRO_SDK_RESUME;
}
});
test('超时(abort)→ 不 resume(预算已耗尽)', async () => {
// 拿到 session 后挂起,直到 abort 信号触发才抛错;配极小 timeout 强制超时
const hangUntilAbort: Script = async function* (args) {
yield { type: 'system', session_id: 'sess-1' };
await new Promise((_res, rej) => {
const ac = args.options.abortController as AbortController;
ac.signal.addEventListener('abort', () => rej(new Error('aborted')));
});
};
const { fn, calls } = makeQuery([hangUntilAbort, successResult]);
const r = await runClaude(baseOpts('r-timeout', 40), fn);
assert.equal(r.ok, false);
assert.equal(calls.length, 1, '超时导致的 abort 不可 resume');
});
test('模型不可用错误 → 走模型回退(非 resume),续跑标记 fellBack', async () => {
// 首攻报模型不可用(终态 result,非中断)→ 不 resume,应触发模型回退
const modelErr: Script = async function* () {
yield { type: 'result', subtype: 'error_during_execution', is_error: true, session_id: 'sess-1', errors: ['model not_found'] };
};
const { fn, calls } = makeQuery([modelErr, successResult]);
const r = await runClaude(baseOpts('r-modelfb'), fn);
assert.equal(r.ok, true, '回退模型后成功');
assert.equal(r.fellBack, true);
assert.equal(calls.length, 2);
assert.equal(calls[1].resume, undefined, '模型回退是重开会话,不带 resume');
assert.notEqual(calls[1].model, calls[0].model, '回退应换了模型');
});