AI Pipeline 编排器:从 880 行巨石到 233 行 Agent 调度引擎的重构实战

> CoordinatorAgent 是 EAMX 的核心 AI 调度器,最初有 880 行。一次架构评审后,它被重构为 233 行的 streamOrchestrator——按 Image Preprocess → Context Build → LLM Reasoning → Learning 四个 Phase 组织,配合 Tool Registry 替代硬编码 Tool 定义。本文详解重构过程与最终架构。

1. 重构背景

coordinator.agent.ts 的问题(来自架构评审 P3 — 协调器只做编排):

重构目标:Coordinator 只负责"按顺序调动各模块",不做具体的决策和执行。

2. 重构后的四阶段架构

```typescript<br>// orchestrator.ts — 主入口,~233 行<br>export async function streamOrchestrator(opts: OrchestratorOptions) {<br> const { ctx, messages, userMessage, imageBase64, res, ... } = opts;

// ── Phase 1: Context Build ──────────────────────<br> const systemPrompt = buildSystemPrompt(ctx, {<br> entities: entityResults,<br> intentContext: intentResult?.contextBlock,<br> recentPatterns: await getRecentPatterns(tenantId, 5),<br> });

// ── Phase 2: Image Preprocess ───────────────────<br> let llmUserText = userMessage;<br> if (imageBase64) {<br> const imgDesc = await describeImageForCoordinator(imageBase64, mimeType);<br> llmUserText = ${llmUserText}\n\n[图片内容摘要]\n${imgDesc};<br> }

// ── Phase 3: LLM Reasoning ───────────────────────<br> const authorizedTools = toolRegistry.getAuthorized(ctx.ability);<br> const tools = createAuthorizedTools(authorizedTools, ctx);

const result = streamText({<br> model: provider.chat(AI_FAST_MODEL),<br> system: systemPrompt,<br> messages: [...chatHistory, { role: "user", content: llmUserText }],<br> tools,<br> stopWhen: stepCountIs(3), // 最多 3 次 Tool Call<br> maxOutputTokens: 8192,<br> });

// Phase 3 流式输出 → SSE<br> const gen = handleSSEStream(result.fullStream, {<br> res, logPrefix, intentResult, tenantId, userMessage,<br> onToolCall: (toolName, action) => {<br> // ── Phase 4: Learning (fire-and-forget) ──────<br> recordIntentPattern({ tenantId, rawExpression, resolvedIntent, confidence })<br> .catch(() => {});<br> },<br> });<br>}<br>```

每个 Phase 的职责单一、边界清楚:

| Phase | 职责 | 复杂度 |<br>|---|---|---|<br>| Phase 1: Context Build | 构造 System Prompt(注入权限、实体、意图上下文) | 委托给 system-prompt-builder.ts |<br>| Phase 2: Image Preprocess | 视觉模型将图片转为文字摘要 | 委托给 describeImageForCoordinator() |<br>| Phase 3: LLM Reasoning | streamText + Tool Calling + SSE | 本文核心 |<br>| Phase 4: Learning | fire-and-forget 记录意图模式 | 异步,不影响主流程 |

3. Phase 3 核心:Tool Calling + SSE
3.1 Tool Registry 按权限过滤

``typescript<br>// orchestrator.ts — 只暴露用户有权使用的 Tool<br>const authorizedTools = toolRegistry.getAuthorized(ctx.ability);<br>const tools = createAuthorizedTools(authorizedTools, ctx);<br>``

toolRegistry.getAuthorized() 内部调用 ability.can(casl.action, casl.subject) 过滤,确保 LLM 的 Tool Choice 列表中不会出现用户无权调用的操作。

3.2 Tool Factory:注入权限、超时、错误处理

```typescript<br>// tool-factory.ts — 每个 McpTool 被包装为 Vercel AI SDK Tool<br>export function mcpToolToAiSdkTool(mcpTool, ctx) {<br> return tool({<br> description: mcpTool.description,<br> inputSchema: mcpTool.inputSchema,<br> execute: async (input) => {<br> // ① CASL 权限检查<br> if (casl && !ctx.ability.can(casl.action, casl.subject))<br> return { success: false, error: { code: "PERMISSION_DENIED" } };

// ② 超时控制<br> return Promise.race([<br> mcpTool.handler(input, ctx),<br> timeout(15000),<br> ]);<br> },<br> });<br>}<br>```

三个注入点:权限检查、超时控制、错误格式化——全部由 Tool Factory 统一处理,Tool 本身不需要关心这些横切关注点。

3.3 SSE 流式输出

``typescript<br>// sse-handler.ts — Generator 函数处理 streamText 事件流<br>export async function* handleSSEStream(fullStream, opts) {<br> for await (const event of fullStream) {<br> switch (event.type) {<br> case "text-delta":<br> // LLM 逐字输出 → 实时推送给前端<br> opts.res.write(data: ${JSON.stringify({ type: "chunk", content: event.text })}\n\n`);<br> break;

case "tool-result":<br> // Tool 执行完毕 → 推送 UI Card<br> uiCard = event.output;<br> opts.res.write(data: ${JSON.stringify({ type: "ui_card", card: uiCard })}\n\n);<br> break;<br> }<br> }<br>}<br>```

前端接收两种 SSE 事件:chunk(流式文字,打字机效果)和 ui_card(结构化卡片——DynamicForm、故障单、设备列表等)。

4. 两阶段管线集成(铭牌 OCR 特殊路径)

当用户明确说"建档/录入设备"并上传图片时,管线不走 LLM Tool Calling,而是直接走两阶段快速通道:

``typescript<br>// coordinator.agent.ts(保留在旧文件中,待迁移到 orchestrator)<br>if (hasEquipmentCreateIntent && imageBase64) {<br> const rawFields = await extractFromNameplate(imageBase64, mimeType, "铭牌");<br> const classified = classifyFields(rawFields);<br> const formResult = await executeEquipmentAction("prepare_equipment_form", formParams, ctx);<br> // → 直接返回 DynamicForm,绕过 LLM Tool Calling<br>}<br>``

这是性能优化:铭牌建档是确定性流程(拍铭牌 → 填表单),不需要 LLM 做路由决策。

5. 防止 Tool Calling 死循环

``typescript<br>streamText({<br> // ...<br> stopWhen: stepCountIs(3), // 最多 3 步 Tool Call<br>});<br>``

stepCountIs(3) 是 Vercel AI SDK 提供的控制机制——当 LLM 连续调用 Tool 超过 3 次时,强制停止并返回已有结果。这在设备关联备件 + 创建保养计划 + 检查库存的连续调用场景中重要:防止 LLM 在一个对话轮次中无休止地调用 Tool。

6. 核心文件索引

| 文件 | 行数 | 职责 |<br>|---|---|---|<br>| ai/orchestrator.ts | 233 | Pipeline 入口,四阶段调度 |<br>| ai/tool-registry.ts | 86 | Tool 注册 + 按权限过滤 |<br>| ai/tool-factory.ts | 93 | McpTool → AI SDK Tool + 权限/超时/错误注入 |<br>| ai/sse-handler.ts | 118 | Generator 函数,streamText → SSE 事件转换 |<br>| ai/system-prompt-builder.ts | — | System Prompt 动态构建 |<br>| agents/coordinator.agent.ts | 880 | 旧 Coordinator(两阶段管道逻辑仍在此处) |

下一篇:基于 Qwen-VL 的铭牌 OCR 两阶段管线——详细拆解已发布的同名文章,包含 130s→3-8s 的性能修复过程。

作者:白杨,十余年设备资产管理从业经验,EAMX 产品负责人。