diff --git a/docs/09-runtime-implementation.md b/docs/09-runtime-implementation.md new file mode 100644 index 0000000..e548358 --- /dev/null +++ b/docs/09-runtime-implementation.md @@ -0,0 +1,204 @@ +# 09 · Runtime 实现设计(Mastra 落地) + +> 回答「Runtime 到底怎么实现」。API 形态基于 Mastra 当前文档核实 +> (`@mastra/core/workflows` 的 suspend/resume、schedule、快照持久化), +> 实现时以安装版本的内嵌文档为准。 + +## 0. 总思路:薄壳 + 三原语 + 四个自建服务 + +``` +Runtime = + Mastra Workflow ← 流程模板 + Proposal 状态机(suspend/resume = 审批悬挂) + Mastra Agent ← 五类智能体(作为工作流中的受控步骤,不自由漫游) + Mastra Tool ← Skill 注册表(包一层血缘记录) + + 自建:触发服务 · 事件总线(outbox) · 血缘/审计服务 · 记忆落库 +``` + +关键反直觉点:**Agent 不是入口,Workflow 才是**。每个业务流程(申报、复盘、纠偏) +是一条 `createWorkflow` 定义的确定性步骤链,Agent 只在其中特定步骤被调用。 +这就是 02 篇「LLM 决定做什么,工作流引擎决定怎么按步骤做」的代码形态。 + +## 1. 触发服务(三源归一) + +三类触发源全部收敛为「实例化某个 workflow 模板 + 参数」: + +```ts +// 周期触发:直接用 workflow 的 schedule(cron),完全不经过 LLM +export const daySituationWorkflow = createWorkflow({ + id: 'day-ahead-situation', + inputSchema: z.object({ targetDate: z.string() }), + outputSchema: situationReportSchema, + schedule: { cron: '0 6 * * *' }, // D-1 06:00 +}).then(fetchContextStep).then(forecastStep).then(analyzeStep).then(publishStep).commit() + +// 事件触发:outbox 消费者做确定性规则判断后实例化 +eventBus.on('MeterDeviationExceeded', async (evt) => { + const run = await mastra.getWorkflow('intraday-correction').createRun() + await run.start({ inputData: { unitId: evt.unitId, deviation: evt.value } }) +}) + +// 人工触发:唯一经过 LLM 的入口——路由 Agent 做意图解析 +// 输出被 zod 约束为「模板 id + 参数」,LLM 不能发明新流程 +const routed = await routerAgent.generate(userMessage, { + output: z.object({ + workflowId: z.enum(['day-ahead-bid', 'review', 'adhoc-analysis', ...]), + params: z.record(z.unknown()), + }), +}) +``` + +要点:意图解析的输出模式(structured output)把 LLM 的不确定性压缩为 +「在已注册模板中选一个」——选错可纠正、可审计,但不可能执行未定义流程。 + +## 2. Proposal 状态机 = 一条独立的生命周期工作流 + +**设计决策**:可信执行链不内联在每个业务流程尾部,而是实现为唯一的 +`proposal-lifecycle` workflow。业务流程的终点是「提交 Proposal」,生命周期流程接管后续。 +状态机只实现一次,所有 Proposal 类型共用;审批恢复(resume)永远指向同一条流程。 + +```ts +const envelopeGate = createStep({ + id: 'envelope-gate', + inputSchema: proposalSchema, + resumeSchema: z.object({ // 人工审批回传的数据 + decision: z.enum(['approve', 'reject']), + approverId: z.string(), + comment: z.string().optional(), + }), + suspendSchema: z.object({ reason: z.string(), proposalId: z.string() }), + execute: async ({ inputData, resumeData, suspend }) => { + if (resumeData) { // —— 恢复路径:人工已裁决 + await audit.record('HUMAN_DECISION', resumeData) // 审批人身份写入血缘 + if (resumeData.decision === 'reject') throw new ProposalRejected(resumeData) + return { ...inputData, status: 'APPROVED' } + } + const env = await envelopeService.match(inputData) // 自建:包络匹配与边界检查 + if (env.within) { + await audit.record('AUTO_APPROVED', { envelopeId: env.id }) + return { ...inputData, status: 'AUTO_APPROVED', envelopeRef: env.id } + } + await notifyApprovers(inputData, env.requiredLevel) // 推送审批工作台 + return await suspend({ reason: env.reason, proposalId: inputData.id }) // 悬挂 + }, +}) + +export const proposalLifecycle = createWorkflow({ + id: 'proposal-lifecycle', + inputSchema: proposalSchema, + outputSchema: releasedProposalSchema, +}) + .then(ruleCheckStep) // 确定性步骤:调用自建规则引擎(含血缘完整性检查) + .then(simulationStep) // 调用仿真 Skill;告警 → 强制走 envelopeGate 的人工分支 + .then(envelopeGate) // 包络内自动 / 包络外 suspend 等人工 + .then(releaseStep) // 投递适配器(申报)或执行引擎(控制计划) + .commit() +``` + +审批工作台的「批准」按钮就是一次 resume 调用: + +```ts +// POST /approvals/:runId/decide (鉴权走既有权限系统,见 06 篇 §4) +const run = await mastra.getWorkflow('proposal-lifecycle').createRun({ runId }) +await run.resume({ + step: 'envelope-gate', + resumeData: { decision: 'approve', approverId: session.userId }, +}) +``` + +Mastra 将悬挂快照持久化到配置的存储(我们配 PostgreSQL), +**审批悬挂数小时、进程重启、多副本部署都不丢状态**——这正是选它承载状态机的原因。 +申报截止倒计时告警(07 篇 T-30min 升级)由触发服务的巡检 cron 扫描悬挂中的 run 实现。 + +多级审批:L2/L3 包络用嵌套的第二个 suspend 步骤或循环恢复(同一 gate 按 +`requiredLevel` 依次收集多个审批人 resume),实现时按框架版本择优。 + +## 3. Agent 作为受控步骤 + +```ts +const bidStrategyStep = createStep({ + id: 'bid-strategy', + inputSchema: bidContextSchema, // 上游步骤已确定性地备好:预测、账本、规则片段 + outputSchema: bidProposalDraftSchema, + execute: async ({ inputData }) => { + const result = await tradingAgent.generate(buildPrompt(inputData), { + output: bidProposalDraftSchema, // 结构化输出:申报说明 + 对工具输出的引用 + // Agent 工具集仅含:报价优化 MILP、收益测算、规则检索 —— 注册表白名单 + }) + return assembleProposal(result, lineage.current()) // 数字从血缘取,不从 LLM 文本取 + }, +}) +``` + +两条纪律,代码层面强制: + +1. **上下文由上游步骤确定性组装**(fetchContext 步骤查库、读账本、RAG 检索), + Agent 不自己「想查什么查什么」——上下文可复现,才谈得上决策可复现; +2. **`assembleProposal` 只从血缘记录取数字**:LLM 输出中的 payload численные字段 + 一律以 `{toolCallId, path}` 引用形式给出,组装器解引用填值。LLM 想编一个数字, + 类型上就写不进 Proposal(P2 的编译期/运行期双保险)。 + +## 4. Skill 注册表:包一层血缘 + +```ts +function registerSkill(def: SkillDef) { + return createTool({ + id: def.id, + description: def.description, + inputSchema: def.inputSchema, + outputSchema: def.outputSchema, + execute: async (input) => { + const inputRef = await snapshots.put(input) // 不可变输入快照 + const out = await httpCall(def.endpoint, input, def.version) // Python Skill 服务 + const outputRef = await snapshots.put(out) + lineage.record({ tool: def.id, version: def.version, inputRef, outputRef }) + return out + }, + }) +} +``` + +`lineage` 上下文用 Node `AsyncLocalStorage` 随工作流 run 传播(runId → 血缘缓冲), +Proposal 组装时一次性固化进 `lineage.tool_calls`。 + +## 5. 事件总线:一期 outbox,二期 Kafka + +- 一期:PostgreSQL 事务性发件箱(业务写库与事件写入同事务)+ 轮询/LISTEN 消费者。 + 单库即可满足一期吞吐(申报级频次),且事件天然与业务数据同库对账; +- 二期接入遥测级事件量后,outbox 中继到 Kafka,消费者接口不变(薄抽象先行); +- 事件即 00 篇 §4 的业务对象,全部 zod schema + 版本号;事件表即审计事实来源之一。 + +## 6. 记忆落库 + +| 层 | 实现 | +|---|---| +| 工作记忆 | Mastra Memory(限当前任务线程);工作流步骤间传参优先,能不进记忆就不进 | +| 情景记忆 | 不用框架——事件表 + 时序库就是情景记忆(P5),查询走确定性 fetchContext 步骤 | +| 语义记忆 | 自建表(结构化经验 + pgvector 检索),复盘流程写入,fetchContext 按场景注入 | + +刻意选择:**不用 Agent 框架的自动语义召回喂生产决策**——注入什么经验由流程模板 +显式声明(如申报流程注入「同类天气日的预测偏差模式」),可测试、可评审。 + +## 7. 模块结构(代码仓库形态) + +``` +packages/ +├── runtime/ # Mastra 实例、触发服务、路由 Agent、outbox +│ ├── workflows/ # 流程模板:day-ahead-bid / proposal-lifecycle / review / ... +│ ├── agents/ # 五类 Agent 定义(instructions + 工具白名单) +│ └── tools/ # registerSkill 封装的工具注册表 +├── domain/ # 业务对象 zod schemas(Proposal/Envelope/Ledger...)— 纯类型,零依赖 +├── services/ # 自建确定性服务:规则引擎 · 包络 · 持仓账本 · 血缘/审计 +├── adapters/ # 交易平台 / 调度 / 计量 防腐层 +└── skills-py/ # Python Skill 服务群(独立部署,HTTP + JSON Schema 契约) +``` + +依赖方向:`runtime → services → domain`;`domain` 被所有层引用但不引用任何层。 +自建服务不是 Mastra 工具(LLM 不可调用),是工作流步骤直接 import 的普通模块—— +规则校核、包络检查永远不暴露给 LLM 决定是否调用。 + +## 8. 一期验证方式 + +- Mastra Studio(`npm run dev` → localhost:4111)作为开发期调试台: + 单独试跑 Agent、观察工作流步骤 IO——但**不作为**运营界面; +- 每条流程模板配历史数据回放测试(07 篇场景即第一条端到端用例); +- 审批悬挂做混沌测试:悬挂中重启进程 / 双副本竞争 resume / 超时升级,全部有预期行为断言。