模型系统¶
jianmu/model 是 Jianmu 的模型系统。它不只是一个“发请求给大模型”的薄封装,而是一层面向框架内部的统一模型抽象:上承 LLM 节点、AgentRuntime 与 Swarm 预设,下接 OpenAI 兼容接口、LiteLLM 以及未来可扩展的自定义 Provider,在中间负责消息编码、工具调用参数装配、流式聚合、重试降级和用量遥测。
如果只从类名看,这一层最显眼的入口是 ModelClient;但如果从架构上看,ModelClient 只是模型系统的外观。真正需要理解的是:jianmu/model 如何把“模型调用”从 SDK 细节提升为框架能力。
这一层到底解决什么问题¶
Jianmu 把模型调用单独抽成一个系统,主要是为了解决五类问题:
- Provider 差异隔离:让上层节点不必分别处理 OpenAI、LiteLLM 或自定义后端的请求格式与返回格式。
- 调用契约统一:让流式与非流式调用、纯文本与工具调用、普通输出与结构化输出都走同一套接口。
- 可靠性治理:把重试、超时、并发槽位、同族模型降级等稳定性逻辑集中在一层。
- 消息协议转换:把
jianmu.message.Message领域模型转换为 Provider API 能接受的 payload,再把响应解码回统一消息对象。 - 可观测性注入:统一发射 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 则将这些能力整理成一个稳定的模块出口,供上层统一导入:
ModelClientModelConfigModelProviderProtocolMessageChunkinvoke_modelOpenAIProviderLiteLLMProvider
入口视角:为什么大家先看到的是 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 用量数据,它会:
- 调用
normalize_usage_dict()将prompt_tokens/input_tokens、completion_tokens/output_tokens等多命名统一为标准格式 - 若提供了
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() 中遵循相同的结构模式:
这种统一的日志埋点确保了无论使用哪个 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 的调用机制后,建议继续阅读以下页面以建立完整的调用链认知:
- LLM 节点:AgentLLMNode 与 SimpleLLMNode 的上下文构建与模型调用 — 了解 LLM 节点如何使用 ModelClient 完成从上下文构建到模型调用的完整流程
- 工具抽象:Tool 基类、@tool 装饰器与 ToolSet 工具集装配 — 理解 ModelConfig.tools 中工具 schema 的生成来源
- 遥测与可观测性:TelemetryHub、多 Sink 输出与 Span 追踪 — 深入 ModelClient 发射的
llm.usage和llm.stream.delta事件的消费端 - 上下文构建器:消息过滤、Token 预算控制与多源 Prompt 装配 — 了解调用前的消息列表是如何组装和约束的