8 章设计文档 HTML(架构决策/状态设计/DB Schema/记忆注入/Agent 规格/ 运维护栏/安全沙箱/前端 API 契约),全图深色内联 SVG,附 docs/index.html 索引。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
18 KiB
Maestro 重构设计文档
状态:草稿 · 2026-06-22
范围:Task 模型 · 调度算法 · 流式通信 · Hook 可观测性 · 前端重构 · 后端拆分
1. 背景与目标
Maestro 是一个本地优先的 Git 任务编排守护进程,驱动多项目的 headless Claude Code agent 自动执行代码任务。当前版本(约 1.0)已具备基础的"创建 → 方案 → 执行 → 审批 → 合并"完整闭环,但存在以下痛点需要通过重构解决:
| 类别 | 问题 |
|---|---|
| Task 模型 | 缺少 task_type/scope/ownedFiles/expected_output,拆解产物元数据不足 |
| 调度算法 | 简单优先级排序,缺少依赖链 unlock 价值、饥饿保护、文件耦合感知 |
| 执行管道 | 并发 Code+Security 审查是串行的;verify 过于单一;审查 verdict=reject 不挡 |
| 流式通信 | agent 执行中实时 token 无法到达前端;无 human-in-the-loop 动态暂停 |
| 可观测性 | outbox phase 不广播 WS;无 Trace ID 跨进程;transcript 无 Web 检索 |
| 代码组织 | store.ts(1227 行) God Object、orchestrator.ts(426 行)、web/app.js(2473 行) 亟需拆分 |
| 前端 | 纯原生 JS SPA,无组件化,无类型安全,与 Claude Design System 脱节 |
重构目标:
- Task 模型增强(分类/文件所有权/期望产出)
- 调度算法升级为 CPM-based rankU + 老化加成 + 耦合感知
- 执行管道三类优化(已在 plan 中详细定稿)
- SSE 流式通信 + Hook 拦截 + Trace ID 传播
- 后端按职责拆分(store Repos / daemon 四组件 / API routes)
- 前端迁移至 React 18 + TypeScript + Vite + Zustand,对接 Claude Design System
2. Task 模型增强
2.1 新增字段
// src/model/types.ts — Task 接口扩展
export type TaskType = 'feature' | 'bugfix' | 'refactor' | 'chore' | 'docs';
export type TaskScope = 'file' | 'module' | 'service' | 'cross-service';
export interface Task {
// === 现有字段(保留)===
id, projectId, parentId, depth, title, complexity, status,
priority, deps, plan, spec, operations, approvals, result,
assignee, retryBaseline, nextEligibleAt, lastRunError, createdAt, updatedAt
// === 新增字段 ===
taskType: TaskType | null; // 任务分类(feature/bugfix/refactor/chore/docs)
scope: TaskScope | null; // 改动范围维度
ownedFiles: string[]; // 声明的文件所有权(冲突检测用)
expectedOutput: string | null; // "done" 的可验证描述,供 exec_review 对照
parentVersionId: string | null; // reject 后新建版本指向前一版,构成版本链
version: number; // 任务版本号(每次 reject 递增)
}
2.2 ownedFiles 冲突检测
在 orchestrator.ts 的 claimable() 中增加文件交集检查:
// 已声明 ownedFiles 的在途任务集合
function hasFileConflict(candidate: Task, inFlight: Task[]): boolean {
if (!candidate.ownedFiles?.length) return false;
const candidateSet = new Set(candidate.ownedFiles);
return inFlight.some(t =>
t.ownedFiles?.some(f => candidateSet.has(f))
);
}
- 文件交集 → 降低 score(-0.5/个重叠文件),不硬 block(防饥饿)
agingBonus兜底:等待超 48h 的任务最多加 1 点,确保不被永久回避
2.3 自动元数据填充(MCP 工具)
新增 MCP 工具 suggest_task_metadata:人工输入 title 后,调用 sonnet 分析仓库上下文自动推断 taskType/scope/ownedFiles/expectedOutput,人工确认后写入。触发:MCP 工具调用 or UI "智能填充"按钮。
2.4 planner 输出扩展
Hard/Medium 任务 planner 的 decompose JSON 格式扩展:
{
"plan": "...(分析正文)...",
"subtasks": [
{
"title": "类型定义与接口",
"complexity": "easy",
"priority": 0,
"deps": [],
"ownedFiles": ["src/model/types.ts"],
"expectedOutput": "类型文件通过 typecheck"
},
{
"title": "TaskRepo 实现",
"complexity": "medium",
"priority": 1,
"deps": [0],
"ownedFiles": ["src/store/taskRepo.ts"],
"expectedOutput": "TaskRepo CRUD 方法通过单测"
}
]
}
- 正文先输出 Markdown 子任务表格(供 plan_review 人审)
- JSON 块作为机器解析源,
deps用子任务数组序号引用 - daemon
ingest.ts在子任务全建好后做"序号→taskId"二次映射
3. 调度算法升级(CPM-based rankU)
3.1 算法设计
用 CPM(Critical Path Method)后向传播 替代简单优先级排序:
// src/model/scoring.ts
/** 递归计算任务向后传播的 unlock 价值 */
function rankU(taskId: string, cache: Map<string, number>): number {
if (cache.has(taskId)) return cache.get(taskId)!;
const task = getTask(taskId);
const base = baseScore(task); // P0=3, P1=2, P2=1
const unlockValue = dependents(taskId)
.reduce((sum, dep) => sum + rankU(dep.id, cache), 0);
const result = base + unlockValue;
cache.set(taskId, result);
return result;
}
/** 老化加成:等待越久加分越多,防饥饿 */
function agingBonus(task: Task): number {
const base = baseScore(task);
const waitHours = (Date.now() - new Date(task.createdAt).getTime()) / 3600000;
return Math.min(base, waitHours / 48); // 48h 达到 base 上限
}
/** 耦合惩罚:文件重叠 */
function couplingPenalty(task: Task, inFlight: Task[]): number {
const overlaps = inFlight.reduce((sum, t) => {
const shared = (task.ownedFiles ?? []).filter(f => t.ownedFiles?.includes(f));
return sum + shared.length;
}, 0);
return overlaps * 0.5;
}
/** 最终调度得分 */
function scheduleScore(task: Task, completedDeps: Task[], inFlight: Task[]): number {
return rankU(task.id, new Map())
+ completedDeps.reduce((s, d) => s + baseScore(d), 0) // 链惯性
+ agingBonus(task)
- couplingPenalty(task, inFlight);
}
3.2 CAS 防双重 claim
-- orchestrator claim 阶段
UPDATE tasks
SET status = 'queued', claimed_at = datetime('now')
WHERE id = ? AND status = 'ready' AND claimed_at IS NULL
重要:所有使 task 回到 ready 的转移(reject/requeue/retry)必须同时清空 claimed_at:
UPDATE tasks SET status = 'ready', claimed_at = NULL WHERE id = ?
否则 retry 任务永远命中 claimed_at IS NULL = false,无法再被 claim。
3.3 调度决策日志
每次 claim 记录结构化日志:
{
"taskId": "t-001",
"score": 4.5,
"breakdown": {
"rankU": 3.0,
"chainInertia": 1.0,
"agingBonus": 0.5,
"couplingPenalty": 0.0
},
"competitors": [...]
}
4. 三类执行管道优化(定稿,已在 plan 中逐维确认)
见 .claude/plans/tidy-jumping-shell.md 的完整定稿小结,此处仅列关键决策:
4.1 Planner(任务拆解)
| 维度 | 决策 |
|---|---|
| 触发/执行期锁 | 在途即锁只读:禁改、禁再调度,run 结束解锁 |
| 模型/档位 | hard=fable-5 / medium=opus-4-8 / easy=sonnet-4-6(env 三档可覆盖) |
| 输出格式 | 子任务表格(人审) + 扩展 JSON{title,complexity,priority,deps,ownedFiles,expectedOutput} |
| 审批 | 默认人审;可选 auto-approved+全 easy+改动小 自动放行(默认关) |
4.2 Executor(代码执行)
| 维度 | 决策 |
|---|---|
| 双复审 | code + security 改 Promise.all 并行(≤15min 替代 ≤30min) |
| 复审模型 | 统一 fable-5,不被 project.model 降档 |
| 新增硬闸 | 分项 checks(lint/typecheck/build)+ diff 越界/体量闸 + verdict=reject 变硬闸 |
| 执行前同步 | createWorktree 后先 merge main;分歧超阈值 → needs_attention 重评估 |
4.3 ConflictResolver(解冲突)
| 维度 | 决策 |
|---|---|
| 架构 | 新增 runKind=conflict;专用 pipeline:真正 git merge → CC 解 → commit → 复审 |
| 模型 | 固定 fable-5,不随原任务复杂度降档 |
| 调度 | 插队/预留名额,优先于普通 executor |
| 白名单 | 仅 conflict pipeline 开放 Bash(git merge:*) |
4.4 通用规则(摘录)
- 执行期锁:在途 run → task 只读,run 结束解锁
- 复审独立性:复审只读、用最强模型;verdict=reject 变硬闸
- 模型回退链:
fable-5 → opus-4-8 → sonnet-4-6(单 run 只重试一次) - Reflect 阶段:连续失败 2 次,agent 先分析失败原因再 retry(非盲目重试)
- stages.json 检查点:pipeline 各阶段写入完成状态,重启后跳过已完成阶段
5. 流式通信(Q7)
5.1 SSE 实时 token 推送
架构:cc.ts for-await → onToken 回调 → daemon EventEmitter (per runId) → SSE /api/tasks/:id/stream
// cc.ts - 已有 for await,增加 onToken 钩子
for await (const message of q) {
out.write(JSON.stringify(message) + '\n'); // 保留:持久化到 transcript
options.onToken?.(message); // 新增:实时推送回调
}
// server.ts - 新增 SSE 端点(以 taskId 为索引,内部映射到当前活跃 runId)
fastify.get('/api/tasks/:taskId/stream', (req, reply) => {
reply.raw.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
});
// 查找该 task 当前 in-progress 的 run,按 runId 注册 emitter
const activeRunId = store.getActiveRunId(req.params.taskId);
const emitter = activeRunId ? tokenEmitters.get(activeRunId) : null;
emitter?.on('token', (msg) => reply.raw.write(`data: ${JSON.stringify(msg)}\n\n`));
});
端点设计:URL 用 taskId(对前端友好),内部 tokenEmitters 以 runId 为键。一个 task 可有多次 run,端点始终映射到最新 in-progress run;run 结束时清理对应 emitter。
选型理由:SSE 天然支持重连+Last-Event-ID(断网续读),无需引入 gRPC/额外 WebSocket。
5.2 类 Chat 任务(执行中注入指令)
POST /api/tasks/:id/inject { message: string }
→ daemon 写入 runs/<runId>/inbox.json
→ worker 在 pipeline 断点轮询 inbox.json
→ cc.ts resume() 带入新内容
不做真正交互式会话——worker 是 subprocess,双向实时通道复杂度过高。实际方案:中途补充指令 → 追加到 next prompt turn。
5.3 动态暂停(Human-in-the-loop 增强)
新增 OutboxRecord 类型:
// src/executor/protocol.ts
type OutboxRecord =
| { type: 'phase'; phase: string }
| { type: 'result'; ... }
| { type: 'clarify'; question: string; runId: string } // 新增
| { type: 'error'; ... }
状态机新增 awaiting_input(在 executing 和 exec_review 之间):
executing → awaiting_input (agent 输出 <ask>...</ask>)
awaiting_input → executing (POST /api/tasks/:id/reply 提供答复)
6. Hook 拦截与可观测性(Q8)
6.1 Hook 契约
// .maestro/hooks.ts(项目级)或 ~/.maestro/hooks.ts(全局级)
export interface MaestroHooks {
'before:task:claim'?: (task: Task) => Promise<void | { cancel: string }>;
'after:planner:output'?: (task: Task, plan: string) => Promise<string>; // 可改 plan
'before:execute'?: (task: Task, job: JobSpec) => Promise<void | { cancel: string }>;
'before:merge'?: (task: Task, branch: string) => Promise<void | { cancel: string }>;
'on:conflict'?: (task: Task, files: string[]) => Promise<'auto' | 'manual'>;
'before:exec-review'?: (task: Task, result: TaskResult) => Promise<void | { cancel: string }>;
}
- 超时 5s → 等同于 cancel
- Shell script 方式:
.maestro/hooks/before-execute.sh(exit != 0 = cancel,stdout = reason) - 项目级优先,全局兜底
6.2 Trace ID 传播
// src/executor/protocol.ts - JobSpec 新增
interface JobSpec {
runId: string;
taskId: string;
traceId: string; // 新增:daemon 写 job.json 时 crypto.randomUUID()
// ...
}
- Worker 所有
appendOutbox记录携带traceId - Dashboard 可按 traceId 聚合 plan → execute → review → merge 的完整链路
6.3 Phase 事件广播
ingest.ts 处理 phase 记录时,增加 WebSocket 广播:
case 'phase':
log.info({ taskId, runId, phase: record.phase }, 'pipeline phase');
store.broadcast({ type: 'run.phase', taskId, runId, phase: record.phase }); // 新增
break;
前端看板实时显示"分析中 / 执行中 / 验证中 / 复审中"进度条。
6.4 Transcript 回放
GET /api/tasks/:id/transcript → 流式返回 runs/<runId>/transcript.jsonl
GET /api/tasks/:id/transcript?q=keyword → 服务端 grep 返回匹配行
无需 ElasticSearch,本地 JSONL grep 即可。
7. 后端代码结构重构
7.1 store.ts 拆分
src/store/
├── db.ts # DBAdapter 接口(SqliteAdapter / future PostgresAdapter)
├── store.ts # 入口(组合所有 Repos,提供 subscribe/broadcast)
├── projectRepo.ts # Project CRUD
├── taskRepo.ts # Task CRUD + 状态机守卫
├── runRepo.ts # Run CRUD
├── approvalRepo.ts # ApprovalRecord
├── eventRepo.ts # Event 追加 + 查询
└── metricsRepo.ts # 聚合指标查询
7.2 daemon 拆分
src/daemon/
├── orchestrator.ts # 入口:tick = Scheduler.claim → WorkerManager.spawn/reap → Ingestor.ingest
├── scheduler.ts # 纯调度逻辑(scheduleScore/claimable/claim CAS)
├── workerManager.ts # spawn/reap/heartbeat 检测
├── ingestor.ts # outbox.ndjson → DB 事件(原 ingest.ts)
└── mergeCoordinator.ts # merge-resolve 任务池 + 收口原任务
7.3 API 拆分
src/api/
├── server.ts # Fastify 初始化 + 路由注册 + WS 挂载
├── middleware/
│ └── auth.ts # No-op 占位(future JWT)
├── routes/
│ ├── projects.ts
│ ├── tasks.ts
│ ├── runs.ts
│ ├── approvals.ts
│ └── metrics.ts
└── schemas/ # Fastify JSON Schema 校验
7.4 executor 拆分
src/executor/
├── protocol.ts # 文件协议(含 clarify 类型、traceId)
├── cc.ts # Claude Agent SDK wrapper(含 onToken 钩子)
├── pipelines/
│ ├── executor.ts # 代码执行 pipeline
│ ├── planner.ts # 任务拆解 pipeline
│ ├── conflict.ts # 解冲突 pipeline(新增)
│ └── reviewer.ts # 复审 pipeline(并行 code+security)
└── worker.ts # 入口(读 job.json → 分发到对应 pipeline)
8. 前端重构
8.1 技术栈
| 层 | 选型 |
|---|---|
| 框架 | React 18 + TypeScript |
| 构建 | Vite |
| 状态管理 | Zustand(全局 store:projects/tasks/ws 连接) |
| 样式 | CSS Modules + IBM Plex Mono(Claude Design System 字体) |
| 国际化 | i18next(5 语言:zh/en/es/ja/fr) |
8.2 目录结构
web/
├── index.html
├── vite.config.ts
├── src/
│ ├── main.tsx
│ ├── App.tsx
│ ├── store/ # Zustand stores
│ ├── api/ # REST + SSE + WS client
│ ├── components/ # 通用组件(Button/Badge/Modal...)
│ ├── screens/ # 页面(Dashboard/TaskDetail/Settings)
│ └── i18n/
└── public/
8.3 布局
3 列布局(对齐 Claude Design System):
- 左侧:项目列表(侧栏)
- 中间:任务看板(按状态分组)
- 右侧:任务详情(审批/流式输出/transcript)
8.4 实时特性
- WebSocket:任务状态变更 + Phase 事件 → 看板实时刷新
- SSE:TaskDetail 右侧面板显示 agent 实时 token 输出
- 审批闸:plan_review / spec_review / exec_review → 内联 approve/reject
9. Schema 变更(src/store/schema.sql)
-- tasks 表新增列
ALTER TABLE tasks ADD COLUMN task_type TEXT; -- feature/bugfix/refactor/chore/docs
ALTER TABLE tasks ADD COLUMN scope TEXT; -- file/module/service/cross-service
ALTER TABLE tasks ADD COLUMN owned_files TEXT; -- JSON string[]
ALTER TABLE tasks ADD COLUMN expected_output TEXT; -- 验收描述
ALTER TABLE tasks ADD COLUMN parent_version_id TEXT; -- 版本链前驱
ALTER TABLE tasks ADD COLUMN version INTEGER DEFAULT 1;
ALTER TABLE tasks ADD COLUMN claimed_at TEXT; -- CAS claim 时间戳
-- runs 表 kind 新增 'conflict'(TEXT 无 CHECK,兼容)
-- runs 表新增 trace_id
ALTER TABLE runs ADD COLUMN trace_id TEXT;
-- 新状态 'awaiting_input' 已在 status.ts TRANSITIONS 中处理,无需 schema 改动
10. 实现优先级
| 阶段 | 内容 | 依赖 |
|---|---|---|
| P0(核心正确性) | CAS claim / 执行期锁 / 分项 checks 闸 / verdict→硬闸 | — |
| P0(Task 模型) | 新增 owned_files/expected_output/task_type 字段 + Schema | — |
| P1(调度升级) | rankU + agingBonus + couplingPenalty | Task 模型 |
| P1(管道优化) | 复审并行 / Planner 分档 / stages.json 检查点 | — |
| P1(解冲突) | conflict runKind + 专用 pipeline | — |
| P2(流式通信) | cc.ts onToken + SSE endpoint + Phase WS | — |
| P2(Hook) | MaestroHooks 契约 + shell/TS 两种实现 | — |
| P2(Trace ID) | JobSpec.traceId 传播 + outbox 携带 | — |
| P3(后端拆分) | store Repos / daemon 四组件 / API routes | P0-P1 稳定后 |
| P3(前端迁移) | React+Vite + Zustand + 实时流 | P2 SSE/WS |
附录 A:状态机新增状态
awaiting_input (新增)
← executing (clarify 记录触发)
→ executing (POST /reply 恢复)
→ cancelled
完整状态机见 src/model/status.ts。
附录 B:调研来源
- Agent 1:Task 分类与自动拆解(Plan-and-Execute / DAG vs Tree / MCP auto-fill)
- Agent 2:文件冲突最小化(ownedFiles / interface-first / madge 分析)
- Agent 3:调度算法(CPM rankU / agingBonus / CAS)
- Agent 4:流式通信与 Human-in-the-loop(SSE / clarify / Hook 契约 / Trace ID)
- Agent 5:AI Agent 编排最佳实践(expected_output / stages.json / Reflect 阶段)