当一个 Agent 只调用一个模型时,代码看起来像一个函数:接收提示词、等待响应、把结果交给用户。系统一旦进入生产环境,真实问题会迅速从“模型会不会回答”变成“哪个模型在什么时间、以什么预算、按照什么协议完成哪个子任务”。文本推理、视觉理解、图片生成和视频生成的延迟曲线、输入格式、限流策略、计费单位都不相同。把这些差异直接散落在业务代码中,得到的不是一个智能系统,而是一组难以升级的厂商适配器。在这里插入图片描述

本文从软件工程和中间件角度拆解 Agentic Workflow。核心观点是:多模型 Agent 的关键组件不是 Prompt,而是一个能表达能力、状态、预算、质量与故障边界的调度层。这个调度层将异构模型包装成可观测的任务执行单元,再用动态 Token 路由、语义缓存和幂等状态机把“模型调用”变成可治理的基础设施。文中的 Benchmark 数值用于说明测量方法,不代表任何供应商的固定性能;模型名也应在部署配置中使用可替换的路由别名。

一、范式演进:单体 LLM 调优到 Agentic 多模型协同的架构跨越

1. 从 Prompt 问题转向系统问题

早期 LLM 应用通常把质量问题归因于 Prompt:加入角色、补充上下文、设置格式约束,再反复调整温度。单模型模式在问答和摘要任务上很有效,但它有三个结构性上限。
在这里插入图片描述

第一,模型的知识、工具权限和上下文窗口被绑定在一次请求内,无法为不同阶段选择不同能力。第二,模型输出既是自然语言又可能是下一步操作,业务代码必须用正则表达式猜测意图。第三,延迟、失败重试和成本都发生在一个黑盒调用里,难以区分是检索慢、推理慢还是外部工具失败。

Agentic Workflow 将任务拆成一组有依赖关系的工作单元。规划器负责把用户目标转换为任务图,执行器负责调用工具或模型,验证器负责检查格式和事实,记忆层负责保存可复用的中间状态,调度器则决定每个节点使用哪一个模型。它类似传统分布式系统中的编排引擎,只是节点的输出带有概率性,需要额外的评估和回滚机制。

这种架构还改变了团队分工。算法团队不再直接把模型 ID 写进业务服务,而是维护提示模板、评估集与能力门槛;平台团队维护路由、密钥、配额和观测;业务团队只声明目标与验收条件。三者通过版本化契约协作。模型升级时,平台先镜像一部分请求,在不影响用户结果的情况下比较新旧输出,再由质量门决定是否扩大流量。升级因而从一次高风险代码发布,变成可回滚的配置变更。

2. 多模态能力的解耦边界

“全能模型”并不意味着所有模态都应该由同一实例处理。文本推理强调长链路一致性和结构化输出,视觉模型更关注图像编码与空间关系,生成模型则受分辨率、采样步数和任务队列影响。将它们拆开以后,可以分别设置超时、并发额度和质量阈值,再在工作流层进行组合。

一个可落地的任务图如下:

用户目标
   |
   v
需求解析器 -> 任务图(Task DAG) -> 能力筛选器
                                |
              +-----------------+-----------------+
              |                 |                 |
         文本推理节点       视觉理解节点       时序生成节点
          (脚本/计划)       (素材/布局)        (图像/视频)
              |                 |                 |
              +-----------------+-----------------+
                                v
                     结果校验 -> 资产归档 -> 交付

这里的边界不是按厂商划分,而是按能力契约划分。一个名为 reasoning.deep 的路由别名可以在不同区域映射到不同模型;vision.extract 可以在低负载时使用轻量模型,在复杂图表时切换到高推理模型。业务代码只依赖能力名、输入输出 Schema 和 SLO,不直接依赖某个模型的版本字符串。

3. Agentic Workflow 的最小闭环

一个健壮闭环至少包含五个动作:计划、执行、观察、反思和提交。计划阶段生成有序任务和验收条件;执行阶段调用模型或工具;观察阶段记录原始响应、耗时、Token、错误码和外部资产;反思阶段用规则或另一个模型检查结果;提交阶段把通过校验的结果写入版本化存储。任何一步都应该产生事件,而不是仅仅返回一个字符串。

事件可以使用如下字段:run_idtask_idparent_idcapabilityattemptstatusinput_refoutput_refusagedeadlinetrace_id。其中 input_refoutput_ref 指向资产仓库或对象存储,避免把大图片和视频塞入日志。事件采用追加写入,重放时可以从最近一个检查点恢复,排查问题时也能还原模型选择的依据。

任务图本身也要版本化。规划器可能在执行中发现缺少素材并追加节点,但不能静默修改原图;它应生成 dag_version=2,保留变更原因和父版本。这样,后续的成本统计才能解释为什么同类需求有时调用五次、有时调用七次。对于不可逆工具,例如发送消息或写入外部系统,节点还要声明副作用类型,并在重放时跳过已经提交的操作。

4. 为什么中间件成为主角

当任务数量增加,Prompt 调优带来的边际收益会下降,而中间件的收益会扩大。统一协议可以减少 SDK 数量,能力目录可以避免调用不存在的模型,预算策略可以阻止异常循环,故障分类可以区分可重试与不可重试错误。

换句话说,模型负责生成能力,中间件负责把能力放入可控的生产流程。

二、异构模型选型矩阵:推理、视觉与时序生成的 Benchmark 实测

在这里插入图片描述

1. 先定义任务,再定义指标

“哪个模型最好”不是可执行的问题。应先定义任务样本和验收条件,再比较同一能力类别中的候选路由。

文本推理可以使用结构化计划、代码修复和长文档问答三类样本;视觉理解需要覆盖表格、低清截图和多图对照;生成任务则固定分辨率、时长、采样步数以及安全过滤策略。样本要脱敏、版本化,并保留一组从未参与提示词调优的隔离集。

在线指标至少包括首字延迟 TTFT、完整响应延迟、每秒输出 Token、错误率、重试率、P95/P99 延迟和单位任务成本。离线指标要根据任务定义,例如 JSON Schema 通过率、事实核验通过率、代码测试通过率、视觉字段召回率和视频约束满足率。稳定性不能只看平均成功率,还要看限流后的恢复时间以及跨区域网络抖动。

2. 一套可复现的测试协议

测试客户端固定并发度、请求顺序和超时,使用同一网络出口,并在请求头中写入测试批次。流式响应要分别记录 TTFT 与总耗时;非流式响应只记录总耗时会掩盖首字体验。每个候选路由运行至少三轮,报告中位数和 P95,避免偶然缓存或冷启动造成误判。对于生成模型,应单独记录排队时间、执行时间和资产下载时间。

为了让比较具备统计意义,测试报告还应给出样本量、置信区间和失败样本分布。若两个路由的平均分只相差零点几个百分点,而置信区间大面积重叠,就不应宣称存在确定优势。

可采用分层抽样保证短文本、长文本、中文、英文、含工具调用和不含工具调用的比例稳定。对随机生成任务,应固定一部分种子用于回归,同时保留随机种子用于观察多样性,防止模型只针对固定样本表现良好。

下面是一组示例数据,数字只用于演示报表结构。真实项目应将 alias 与模型控制台中的可用 ID 绑定,并把价格表作为独立配置版本管理。

能力类别路由别名TTFT 中位数P95 总延迟吞吐/任务结构化通过率单任务成本指数
深度推理reasoning.deep1.9 s12.4 s41 tok/s96.2%1.00
通用文本reasoning.fast0.8 s5.1 s78 tok/s91.5%0.32
视觉抽取vision.extract1.4 s7.8 s6 img/min94.1%0.68
图片生成image.storyboard2.6 s18.7 s1.9 img/min88.0%0.74
视频生成video.short4.2 s46.5 s0.6 clip/min82.7%2.40

3. 模型实体如何进入选型矩阵

DeepSeek-R1、DeepSeek-V3、Claude 3.5 Sonnet、GPT-4o 等名称可以作为候选路由的历史基线或控制台别名,但不应写死在业务逻辑中。推理任务通常关注复杂规划、工具调用和长上下文一致性;Qwen2.5-72B、Llama-3.3 一类候选可用于分类、改写和短摘要,但是否合适仍由本地评测决定。

对于视觉和生成任务,Flux.1 Pro、Midjourney V6、Kling V3 API、MiniMax Video-01 等可作为能力样本进行对比,实际可用性取决于所在区域、账户权限、输入协议和服务状态。某个模型在离线集上领先,并不意味着它在高并发或长时间运行时仍然领先。

模型命名需要一层明确的命名空间。诸如 GPT-5.6、Qwen 3.7 Max、Sora API 映射等字符串如果来自内部路由配置,只能被解释为别名,不能据此推断供应商已经公开相同名称或能力。

配置应同时保存 route_aliasproviderprovider_model_idrevisioneffective_at,运行日志记录真正解析后的值。这样既能保留业务侧稳定名称,又能避免文档、监控与实际调用对象不一致。

选型矩阵还应包含“不可用条件”。例如,视频节点可能不支持同步响应,图像节点可能要求对象存储地址,某些文本节点不返回 Token 明细。调度器读取这些约束后,才知道哪些任务可以并行,哪些任务必须等待资产落盘。把限制条件显式化,通常比再增加一个模型更能提升整体稳定性。

Benchmark 不是上线前只运行一次的报告,而应成为持续任务。平台每天从脱敏流量中抽取代表性样本,用当前生产版本和候选版本做盲测;每周分析质量漂移,每次价格、上下文限制或安全策略变化时触发完整回归。

评测结果带上数据集版本和执行环境,过期结果自动降低决策权重。这样,路由依据反映的是近期真实负载,而不是几个月前的一张静态排行榜。

4. 从分数转向决策边界

假设一个任务同时要求事实准确、响应低于十秒、成本不超过预算,那么评分可以写成:

[
Score(m)=w_q Q(m)+w_l L(m)+w_c C(m)+w_s S(m)-P(m)
]

其中 Q 是质量分,L 是延迟归一化得分,C 是成本得分,S 是稳定性得分,P 是违反硬约束的惩罚。

权重只能在业务验收后确定,不能为了让某个模型胜出而反向调整。更稳妥的做法是先用硬过滤器排除不满足上下文长度、输出格式、区域和截止时间的候选,再对剩余路由排序。

三、Agentic Workflow 中的状态机与 Memory 机制设计

在这里插入图片描述

1. 用状态机取代隐式循环

Agent 常见的隐式循环是“调用模型,发现结果不对,再把结果塞回 Prompt”。这种写法无法限制循环次数,也无法判断重试是否会产生重复扣费。

应将每个任务建模为有限状态机:PENDINGRUNNINGWAITING_TOOLVALIDATINGSUCCEEDEDRETRYABLEFAILEDCANCELLED。状态转换必须记录原因和时间戳,只有允许的边才能被执行器接受。

例如,模型返回格式错误时可以从 VALIDATING 转到 RETRYABLE,并附带修复提示;鉴权失败则直接进入 FAILED,防止无意义重试;视频生成进入 WAITING_TOOL 后,由回调或轮询事件推动状态,而不是占用一个 HTTP 连接。状态机使每次执行都可恢复,也让运营人员能够按失败类别统计。

2. 短期记忆、长期记忆与工作记忆

短期记忆保存当前任务最近几轮消息,适合维持指代和格式约束;长期记忆保存经过筛选的事实、用户偏好和历史决策,通常需要向量化后检索;工作记忆保存当前 DAG 的中间产物,例如脚本草稿、素材清单和校验报告。

三者不能混为一个聊天历史,否则上下文会快速膨胀,旧信息还可能覆盖新约束。

长期记忆写入前要经过去重、敏感信息脱敏和来源标注。每条记忆至少有 memory_idtenant_idsource_runcreated_atexpires_atconfidence。检索结果不仅返回文本,还要返回来源和置信度;当置信度低于阈值时,Agent 应主动询问或调用检索工具,而不是把相似句子当成事实。

3. 上下文剪裁算法

上下文窗口的核心不是“尽量塞满”,而是让每个 Token 对当前任务有可解释的价值。可以为消息定义优先级:系统约束为 1.0,当前任务输入为 0.9,未完成工具结果为 0.85,最近对话为 0.6,历史闲聊为 0.2。

对每条候选消息计算:

[
Value(x)=a\cdot Relevance(x)+b\cdot Recency(x)+c\cdot Reliability(x)-d\cdot Length(x)
]

Value 排序后,在预算内选择消息,并保留父任务的约束摘要。摘要不能简单截断句尾,应由结构化字段表达目标、禁止事项、已确认事实、待解决问题和引用资产。这样即便切换到另一家模型,执行语义也不会随着自然语言风格变化。

剪裁之后必须做“约束存活检查”。系统逐项确认输出格式、禁止操作、时间范围、用户已确认选择等关键字段仍在上下文中;若缺失,就优先移除低价值历史再补回约束。

对于跨十几个节点的长流程,可以把工作记忆压缩为事件摘要,但摘要必须引用原始事件 ID。任何模型生成的摘要都只是派生数据,发生争议时以不可变的原始事件和资产哈希为准。

4. LangChain 与 LlamaIndex 的边界

LangChain、LlamaIndex、AutoGPT 等框架可以快速搭建链路,但生产系统仍需自建执行边界。框架的链式抽象适合表达 Prompt 和工具组合,企业系统还要补充租户隔离、密钥托管、幂等键、配额、审计和版本迁移。

更合适的做法是把框架当作“任务定义层”,由自己的运行时负责实际的调度、超时和事件写入。

5. 并发与一致性

同一任务图内没有依赖的节点可以并行,但并行不等于无限创建协程。执行器要按租户、能力类别和供应商分别设置 semaphore,并用队列承载超过额度的任务。

写入状态时使用乐观锁:更新条件包含旧版本号,若版本不一致则重新读取,防止回调和超时线程同时覆盖结果。资产引用也需要版本号,否则同名文件可能在重试后被覆盖,导致最终交付无法复现。

恢复策略要区分“计算可重做”和“副作用不可重做”。纯推理节点可以从检查点重新执行,代价是额外 Token;发信、扣款或发布节点必须依赖幂等收据确认是否已经完成。

调度器启动时扫描长期处于 RUNNING 的租约,只有租约过期且不存在有效执行者时才接管。这个租约机制能够防止多实例部署中两个 worker 同时处理同一任务。

四、成本与延迟优化:动态 Token 路由算法与缓存策略

在这里插入图片描述

1. Token 预算是任务级资源

成本核算不能只看一次 API 返回的 usage。一个任务可能包含规划、检索、工具调用、修复和最终生成多个阶段,需要把每次调用关联到同一个 run_id

任务成本可以近似为:

[
Cost(run)=\sum_i(P^{in}i\cdot T{in}_i+P{out}i\cdot T^{out}i)+C{tool}+C{storage}+C{egress}
]

P 是价格配置,T 是输入或输出 Token,工具、存储和出网是额外成本。路由器在计划阶段先估算上限,再为每个节点分配预算。若执行过程中预测成本超过上限,应降低上下文、切换轻量模型或暂停并请求人工确认。

2. 任务复杂度评分

可以从输入长度、约束数量、工具数量、历史失败率和输出结构复杂度计算 complexity。例如,代码修复任务的约束数、需要读取的文件数和测试命令数都较高,适合深度推理路由;标题改写或分类则可以使用快速路由。复杂度只用于候选排序,不能替代质量门槛。

from dataclasses import dataclass

@dataclass
class TaskProfile:
    input_tokens: int
    constraints: int
    tools: int
    output_schema_depth: int
    failure_rate: float

def complexity(p: TaskProfile) -> float:
    length_score = min(p.input_tokens / 4000, 1.0)
    constraint_score = min(p.constraints / 12, 1.0)
    tool_score = min(p.tools / 6, 1.0)
    schema_score = min(p.output_schema_depth / 8, 1.0)
    return (
        0.30 * length_score
        + 0.25 * constraint_score
        + 0.20 * tool_score
        + 0.15 * schema_score
        + 0.10 * p.failure_rate
    )

当分数低于 0.35 时优先选择低延迟路由,介于 0.35 和 0.7 时比较质量与成本,高于 0.7 时优先深度推理。阈值应通过离线回放校准,并按租户和任务类型分桶,避免一个全局阈值把所有请求推向昂贵模型。

3. 动态 Token 路由

Token Dynamic Routing(动态 Token 路由)可以使用“硬约束加软评分”两阶段算法。硬约束检查上下文长度、输出格式、模态支持、区域合规、剩余预算和截止时间;软评分综合质量预测、实时 P95、失败率、Token 价格和队列长度。

实时指标要设置平滑窗口,避免单个尖峰让流量在几个路由之间来回抖动。切换模型时应记录 route_reason,否则事后无法解释成本变化。

为了进一步抑制路由抖动,可以引入滞回区间:当前路由的得分只要没有低于替代路由一个显著差值,就保持原选择;只有连续多个窗口越过阈值才切换。

新路由先接收少量影子流量,质量、错误率和成本同时满足门槛后再逐步扩容。若线上质量探针下降,即使接口健康检查仍然成功,也应自动缩小权重。这与传统负载均衡仅看连接健康的做法不同。

4. Semantic Cache 的正确边界

Semantic Cache 不等于把所有答案永久复用。缓存键应包含规范化 Prompt、系统版本、工具版本、知识库快照、输出 Schema 和安全策略版本。先用向量相似度找到候选,再用关键约束校验和小模型复核,最后才返回缓存结果。

对实时数据、个性化权限、随机创作和包含时间条件的查询,默认禁用语义缓存。

缓存还要设置新鲜度与失效事件。知识库更新、模型升级或安全规则变化都可以广播失效消息。命中缓存时仍需记录一次“虚拟调用”事件,标记节省的 Token 和命中原因,才能在 FinOps 报表中区分真实调用与缓存收益。缓存内容按租户隔离,避免相似问题跨租户泄露信息。

5. 延迟预算与背压

端到端延迟可以分解成排队、网络、首 Token、生成、工具等待和资产下载。调度器为每一段分配预算,并在剩余时间不足时停止低价值步骤。

例如,视频生成已经超出主请求截止时间,就应返回可轮询的任务句柄和当前资产,而不是继续阻塞。队列长度超过阈值时启用背压:降低并发、延迟低优先级任务或切换到异步模式。背压是保护系统的机制,不应依赖供应商的偶然限流。

成本治理还需要异常检测。一个正常摘要节点突然输出数万 Token,可能是停止条件失效;一个 Agent 在同一工具上反复调用,可能是规划循环。

可按任务类型建立 Token 分位数基线,超过 P99 时立即停止当前节点并记录上下文。预算告警应同时包含预计完成成本和已经沉没的成本,使决策者能够判断继续、降级还是终止,而不是在月底账单中才发现异常。

五、多模态 Agent 开发中的“接口碎片化”工程困境

在这里插入图片描述

1. Schema 不一致

文本接口可能接受 messages,图像接口可能要求 promptsizesteps,视频接口则常常返回一个异步任务 ID。即使都声称兼容 OpenAI Compatible Protocol,流式事件名、错误字段和 Token 统计也可能不同。

若业务层直接拼装这些字段,新增一个模型就要复制一套条件分支,最终形成“if 厂商”的不可维护代码。

统一层应定义内部协议:输入由 content[] 组成,每个内容块带 kindurimimemetadata;输出统一为 parts[],并区分文本、结构化对象、资产引用和异步句柄。适配器只负责把内部协议转换为供应商协议,不能在适配器里偷偷改变任务语义。

协议版本升级要遵循向后兼容原则。新增字段默认可选,删除或改变语义则发布新的主版本;网关在入口处完成版本协商,并在响应中回传实际使用的协议版本。

对工具调用尤其要避免把供应商私有字段直接透传,因为下游一旦依赖这些字段,统一层就失去替换能力。适配器契约测试应覆盖正常响应、流式中断、空结果、超大资产和未知字段。

2. 鉴权、计费与限流耦合

多个 API Key 分散在环境变量、配置文件和脚本中,会造成轮换困难和审计盲区。密钥应放在专用的 Secret Store,网关只接收短时凭证或内部租户令牌。

每次调用写入 provider_costtenant_costbudget_remaining,并支持按日、按项目和按能力类别聚合。RPS、RPM、并发任务数和输出 Token 上限要分别计数,不能用一个“请求次数”覆盖所有限流维度。

3. 跨区域网络与错误语义

跨国网络延迟会放大流式首包时间,TLS 握手和连接复用对短请求尤其明显。HTTP 客户端应使用连接池、合理的 keep-alive 和指数退避,但重试只针对超时、暂时不可用和明确的速率限制。

鉴权失败、参数错误、内容安全拒绝和上下文过长不能自动重试。错误对象至少携带 categoryretry_afterprovider_request_id 和用户可读的 safe_message,从而让上层决定降级或中止。

4. 可观测性缺口

如果只记录最终答案,无法解释一次 Agent 为什么昂贵或失败。每个节点都应暴露 RED 指标(请求率、错误率、延迟),并附加 Token、缓存命中、队列等待和路由变更。

Trace 中不应存放原始敏感 Prompt,可记录脱敏摘要、哈希和资产引用。对视频等长任务,要把异步回调和轮询事件串到同一条 trace,否则用户看到的是“网关成功、业务超时”的假象。

内容安全也属于接口碎片化的一部分。不同上游对拒绝原因、风险类别和可申诉信息的表达不同,统一层应映射成有限的安全事件类型,并保留上游原始代码供内部审计。

网关不能为了追求成功率而把被一个路由拒绝的请求无条件转发给另一个路由;是否允许替代执行必须由租户策略和风险类别共同决定。

六、架构破局:异构模型 API 统一路由网关的工程实现

在这里插入图片描述

1. 网关分层

统一网关(API Gateway)可以拆成入口、策略、适配器和运行时四层。入口负责鉴权、租户识别、请求限流和幂等键校验;策略层读取能力目录,执行硬过滤、预算检查和动态排序;适配器把内部协议转换为各供应商格式;运行时负责连接池、流式事件、重试、状态写入和指标上报。

四层之间只传递版本化的数据结构,不直接传递某厂商 SDK 对象。

Client -> Auth/Quota -> Policy Router -> Adapter Pool -> Provider API
             |              |               |
             v              v               v
          Audit Log     Route Reason     Metrics/Trace
             |              |               |
             +-------- State Store <--------+

能力目录是网关的“类型系统”。每条记录包含 capability、输入模态、最大上下文、是否支持流式、是否异步、区域、价格版本、健康状态和质量门槛。健康检查只验证连接,不代表质量合格;质量门槛来自隔离评估集和线上抽样复核,两者需要分别标记。

策略层最好保持无状态,把实时计数放在独立配额服务,把模型元数据放在版本化目录,把任务状态放在事务数据库。这样多个网关实例可以水平扩展,任何实例都能处理同一租户的下一次请求。

配置发布采用“校验、预热、切换、观察、回滚”五步:先验证 Schema,再建立上游连接池,然后切换少量流量,观察错误和质量,最后扩大权重。配置错误不应等到真实请求到来才暴露。

这类网关也不同于只做轮询的 Load Balancing(负载均衡)。传统负载均衡假设后端实例功能近似相同,而模型路由面对的是能力、质量、价格和上下文上限均不相同的候选。

FastAPI 适合承载类型明确的入口,Pydantic 负责边界校验,Asyncio 负责并发 I/O;真正的路由决策仍应由独立策略模块完成,避免 Web 框架、供应商 SDK 和业务规则互相耦合。

2. OpenAI 兼容的内部请求

内部请求可保持类似 Chat Completions 的形状,但增加 task_idcapabilitydeadline_msbudgetidempotency_key

流式响应统一为事件序列:response.startedresponse.deltaresponse.tool_callresponse.completedresponse.failed。供应商的增量字段在适配器内转换,客户端无需知道是 SSE、轮询还是 WebSocket。

在本地或云端测试环境中,网关可以把一个统一协议映射节点配置为测试沙盒的上游,例如 base_url="https://178.nz/yinc",并用环境变量替换真实凭证;该配置只用于验证协议转换、限流和观测链路,不能把它写成业务逻辑中的固定依赖。

3. FastAPI 异步代理骨架

下面的示例展示入口校验、路由选择、超时和统一响应。真实部署还需要 Secret Store、持久化状态库、分布式限流器以及更完整的错误映射。

import asyncio
import os
import time
from typing import Any

import httpx
from fastapi import FastAPI, Header, HTTPException
from pydantic import BaseModel, Field

app = FastAPI()

class GatewayRequest(BaseModel):
    capability: str
    messages: list[dict[str, Any]] = Field(default_factory=list)
    deadline_ms: int = Field(default=30000, ge=500, le=120000)
    budget: float = Field(default=1.0, ge=0)
    stream: bool = False

ROUTES = {
    "reasoning.deep": {
        "url": os.environ["PROVIDER_A_URL"],
        "timeout": 45.0
    },
    "reasoning.fast": {
        "url": os.environ["PROVIDER_B_URL"],
        "timeout": 20.0
    },
}

async def call_upstream(
    route: dict[str, Any],
    payload: dict[str, Any]
) -> dict[str, Any]:
    timeout = httpx.Timeout(route["timeout"], connect=5.0)
    async with httpx.AsyncClient(timeout=timeout, http2=True) as client:
        response = await client.post(route["url"], json=payload)
        if response.status_code == 429:
            raise RuntimeError("rate_limited")
        response.raise_for_status()
        return response.json()

@app.post("/v1/agent/run")
async def run(
    req: GatewayRequest,
    x_request_id: str | None = Header(default=None)
):
    route = ROUTES.get(req.capability)
    if route is None:
        raise HTTPException(
            status_code=400,
            detail="unsupported_capability"
        )

    started = time.perf_counter()
    payload = {
        "messages": req.messages,
        "stream": req.stream
    }
    deadline = req.deadline_ms / 1000

    try:
        result = await asyncio.wait_for(
            call_upstream(route, payload),
            timeout=deadline
        )
    except asyncio.TimeoutError as exc:
        raise HTTPException(
            status_code=504,
            detail="upstream_timeout"
        ) from exc
    except RuntimeError as exc:
        raise HTTPException(
            status_code=503,
            detail=str(exc)
        ) from exc

    elapsed_ms = round(
        (time.perf_counter() - started) * 1000
    )
    return {
        "request_id": x_request_id,
        "capability": req.capability,
        "latency_ms": elapsed_ms,
        "data": result
    }

代码中的 ROUTES 只是静态演示。生产路由表应来自带版本的配置中心,更新时采用双版本加载和灰度比例。适配器要注入 Authorization、幂等键和 trace 头,但不能把上游返回的内部错误原样暴露给客户端。

对流式响应,入口应使用 StreamingResponse,在每个 delta 事件到达时刷新,并在客户端断开后取消上游任务,避免无主请求继续消耗 Token。

流式代理还必须处理半关闭与断点语义。客户端断开不一定意味着上游已取消,网关需要主动发送取消信号并等待确认;若供应商不支持取消,则记录孤儿任务并在后台回收。

SSE 事件要带单调递增序号,客户端可以检测丢包但不应盲目重放已产生副作用的工具调用。完整响应落盘后,再把最终资产引用写入幂等结果缓存。

4. 幂等、重试与降级

幂等键由客户端生成,网关把它和租户、能力、请求摘要绑定。第一次执行写入 RUNNING,重复请求读取相同任务的状态或结果。重试要带 attempt 编号和退避时间,且不能突破原始截止时间。

降级链应按能力定义,例如 reasoning.deep -> reasoning.fast -> human_review,而不是全局地把所有错误都切到同一个模型。降级结果必须标记 degraded=true,让业务决定是否继续发布。

七、端到端实战:构建一个可自治运行的“全栈内容生成” Agent 系统

在这里插入图片描述

1. 任务分解与资产契约

假设输入是“为一款新产品生成一条 30 秒介绍视频”。Agent 不应直接向模型发送这一句话,而是先形成可审计的任务图:解析受众和禁用词,生成脚本,提取镜头清单,生成分镜图,提交视频任务,等待回调,最后做字幕与时长校验。

每个节点声明输入资产、输出资产、最大预算和验收器。

brief
  -> parse_brief
  -> write_script
  -> make_storyboard (parallel shots)
  -> submit_video
  -> wait_video
  -> validate_duration_and_safety
  -> package_result

2. 一个可测试的异步执行器

下面的代码使用路由别名和协议化客户端,重点是任务状态、预算和并发控制。上游 URL、模型 ID 和凭证均由部署配置注入,代码本身不绑定某一家厂商。

import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable

@dataclass
class Job:
    job_id: str
    kind: str
    payload: dict[str, Any]
    budget: float
    status: str = "PENDING"
    result: dict[str, Any] | None = None
    errors: list[str] = field(default_factory=list)

class ModelClient:
    async def invoke(
        self,
        capability: str,
        payload: dict[str, Any],
        budget: float
    ) -> dict[str, Any]:
        # 实际实现调用统一网关,并返回 usage、资产引用和任务句柄。
        await asyncio.sleep(0)
        return {
            "capability": capability,
            "payload": payload,
            "usage": {
                "input_tokens": 0,
                "output_tokens": 0
            }
        }

class AgentRunner:
    def __init__(
        self,
        client: ModelClient,
        max_parallel: int = 3
    ):
        self.client = client
        self.limit = asyncio.Semaphore(max_parallel)
        self.jobs: dict[str, Job] = {}

    async def run_job(
        self,
        job: Job,
        capability: str,
        validator: Callable[[dict[str, Any]], bool]
    ) -> Job:
        async with self.limit:
            job.status = "RUNNING"
            try:
                result = await self.client.invoke(
                    capability,
                    job.payload,
                    job.budget
                )
                job.status = "VALIDATING"

                if not validator(result):
                    job.status = "FAILED"
                    job.errors.append("validation_failed")
                else:
                    job.result = result
                    job.status = "SUCCEEDED"

            except TimeoutError:
                job.status = "RETRYABLE"
                job.errors.append("timeout")
            except Exception as exc:
                job.status = "FAILED"
                job.errors.append(type(exc).__name__)

            self.jobs[job.job_id] = job
            return job

async def build_content_agent(
    brief: str
) -> dict[str, Any]:
    runner = AgentRunner(ModelClient())

    script = Job(
        "script",
        "text",
        {"brief": brief},
        budget=0.25
    )
    await runner.run_job(
        script,
        "reasoning.deep",
        lambda x: "usage" in x
    )

    if script.status != "SUCCEEDED":
        return {
            "status": script.status,
            "errors": script.errors
        }

    shots = [
        Job(
            f"shot-{i}",
            "image",
            {
                "script": script.result,
                "index": i
            },
            0.10
        )
        for i in range(1, 4)
    ]

    await asyncio.gather(*[
        runner.run_job(
            job,
            "vision.extract",
            lambda x: "usage" in x
        )
        for job in shots
    ])

    if any(
        job.status != "SUCCEEDED"
        for job in shots
    ):
        return {
            "status": "FAILED",
            "errors": [
                error
                for job in shots
                for error in job.errors
            ]
        }

    video = Job(
        "video",
        "video",
        {
            "shots": [
                job.result
                for job in shots
            ]
        },
        0.80
    )
    await runner.run_job(
        video,
        "video.short",
        lambda x: "usage" in x
    )

    return {
        "status": video.status,
        "script": script.result,
        "shots": [
            job.result
            for job in shots
        ],
        "video": video.result
    }

这段执行器仍然是最小实现,但已经体现了三个生产原则。

第一,状态更新和结果写入在同一个任务对象上完成,便于接入持久化;第二,分镜节点通过 gather 并行,却受到全局 semaphore 限制;第三,验证失败会进入明确的失败状态,而不是把错误内容继续传给下游。

真实系统还要为视频节点增加轮询状态、回调签名校验、资产过期时间和人工审核分支。

3. 质量门与人工接管

自治不等于无人值守。脚本节点可以检查禁用词、长度和结构化字段;分镜节点可以检查宽高比、主体一致性和版权标签;视频节点可以检查时长、音轨、字幕时间轴和安全策略。

任何一项不通过,都应生成可读的修复建议。连续两次自动修复仍失败时,转入 HUMAN_REVIEW,并附带完整 trace、输入资产和失败原因,人工只需处理局部决策。

人工接管之后也要回到同一状态机,而不是在聊天工具里完成后手工改数据库。审核界面提交结构化决策:通过、驳回、修改字段或从指定节点重跑,并携带审核人、理由和版本号。

重跑只使受影响节点及其后继失效,已经验证且输入未变化的节点可以复用。这个局部重算机制能够显著减少长链路中的重复 Token 和生成等待。

4. 运行数据与回放

每次运行保存 Prompt 模板版本、路由别名、模型配置版本、工具版本、输入资产哈希和验收结果。回放时先固定这些版本,再替换一个变量,例如只改变路由或上下文剪裁策略。这样可以把“感觉变好了”转换为可比较的实验。

对含个人数据的任务,回放集应做字段级脱敏,并设置访问审计和自动过期。

上线前应执行三类故障演练。第一类注入网络超时和 429,确认退避、截止时间和降级链正确;第二类注入格式错误和空资产,确认验证器不会让坏结果进入下游;第三类模拟 worker 崩溃和重复回调,确认租约、幂等键和乐观锁能阻止重复执行。

演练的通过条件不是“最终有结果”,而是状态可解释、成本不越界、无重复副作用且告警包含足够诊断信息。
在这里插入图片描述

八、总结与趋势:云原生 AI 中继与未来 Middleware 的形态

Agentic Workflow 的竞争点正在从单模型能力转向运行时质量。一个面向生产的多模态系统,需要用能力目录描述“能做什么”,用状态机描述“做到哪一步”,用 Token 路由决定“花多少资源”,用统一网关处理“如何调用”,再用评估和审计回答“结果是否可信”。

这些组件共同构成 LLM Middleware,而不是某个 SDK 的附属功能。

未来的中继层会更像云原生控制面:策略以代码或声明式配置发布,路由按实时质量和队列状态自适应,异步任务通过事件总线扩展,模型升级通过灰度和回放验证。

Serverless LLM 托管可以降低闲置成本,但不会消除协议、预算和数据治理问题。对开发团队而言,最值得投入的是清晰的能力契约、可重放的观测数据和可解释的失败处理。只要这些边界稳定,底层模型可以持续替换,Agent 仍能保持可测试、可审计和可演进。在这里插入图片描述

更多推荐