# 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 脱节 | **重构目标**: 1. Task 模型增强(分类/文件所有权/期望产出) 2. 调度算法升级为 CPM-based rankU + 老化加成 + 耦合感知 3. 执行管道三类优化(已在 plan 中详细定稿) 4. SSE 流式通信 + Hook 拦截 + Trace ID 传播 5. 后端按职责拆分(store Repos / daemon 四组件 / API routes) 6. 前端迁移至 React 18 + TypeScript + Vite + Zustand,对接 Claude Design System --- ## 2. Task 模型增强 ### 2.1 新增字段 ```typescript // 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()` 中增加文件交集检查: ```typescript // 已声明 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 格式扩展: ```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)后向传播** 替代简单优先级排序: ```typescript // src/model/scoring.ts /** 递归计算任务向后传播的 unlock 价值 */ function rankU(taskId: string, cache: Map): 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 ```sql -- 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`: ```sql UPDATE tasks SET status = 'ready', claimed_at = NULL WHERE id = ? ``` 否则 retry 任务永远命中 `claimed_at IS NULL` = false,无法再被 claim。 ### 3.3 调度决策日志 每次 claim 记录结构化日志: ```json { "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` ```typescript // 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//inbox.json → worker 在 pipeline 断点轮询 inbox.json → cc.ts resume() 带入新内容 ``` 不做真正交互式会话——worker 是 subprocess,双向实时通道复杂度过高。实际方案:中途补充指令 → 追加到 next prompt turn。 ### 5.3 动态暂停(Human-in-the-loop 增强) 新增 OutboxRecord 类型: ```typescript // 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 输出 ...) awaiting_input → executing (POST /api/tasks/:id/reply 提供答复) ``` --- ## 6. Hook 拦截与可观测性(Q8) ### 6.1 Hook 契约 ```typescript // .maestro/hooks.ts(项目级)或 ~/.maestro/hooks.ts(全局级) export interface MaestroHooks { 'before:task:claim'?: (task: Task) => Promise; 'after:planner:output'?: (task: Task, plan: string) => Promise; // 可改 plan 'before:execute'?: (task: Task, job: JobSpec) => Promise; 'before:merge'?: (task: Task, branch: string) => Promise; 'on:conflict'?: (task: Task, files: string[]) => Promise<'auto' | 'manual'>; 'before:exec-review'?: (task: Task, result: TaskResult) => Promise; } ``` - 超时 5s → 等同于 cancel - Shell script 方式:`.maestro/hooks/before-execute.sh`(exit != 0 = cancel,stdout = reason) - 项目级优先,全局兜底 ### 6.2 Trace ID 传播 ```typescript // 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 广播: ```typescript 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//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) ```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 阶段)