SYSTEM OK | ACTION CONTRACTS 247 DOMAINS 35 DEMO ENVIRONMENT 无需申请 | 服务热线 18926139835
首页/洞察/技术系列/AI Pipeline 编排器

AI Pipeline 编排器

CoordinatorAgent 是 EAMX 的核心 AI 调度器,最初有 880 行。一次架构评审之后, 它被重构为 233 行的 streamOrchestrator,按四个 Phase 组织, 并以 Tool Registry 替代硬编码的 Tool 定义。 本文给出此次重构的问题清单、四个阶段的划分与执行顺序、每一阶段的失败处理方式,以及迁移尚未完成的部分。

作者 白杨 · EAMX 产品负责人 发布 2026-09-28 阅读 约 10 分钟 专题 技术系列

01一次架构评审引出的重构

CoordinatorAgent 是 EAMX 的核心 AI 调度器:用户输入、System Prompt 拼装、图片预处理、 模型调用、结果推送与意图学习都经由它。重构之前,该文件为 880 行。 一次架构评审给出结论 P3——协调器只做编排。此后该文件被拆解为 233 行的 streamOrchestrator,以及若干职责单一的模块。

评审记录中的问题共四处,分列如下。

  • 硬编码 Tool 定义。prepare_equipment_form 的 Zod Schema 直接写在协调器内部,约 120 行。 每新增或调整一个 Tool,改动都落在协调器本体上。
  • 硬编码路由逻辑。以关键词列表(EQUIPMENT_CREATE_KEYWORDS、FAULT_INTENT_KEYWORDS) 静态判断走哪条管线。判断依据写在代码之中,与业务表达之间的对应关系需要人工维护。
  • 职责不清。SSE 输出、System Prompt 构建、图片预处理、LLM 调用与意图学习全部耦合在同一处。
  • 多角色冲突。协调器同时承担路由器、翻译器、执行器与学习器四种角色; 任一角色的变化都会影响其余三者。
表 1 · 旧 Coordinator 的四处问题。四处问题指向同一处成因:编排与决策、执行、学习写在了一起。
位置表现对工程的实际影响
Tool 定义 prepare_equipment_form 的 Zod Schema(约 120 行)写在 coordinator 内 Tool 定义与调度逻辑同处一个文件,新增 Tool 须修改协调器本体
路由逻辑 EQUIPMENT_CREATE_KEYWORDS、FAULT_INTENT_KEYWORDS 静态判断管线走向 管线选择取决于代码内的关键词列表,列表与真实表达之间的差异需人工维护
职责边界 SSE 输出、System Prompt 构建、图片预处理、LLM 调用、意图学习耦合在一起 任一环节的改动都需回归验证整条管线
角色数量 同一处同时承担路由、翻译、执行与学习 角色边界不可见,修改的依据不明确

依据:架构评审 P3 的结论(协调器只做编排)与重构提交中的文件行数。

02重构目标:协调器只做编排

评审给出的目标为一句话:协调器只负责按顺序调动各模块,不做具体的决策与执行。

该目标可拆为两条可检验的判据。其一,协调器只决定阶段之间的顺序与数据传递,不决定阶段内部如何实现; 其二,每个阶段的输入与输出均在协调器内可见,阶段内部的处理细节不外溢到协调器。 这两条判据同时构成复核标准:改动某个阶段时,只需确认该阶段的输入输出未变,不需要通读协调器全文。

需要一并说明的是本次重构的边界。管线对外的行为没有变化:同一批动作契约、同一套权限规则、 同一套审计口径均保持原状,改变的是代码内部的职责划分。因此本次重构不涉及任何业务语义的调整, 亦不引入新的对外能力。

本文的核心判断

协调器的规模问题,实质上是职责边界问题。880 行之中真正属于编排的部分只有阶段顺序与数据传递; 其余部分——Tool 定义、关键词路由、System Prompt 拼装、图片摘要、模式记录—— 各自都可以独立成模块。边界划出之后,协调器本身回到 233 行, 而编排能力并未因此减少。

03四个阶段与执行顺序

重构后的主入口为 streamOrchestrator,从调用参数中解构上下文、会话消息、用户文本、 图片数据与响应对象。四个阶段按编号顺序执行:Phase 1 Context Build、Phase 2 Image Preprocess、 Phase 3 LLM Reasoning、Phase 4 Learning。

阶段之间的数据传递为单向:Phase 1 产出的 System Prompt 供 Phase 3 使用; Phase 2 产出的用户文本(含图片内容摘要)供 Phase 3 使用; Phase 3 的流式结果经 SSE 推送至前端,并在 Tool 调用发生时触发 Phase 4; Phase 4 不向主流程返回任何结果。

四条阶段的执行顺序
Context Build → Image Preprocess → LLM Reasoning → Learning

Phase 1 与 Phase 2 均在进入模型之前完成,二者产出 System Prompt 与用户文本两个入参。 Phase 4 在 Phase 3 的 Tool 调用回调中触发,异步执行,不进入主流程的返回值。

orchestrator.ts · 主入口,约 233 行
export async function streamOrchestrator(opts: OrchestratorOptions) {
  const { ctx, messages, userMessage, imageBase64, res, ... } = opts;

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

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

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

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

  // Phase 3 流式输出 → SSE
  const gen = handleSSEStream(result.fullStream, {
    res, logPrefix, intentResult, tenantId, userMessage,
    onToolCall: (toolName, action) => {
      // ── Phase 4: Learning(fire-and-forget) ──────
      recordIntentPattern({ tenantId, rawExpression, resolvedIntent, confidence })
        .catch(() => {});
    },
  });
}
表 2 · 每个阶段的职责单一、边界清楚。第三列记录该阶段委托的实现位置。
阶段职责复杂度与委托
Phase 1: Context Build 构造 System Prompt(注入权限、实体、意图上下文) 委托给 system-prompt-builder.ts
Phase 2: Image Preprocess 由视觉模型将图片转为文字摘要 委托给 describeImageForCoordinator()
Phase 3: LLM Reasoning streamText + Tool Calling + SSE 本文核心
Phase 4: Learning 以 fire-and-forget 方式记录意图模式 异步,不影响主流程

口径说明:源文摘要按「图片预处理 → 上下文构建 → 推理 → 学习」排列;代码中的阶段编号为 Phase 1 Context Build、Phase 2 Image Preprocess。本文的顺序与编号以代码为准。

04Phase 1 与 Phase 2:进入模型之前的两个入参

Phase 1 · Context Build:System Prompt 的动态构造

该阶段调用 buildSystemPrompt(ctx, { ... }),注入三类信息: 实体识别结果 entities、意图识别模块给出的上下文块 intentContext, 以及 await getRecentPatterns(tenantId, 5) 取得的最近 5 条模式。 三者与当前登录身份的权限信息一并进入 System Prompt。

该阶段不在协调器内实现,委托给 system-prompt-builder.ts。 由此产生的直接结果是:权限、实体、意图三类信息在进入模型之前已经确定, 协调器不参与判断这些信息如何生成,只负责将它们的产出交给下一个阶段。

Phase 2 · Image Preprocess:视觉输入转为文字

该阶段是有条件执行的。未携带图片时直接跳过,用户文本原样作为模型入参; 携带图片时,先由 describeImageForCoordinator(imageBase64, mimeType) 生成文字摘要, 再以「[图片内容摘要] + 摘要」的形式追加到用户文本之后。

该阶段同样不在协调器内实现。协调器在此只做两件事:判断是否需要预处理, 以及决定摘要附加在用户文本的哪个位置。图片内容的解析由视觉模型完成, 协调器不持有解析逻辑。

依据:orchestrator.ts 中的 buildSystemPrompt 调用与 describeImageForCoordinator 调用;模式条数为 getRecentPatterns 的入参常数 5。

05Phase 3 之一:Tool Registry 按权限过滤

Phase 3 在进入模型之前先确定本次请求可用的 Tool 集合。协调器调用 toolRegistry.getAuthorized(ctx.ability) 取得授权工具,再由 createAuthorizedTools 将其包装为模型可调用的 Tool,最后随 streamText 一并提交。

getAuthorized 内部调用 ability.can(casl.action, casl.subject) 完成过滤。 这一处过滤决定了模型看到什么:LLM 的 Tool Choice 列表中不会出现用户无权调用的操作。 换言之,权限判定发生在工具清单的生成环节,属前置约束,而非在调用发生后拦截。

tool-registry.ts · 只暴露用户有权使用的 Tool
const authorizedTools = toolRegistry.getAuthorized(ctx.ability);
const tools = createAuthorizedTools(authorizedTools, ctx);
这一节的权限口径

权限过滤的结果与演示环境的口径一致:权限判定在服务端完成,工具清单由服务端按账号权限生成, 客户端拿到的清单即为服务端的授权结果,模型的选择范围里不会出现该账号无权调用的操作。 演示账号为全部试用权限,其清单覆盖全部 247 个动作;清单可由贵方自有 AI 助手连接演示环境后直接报出, 与站内的动作契约总数与业务域覆盖逐项比对,方式见 AI 原生与 Agent 接入 与 可核验性。

06Phase 3 之二:Tool Factory 的三个注入点

Registry 提供 Tool 的定义,Factory 负责把 MCP Tool 包装为 AI SDK Tool。 mcpToolToAiSdkTool(mcpTool, ctx) 返回 tool({ description, inputSchema, execute }), 并在 execute 内部完成三件事。

① 权限检查。execute 内再次判断 ctx.ability.can(casl.action, casl.subject); 不通过时返回 { success: false, error: { code: "PERMISSION_DENIED" } }, 即以结构化错误返回给模型,而非抛出异常中断整个流。

② 超时控制。以 Promise.race([mcpTool.handler(input, ctx), timeout(15000)]) 包裹调用, 单次 Tool 执行的等待上限为 15000 毫秒。

③ 错误格式化。失败结果的结构由 Factory 统一给出,Tool 自身不定义错误结构。

源文对这三个注入点的概括是:权限检查、超时控制、错误格式化全部由 Tool Factory 统一处理, Tool 本身不需要关心这些横切关注点。这一安排的意义在于,新增 Tool 只需提供 description、inputSchema 与 handler, 权限、超时与错误格式三项不随 Tool 数量增长而重复实现。

表 3 · Tool Factory 的三个注入点,均在包装层内完成,对 Tool 定义本身透明。
注入点实现Tool 侧的负担
权限检查 execute 内调用 ability.can(casl.action, casl.subject),不通过则返回 PERMISSION_DENIED 无需实现权限判断
超时控制 Promise.race 包裹 handler 与 timeout(15000) 无需实现超时逻辑
错误格式化 失败结果以结构化错误返回模型,权限不通过时取 { success: false, error } 结构 无需定义错误结构
tool-factory.ts · 每个 McpTool 被包装为 AI SDK Tool
export function mcpToolToAiSdkTool(mcpTool, ctx) {
  return tool({
    description: mcpTool.description,
    inputSchema: mcpTool.inputSchema,
    execute: async (input) => {
      // ① CASL 权限检查
      if (casl && !ctx.ability.can(casl.action, casl.subject))
        return { success: false, error: { code: "PERMISSION_DENIED" } };

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

07Phase 3 之三:SSE 流式输出

handleSSEStream 是一个异步 Generator 函数,遍历 streamText 返回的 fullStream,按 event.type 分派处理。该文件共 118 行。目前处理两类事件。

text-delta。模型逐字输出时触发,管线实时推送 { type: "chunk", content: event.text },前端按打字机效果逐段呈现。

tool-result。Tool 执行完毕时触发,管线推送 { type: "ui_card", card: uiCard },卡片内容取 event.output。

两类事件的到达节奏不同:chunk 随模型输出逐段到达,ui_card 在 Tool 执行完毕后一次性到达。 前端据此分别处理:文字部分增量追加,卡片部分整体渲染。

表 4 · 前端接收的两类 SSE 事件。二者的推送时机与渲染方式均不同。
事件触发来源前端处理
chunk streamText 的 text-delta 事件 逐字追加,形成打字机效果
ui_card Tool 执行完毕后的 tool-result 事件,内容取 event.output 渲染结构化卡片:DynamicForm、故障单、设备列表
sse-handler.ts · Generator 函数处理 streamText 事件流,共 118 行
export async function* handleSSEStream(fullStream, opts) {
  for await (const event of fullStream) {
    switch (event.type) {
      case "text-delta":
        // 模型逐字输出 → 实时推送给前端
        opts.res.write(`data: ${JSON.stringify({ type: "chunk", content: event.text })}\n\n`);
        break;

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

08Phase 4:Learning 的异步执行

Phase 4 在 Phase 3 的 Tool 调用回调中触发,调用 recordIntentPattern({ tenantId, rawExpression, resolvedIntent, confidence }), 记录本次表达与已解析意图之间的对应关系,为后续的意图识别提供依据。

该阶段采用 fire-and-forget 方式:调用后不等待结果,异常由 .catch(() => {}) 吞掉。 由此产生两点性质。其一,记录过程不阻塞 SSE 输出,前端的响应速度不受记录耗时影响; 其二,记录失败不影响本次对话的结果,用户在本次会话中得到的是完整的模型输出与卡片。

需要说明该阶段的时序边界:记录发生在 Tool 调用之时,此时本次的 Tool 选择已经完成。 因此这类记录作用于后续请求,不改变当前轮次已经作出的决策。

依据:orchestrator.ts 中 onToolCall 回调内的 recordIntentPattern 调用;异步性质由 .catch(() => {}) 与源文的 fire-and-forget 描述给出。

09两阶段管线集成:铭牌建档的特殊路径

当用户明确表达建档或录入设备(hasEquipmentCreateIntent)并上传图片时, 管线不进入 LLM Tool Calling,而是直接走一条两阶段快速通道: 先由 extractFromNameplate(imageBase64, mimeType, "铭牌") 提取原始字段, 再经 classifyFields(rawFields) 完成字段分类, 最后由 executeEquipmentAction("prepare_equipment_form", formParams, ctx) 生成表单, 直接返回 DynamicForm。

该分支的判断依据是流程的确定性:铭牌建档的输入是铭牌照片,输出是设备表单, 中间不需要模型参与路由决策。源文将其定位为性能优化—— 对这类确定性流程,省去一次模型往返即为最直接的收益。

需要如实说明的是,这一段逻辑仍保留在旧文件 coordinator.agent.ts(880 行)之中, 源文在代码注释中标注为「保留在旧文件中,待迁移到 orchestrator」。 也就是说,本次重构完成了主入口与四个阶段的拆分,两阶段快速通道的迁移尚未完成。 该部分的完整实现(含铭牌识别的两阶段划分与调优过程)见本专题第 04 篇。

coordinator.agent.ts · 铭牌建档的快速通道(仍留在旧文件中)
if (hasEquipmentCreateIntent && imageBase64) {
  const rawFields = await extractFromNameplate(imageBase64, mimeType, "铭牌");
  const classified = classifyFields(rawFields);
  const formResult = await executeEquipmentAction("prepare_equipment_form", formParams, ctx);
  // → 直接返回 DynamicForm,绕过 LLM Tool Calling
}

10防止 Tool Calling 死循环

streamText 的配置中给出 stopWhen: stepCountIs(3), 即把单个对话轮次内的 Tool 调用步数限制为 3 次。 这一控制机制由 Vercel AI SDK 提供:当模型连续调用的次数达到上限时,强制停止并返回已有结果。

该上限在连续调用的场景中生效。设备关联备件、创建保养计划、检查库存三件事在同一轮次内依次触发时, 若模型继续发起调用,管线在第 3 步之后停止,不会在同一轮次内无休止地调用 Tool。

与 Phase 3 的其他控制手段相比,二者的作用层级不同:本节的 stepCountIs(3) 作用于步数, 限制一轮对话内可调用多少次;Tool Factory 中的超时控制作用于单次执行的时长,限制每一次调用等待多久。 两者需要同时存在:只限步数不足以约束单次执行的等待时间,只限时长不足以约束同一轮次内的调用次数。

stopWhen · 单轮次内最多 3 步 Tool Call
streamText({
  // ...
  stopWhen: stepCountIs(3),  // 最多 3 步 Tool Call
});

11核心文件索引

重构涉及的文件与行数列示如下。其中三处为本次拆分后新增的模块,一处为拆分的落点, 一处为仍待迁移的旧文件。

表 5 · 核心文件索引。行数为重构当次提交的文件长度,随迭代变动。
文件行数职责
ai/orchestrator.ts 233 Pipeline 入口,四阶段调度
ai/tool-registry.ts 86 Tool 注册与按权限过滤
ai/tool-factory.ts 93 McpTool 转 AI SDK Tool,并注入权限、超时与错误处理
ai/sse-handler.ts 118 Generator 函数,streamText 事件流转换为 SSE 事件
ai/system-prompt-builder.ts — System Prompt 动态构建
agents/coordinator.agent.ts 880 旧 Coordinator,两阶段管线的逻辑仍在此处

依据:重构提交中的文件清单与行数;system-prompt-builder.ts 不计入本次的行数统计。

12作者与依据

作者
EAMX 产品团队 · 产品与解决方案
依据
架构评审 P3 的结论(协调器只做编排);coordinator.agent.ts 与 orchestrator.ts 的重构记录
数据来源
核心文件索引中的文件行数;Tool Registry 与 Tool Factory 的实现代码;SSE 事件处理的实现代码
边界
本文不含并发数、耗时与重试次数的测试口径;铭牌 OCR 的调优过程与实测口径另见技术系列 04
更新日期
2026-09-28
01相关产品能力

本文所述的问题,系统如何解决

编排在系统内部完成,其对外结果体现为三条通道下的权限与审计口径。下列各页说明能力本身与核验方式。

02相关文章

同专题与跨专题的延伸

本文所述的每一项,均可由贵方自有 Agent 核对

四个阶段的划分、按权限过滤的工具清单、该账号可见工具数由端点实时返回—— 上述各项无需采信本文陈述,可连接演示环境后由贵方助手自行报出。

信息中心 / 技术评估
关心编排方式与接入细节
业务部门 / 执行层
关心这些编排在日常业务中的落点