跳转至

模型系统

jianmu/model 是 Jianmu 的模型系统。它不只是一个“发请求给大模型”的薄封装,而是一层面向框架内部的统一模型抽象:上承 LLM 节点、AgentRuntime 与 Swarm 预设,下接 OpenAI 兼容接口、LiteLLM 以及未来可扩展的自定义 Provider,在中间负责消息编码、工具调用参数装配、流式聚合、重试降级和用量遥测。

如果只从类名看,这一层最显眼的入口是 ModelClient;但如果从架构上看,ModelClient 只是模型系统的外观。真正需要理解的是:jianmu/model 如何把“模型调用”从 SDK 细节提升为框架能力。

这一层到底解决什么问题

Jianmu 把模型调用单独抽成一个系统,主要是为了解决五类问题:

  1. Provider 差异隔离:让上层节点不必分别处理 OpenAI、LiteLLM 或自定义后端的请求格式与返回格式。
  2. 调用契约统一:让流式与非流式调用、纯文本与工具调用、普通输出与结构化输出都走同一套接口。
  3. 可靠性治理:把重试、超时、并发槽位、同族模型降级等稳定性逻辑集中在一层。
  4. 消息协议转换:把 jianmu.message.Message 领域模型转换为 Provider API 能接受的 payload,再把响应解码回统一消息对象。
  5. 可观测性注入:统一发射 token usage 与 stream delta 事件,让 telemetry、日志、成本统计都能挂在标准位置。

换句话说,jianmu/model 的职责不是“替你选一个模型 SDK”,而是“把模型能力变成 Jianmu 可以组合、注入、追踪、替换的一层基础设施”。

模型系统在整体架构中的位置

从整体分层上看,模型系统位于“节点组合层”和“外部模型服务”之间:

flowchart LR
    subgraph 上层消费方
        NODES["LLM 节点 / ReAct / Plan-Execute / Swarm"]
        RUNTIME["AgentRuntime / PromptRuntime"]
    end

    subgraph jianmu.model
        CLIENT["ModelClient"]
        CONFIG["ModelConfig"]
        CODECS["Encoder / Decoder"]
        STREAM["Stream Aggregation"]
        PROVIDERS["Providers"]
    end

    subgraph 外部模型后端
        OPENAI["OpenAI-Compatible APIs"]
        LITELLM["LiteLLM Routed Models"]
        CUSTOM["Custom Providers"]
    end

    NODES --> CLIENT
    RUNTIME --> CLIENT
    CLIENT --> CONFIG
    CLIENT --> CODECS
    CLIENT --> STREAM
    CLIENT --> PROVIDERS
    PROVIDERS --> OPENAI
    PROVIDERS --> LITELLM
    PROVIDERS --> CUSTOM

这里有一个很重要的认知边界:

  • jianmu.node 决定“何时调用模型、带什么上下文调用模型”
  • jianmu.context / ContextBuilder 决定“传给模型的消息列表如何被装配”
  • jianmu.tool 决定“哪些工具 schema 会暴露给模型”
  • jianmu.model 决定“这些消息和工具如何真正变成一次稳定、可观测、可替换的模型调用”

所以它既不属于运行时,也不只是 API wrapper,而是一层独立的模型基础设施。

模块地图:jianmu/model 由哪些部分组成

当前 jianmu/model 目录可以粗分为六个职责块:

模块 文件 核心职责
协议定义 base.py 定义 ModelProviderProtocol,约束所有 Provider 的 generate() 和 stream() 签名
外观与解析 client.py ModelClient 本体、Provider 解析、注册、重试、降级逻辑
调用分发 invoke.py invoke_model() 根据 stream 标志统一分发到 generate 或 stream 路径
流式聚合 stream.py consume_stream() 逐块合并文本和 tool_calls,恢复完整 assistant 消息
编解码 codecs.py OpenAI / LiteLLM 编码器与解码器,处理 Message 与 Provider payload 的双向转换
类型定义 types.py ModelConfig 与 MessageChunk,描述一次调用和一次流式增量的标准结构

此外,providers/ 下提供当前内置的两类后端实现:

  • OpenAIProvider:面向 OpenAI 兼容 chat-completions 接口
  • LiteLLMProvider:面向 LiteLLM 的多后端路由能力

__init__.py 则将这些能力整理成一个稳定的模块出口,供上层统一导入:

  • ModelClient
  • ModelConfig
  • ModelProviderProtocol
  • MessageChunk
  • invoke_model
  • OpenAIProvider
  • LiteLLMProvider

入口视角:为什么大家先看到的是 ModelClient

虽然模型系统包含多个子模块,但调用方通常只会先接触 ModelClient。原因很简单:它把整个调用链压缩成了一个稳定入口:

  • ModelClient.resolve():按环境和偏好解析可用 Provider,并返回统一外观
  • ModelClient.resolve_provider():只拿底层 Provider 实例,适合框架注入
  • ModelClient.invoke():以统一契约执行一次模型调用

因此更准确的表述应该是:

  • ModelClient 是模型系统的门面
  • jianmu/model 才是完整的模型能力层

下面再进入 ModelClient 这条主路径,就会更容易理解它为什么能承载这么多职责。

架构概览:六个模块如何围绕 ModelClient 协作

ModelClient 并非一个孤立类,而是由六个模块协同构成的调用管道。从调用方的视角看,入口只有一个 ModelClient.invoke(),但该方法背后串联了 Provider 工厂、重试引擎、调用分发器、流消费器、编解码器和遥测中心。

flowchart TB
    subgraph 调用方
        NODE["LLM 节点 / AgentRuntime"]
    end

    subgraph ModelClient 外观
        MC["ModelClient.invoke()"]
        RESOLVE["resolve_provider() 解析"]
        REGISTRY["_REGISTRY 自定义工厂"]
        RETRY["重试循环 + 降级"]
    end

    subgraph invoke_model 分发
        IM["invoke_model()"]
    end

    subgraph 传输层
        GEN["provider.generate()"]
        STR["provider.stream()"]
        CS["consume_stream() 流聚合"]
    end

    subgraph 编解码
        ENC["Encoder 编码消息"]
        DEC["Decoder 解码响应"]
    end

    subgraph 可观测
        LOG["结构化日志"]
        TELE["TelemetryHub emit()"]
    end

    NODE --> MC
    MC --> RESOLVE --> REGISTRY
    MC --> RETRY --> IM
    IM -->|"stream=False"| GEN
    IM -->|"stream=True"| STR --> CS
    GEN --> DEC
    CS --> DEC
    IM --> TELE
    RETRY --> LOG

    style MC fill:#e1f5fe,stroke:#0288d1,stroke-width:3px
    style IM fill:#fff3e0,stroke:#f57c00,stroke-width:2px

ModelProviderProtocol:所有 Provider 的契约

ModelProviderProtocol 是一个 typing.Protocol,定义了 Jianmu 期望每个 Provider 必须实现的两个异步方法:

class ModelProviderProtocol(Protocol):
    async def generate(self, messages: List[Message], config: ModelConfig) -> Message: ...
    async def stream(self, messages: List[Message], config: ModelConfig) -> AsyncIterator[MessageChunk]: ...

选择 Protocol 而非 ABC 是刻意的架构决策:它允许 OpenAIProvider 和 LiteLLMProvider 无需共同基类即可被 ModelClient 统一消费,同时保留各自独立的 SDK 对象(AsyncOpenAI client vs litellm module)和初始化参数。类型检查器在静态期验证结构兼容性,运行时无侵入。

ModelConfig:一次调用的完整参数快照

ModelConfig 是一个 @dataclass,聚合了一次模型调用的全部参数。它设计为传输中立——同一个 ModelConfig 实例可同时用于 OpenAI 和 LiteLLM Provider,无需修改:

参数组 字段 默认值 说明
采样 temperature 0.7 采样温度
top_p 0.95 核采样截断
top_k 40 Top-k 截断
max_tokens None 最大生成 token 数
连接 timeout 120.0 超时秒数
工具调用 tools None 工具 schema 列表
tool_choice None Provider 特定工具选择指令
strict_tools False 是否强制参数匹配 schema
结构化输出 response_format None JSON mode / structured output 配置
重试降级 max_retries None None 时回退全局默认
fallback_model None 重试耗尽后降级模型
disable_fallback False 显式禁用降级
扩展 extra {} Provider 特定参数透传

with_tools() 和 with_extra() 是两个不可变辅助方法,返回新实例而非修改原对象,确保并发安全。

Provider 解析:从偏好列表到实例的决策链

ModelClient.resolve_provider() 是 Provider 的工厂入口,其解析逻辑遵循严格的优先级链路:

flowchart TD
    START["resolve_provider(name, preference, env_override)"] --> ENV["load_project_env(override=env_override)"]
    ENV --> NAME{"name 是否指定?"}
    NAME -->|是| CREATE["_create_provider(name)"]
    NAME -->|否| ORDER["按 preference 或默认 ['openai', 'litellm'] 遍历"]
    ORDER --> CHECK{"_has_credentials(provider)?"}
    CHECK -->|是| TRY_CREATE["尝试 _create_provider()"]
    TRY_CREATE -->|成功| RETURN["返回 provider 实例"]
    TRY_CREATE -->|异常| NEXT["continue 下一个"]
    CHECK -->|否| NEXT
    NEXT --> FALLBACK["遍历 _REGISTRY 中未在 order 中的自定义 provider"]
    FALLBACK -->|全部失败| ERROR["RuntimeError: No model provider configured"]
    CREATE --> RETURN

凭证检测是决策的关键:_has_credentials() 检查环境变量中是否存在对应 Provider 的 API Key。OpenAI 只需 OPENAI_API_KEY 或 API_KEY;LiteLLM 则扩展检查 GOOGLE_API_KEY、ANTHROPIC_API_KEY、GEMINI_API_KEY 等多个变量,以覆盖其多后端路由特性。

自定义 Provider 注册通过 ModelClient.register_provider(name, factory) 将工厂函数或类存入模块级 _REGISTRY 字典。已注册的 Provider 在解析时优先级低于 preference 列表中的内置 Provider,但可作为兜底被遍历。unregister_provider() 和 list_providers() 提供了完整的生命周期管理。

invoke():重试、降级与统一调用路径

ModelClient.invoke() 是本外观模式的核心方法,在单次调用中实现了重试循环 → 降级链 → 异常传播的完整可靠性保障:

flowchart TD
    INVOKE["invoke(messages, config, stream, on_text_update, trace_node)"]
    INVOKE --> MAX["max_retries = _resolve_max_retries(config)"]
    MAX --> LOOP{"attempt in 0..max_retries"}
    LOOP --> CALL["invoke_model(provider, ...)"]
    CALL -->|成功| RETURN["返回 (Message, str)"]
    CALL -->|异常| CLASSIFY{"_is_retryable_model_error(exc)?"}
    CLASSIFY -->|是 且 有剩余重试| WARN["logger.warning 记录重试"] --> LOOP
    CLASSIFY -->|否 或 重试耗尽| FALLBACK{"满足降级条件?"}
    FALLBACK -->|是| FALLBACK_CALL["以 fallback_model 调用 invoke_model"]
    FALLBACK_CALL --> RETURN
    FALLBACK -->|否| RAISE["raise last_exc"]

可重试错误分类:_is_retryable_model_error() 通过两层判断——异常类型(TimeoutError、ConnectionError)和异常消息关键词(timeout、rate limit、429、5xx 等)——来识别瞬时故障。

降级条件十分严谨,必须同时满足四项: 1. 存在 fallback_model(来自 ModelConfig.fallback_model 或全局 ModelsConfig.fallback) 2. config.disable_fallback 不为 True 3. fallback 模型与主模型属于同一 Provider 族(_is_same_provider_model_family()) 4. 最后一次异常是可重试的

Provider 族检测:_model_provider_family() 通过模型名前缀(gpt- → openai,claude- → anthropic,gemini- → google)或 / 分隔符来推断 Provider 归属,防止跨 Provider 降级(如 OpenAI → Anthropic)导致请求格式不兼容。

invoke_model:流/非流统一分发与遥测

invoke_model() 是调用链中承上启下的枢纽。上游来自 ModelClient.invoke() 的重试循环,下游分叉到 provider.generate()(非流式)或 consume_stream()(流式),最终统一返回 tuple[Message, str],消除了流与非流调用在返回类型上的差异。

sequenceDiagram
    participant MC as ModelClient.invoke
    participant IM as invoke_model
    participant P as provider
    participant CS as consume_stream
    participant T as TelemetryHub

    MC->>IM: invoke_model(provider, messages, config, stream, ...)
    alt stream=False
        IM->>P: generate(messages, config)
        P-->>IM: Message
    else stream=True
        IM->>P: stream(messages, config)
        P-->>CS: AsyncIterator[MessageChunk]
        CS-->>IM: (Message, content_str)
    end
    IM->>IM: normalize_usage_dict(metadata.usage)
    IM->>T: emit("llm.usage", {node, usage})
    IM-->>MC: (response_msg, content)

遥测发射是 invoke_model() 的固定后置步骤:每当响应中包含 token 用量数据,它会:

  1. 调用 normalize_usage_dict() 将 prompt_tokens / input_tokens、completion_tokens / output_tokens 等多命名统一为标准格式
  2. 若提供了 trace_node 参数,通过 emit("llm.usage", ...) 将用量事件推入 TelemetryHub

consume_stream:流式输出的增量聚合

consume_stream() 消费 provider.stream() 返回的 AsyncIterator[MessageChunk],将其还原为完整的 Message 对象。它解决的核心问题是流式 tool_calls 的按索引合并:

flowchart LR
    CHUNK1["chunk: {text: '我将', tool_calls: null}"] --> PARTS
    CHUNK2["chunk: {text: '调用', tool_calls: [{index:0, name:'search'}]}"] --> MERGE
    CHUNK3["chunk: {text: '搜索', tool_calls: [{index:0, arguments:'{\"q\"'}]}"] --> MERGE
    CHUNK4["chunk: {text: '', tool_calls: [{index:0, arguments:':\"hello\"}'}]}"] --> MERGE

    subgraph 聚合
        PARTS["parts: ['我将', '调用', '搜索']"]
        MERGE["indexed_calls: {0: {name:'search', arguments:'{\"q\":\"hello\"}'}}"]
    end

    PARTS --> CONTENT["content = '我将调用搜索'"]
    MERGE --> FINALIZE["finalize_stream_tool_calls()"]
    FINALIZE --> TC["tool_calls: [{name:'search', arguments:'{...}'}]"]

    CONTENT --> RESULT["Message(role='assistant', content=..., tool_calls=...)"]
    TC --> RESULT

merge_stream_tool_calls() 处理两种 tool_calls 格式: - 有索引的调用(index 字段存在):按 index 合并到 indexed_calls 字典,同名调用拼接 arguments 字符串 - 无索引的调用:直接追加到 direct_calls 列表

finalize_stream_tool_calls() 在流结束后按索引顺序重组调用列表,并对完全相同的 (name, arguments) 组合去重,防止流式传输中的重复帧导致工具重复执行。

编解码层:Jianmu Message ↔ Provider 格式的双向转换

消息格式转换是 Provider 差异的主要来源,Jianmu 通过 Encoder + Decoder 对 将其隔离在 codecs.py 中:

组件 方向 关键处理
OpenAIChatEncoder Jianmu → OpenAI 多 system 消息压缩为首条 system + 后续转 user;tool_calls 转为 {id, type:"function", function:{name, arguments}};tool 消息添加 tool_call_id
OpenAIChatDecoder OpenAI → Jianmu tc.function.{name,arguments} 展平为 {id, name, arguments};提取 model/id/finish_reason/usage 到 metadata
LiteLLMChatEncoder 继承 OpenAIChatEncoder 与 OpenAI 编码完全一致,保持语义对等
LiteLLMChatDecoder LiteLLM → Jianmu 独立实现但逻辑对称,保留独立演进空间

关键设计决策:LiteLLMChatEncoder 直接继承 OpenAIChatEncoder 而不做任何覆盖,因为 LiteLLM 接受 OpenAI 兼容的消息格式。但 LiteLLMChatDecoder 是独立类而非继承,以便未来 LiteLLM 后端出现与 OpenAI 不一致的响应格式时可以独立修改。

_normalize_system_roles() 是 Encoder 中最精妙的部分:许多 LLM API 要求 system 消息只能出现在对话开头。Jianmu 将位置 > 0 的 system 消息自动转换为 user 角色内容,并添加 [System note preserved in chronological order] 前缀,保留其语义但符合 API 约束。

Provider 实现:并发控制与结构化日志

OpenAIProvider

OpenAIProvider 封装 openai.AsyncOpenAI 客户端,核心特征包括:

  • 构造函数:接受 api_key、base_url、自定义 encoder/decoder 和 **default_kwargs。API key 优先级为显式参数 > OPENAI_API_KEY 环境变量 > API_KEY 环境变量
  • _build_tools():将 Jianmu 的工具 schema 包装为 OpenAI 的 {"type": "function", "function": {...}} 格式
  • _sanitize_strict_system_position():在已编码的消息上再做一层 system 位置净化,确保兼容 strict API
  • 并发控制:通过 asyncio.Semaphore 实现共享槽位管理,key 由 (base_url, model, limit) 三元组构成,跨调用方共享。默认最大并发为 1(可通过 ProviderLimitsConfig.openai_max_concurrency 配置),设为 0 或负数则完全跳过限制

LiteLLMProvider

LiteLLMProvider 封装 litellm 模块,独特之处在于 Provider 推断:

  • _infer_custom_provider() 按优先级扫描环境变量(LITELLM_CUSTOM_LLM_PROVIDER → OPENROUTER_API_KEY → ANTHROPIC_API_KEY → GOOGLE_API_KEY → ...),自动为无前缀的模型名推断 custom_llm_provider
  • _model_has_provider_prefix() 检查模型名是否已包含已知的 LiteLLM provider 前缀(如 openai/gpt-4、anthropic/claude-3),避免重复推断
  • 构造函数无 api_key 参数——LiteLLM 自行从环境变量读取对应后端的凭证

共同模式

两个 Provider 在 generate() 和 stream() 中遵循相同的结构模式:

获取信号量槽位 → 记录 🚀 start 日志 → 执行 SDK 调用 → 记录 ✅ done / ❌ failed 日志 → 释放信号量

这种统一的日志埋点确保了无论使用哪个 Provider,调用方可获得一致的可观测体验。

遥测集成:从调用到可观测事件

Jianmu 的遥测通过 telemetry 模块的 emit() 函数以字符串事件名 + 字典 payload 模式工作。Model 层发射两个关键事件:

事件名 发射位置 Payload 触发条件
llm.stream.delta consume_stream() {"content": chunk.text} 每个流式文本块到达时
llm.usage invoke_model() {"node": trace_node, "usage": {prompt_tokens, completion_tokens, total_tokens}} 每次调用完成且有用量数据时

llm.usage 事件的数据经过 normalize_usage_dict() 标准化,兼容 OpenAI(prompt_tokens / completion_tokens)、Google(prompt_token_count / candidates_token_count)和 Anthropic(input_tokens / output_tokens)等不同命名。标准化后统一输出 prompt_tokens / completion_tokens / total_tokens 三个 key,若 total_tokens 缺失则自动计算 prompt + completion。

这些事件被 TelemetryHub 消费后,可通过注册的 Sink(如 ConsoleStreamSink、JsonlFileSink、UsageAggregationSink)输出到不同目标。详见 遥测与可观测性:TelemetryHub、多 Sink 输出与 Span 追踪。

使用示例

最简调用(Provider 自动解析):

from jianmu.model import ModelClient, ModelConfig
from jianmu.message import Message

client = ModelClient.resolve()  # 自动按凭证选择 openai 或 litellm
response, text = await client.invoke(
    messages=[Message(role="user", content="你好")],
    config=ModelConfig(model="gpt-4.1"),
)

显式指定 Provider:

client = ModelClient.resolve(name="litellm")
response, text = await client.invoke(
    messages=[...],
    config=ModelConfig(model="anthropic/claude-3-5-sonnet-20241022"),
)

流式调用带实时回调:

def on_text(text: str):
    print(f"\r{text}", end="", flush=True)

response, full_text = await client.invoke(
    messages=[...],
    config=ModelConfig(model="gpt-4.1"),
    stream=True,
    on_text_update=on_text,
)

带重试与降级的健壮调用:

response, text = await client.invoke(
    messages=[...],
    config=ModelConfig(
        model="openai/gpt-4.1",
        max_retries=2,
        fallback_model="openai/gpt-4.1-mini",
    ),
)

阅读后续

理解 ModelClient 的调用机制后,建议继续阅读以下页面以建立完整的调用链认知: