Ollama本地Provider
运行 codex --oss --local-provider ollama -m gpt-oss:20b 后,第一条模型请求之前发生了什么? 服务能连接,模型列表却查询失败时,Codex 是否还会继续?下载进度已经显示 success,为何直接收集底层事件又会看到两次成功?这些问题都需要沿调用方、流解析器和消费者一起读,单看 OllamaClient 的方法名无法回答。
本文面向了解 Rust 借用、Result 和 async/await 的读者。可先读内置与自定义 Provider 解析,理解 Provider ID 如何进入配置;推理阶段的事件协议由 Responses 流式事件解析展开。这里把重点放在启动时的本地模型准备,读完应能追踪一次缺失模型的拉取,并定位地址、版本、分片和终止事件导致的失败。
先区分三种“可用”:HTTP 探测成功说明一个端点返回成功状态;本地模型列表命中说明服务报告了指定名称;模型推理成功还要求服务理解请求、模型支持所需能力,并产生合法的 Responses 事件。codex-ollama 负责前两层的准备过程,不包含 Ollama 服务端的权重存储、GPU 调度或 token 生成实现。
1. 启动入口
1.1 开关与模型
OSS 模式在这里是启动入口的本地 Provider 准备模式。Provider ID是配置目录中的键,例如 ollama;模型名是提交给该服务的字符串,例如 gpt-oss:20b。二者不共用命名空间。
源码文件:codex-rs/utils/cli/src/shared_options.rs
相关函数/类型:SharedCliOptions
// ...
// 三个字段分别控制模型、准备模式和 Provider ID;local-provider 本身不会开启 oss。
/// Model the agent should use.
#[arg(long, short = 'm')]
pub model: Option<String>,
/// Use open-source provider.
#[arg(long = "oss", default_value_t = false)]
pub oss: bool,
/// Specify which local provider to use (lmstudio or ollama).
/// If not specified with --oss, will use config default or show selection.
#[arg(long = "local-provider")]
pub oss_provider: Option<String>,
// ...--local-provider ollama 填充 oss_provider,--oss 才使入口进入本地准备分支。只配置 model_provider = "ollama" 可以选择推理 Provider,但不会自动令 cli.oss 变成 true,因此不能据此断言一定发生下载前检查。
源码文件:codex-rs/core/src/config/mod.rs
相关函数/类型:resolve_oss_provider
// 选择发生在正式 Config 构造前;此函数只取 ID,不探测服务器。
pub fn resolve_oss_provider(
explicit_provider: Option<&str>,
config_toml: &ConfigToml,
) -> Option<String> {
if let Some(provider) = explicit_provider {
// Explicit provider specified (e.g., via --local-provider)
Some(provider.to_string())
} else {
config_toml.oss_provider.clone()
}
}解析函数只决定 ID:显式 --local-provider 优先,否则取加载后的 ConfigToml.oss_provider。它既不检查端口,也不验证模型是否存在。TUI 与非交互的 codex exec 在没有 ID 时走不同分支:前者尝试检测本地服务并让用户选择,后者直接返回缺少默认 OSS Provider 的错误。
TUI 的兜底检测在 codex-rs/tui/src/oss_selection.rs 的 detect_oss_provider:先探测 LM Studio,再探测 Ollama;只有一方运行时自动选择,否则打开选择界面。check_port_status 对固定的 http://localhost:1234 与 http://localhost:11434 根地址执行 GET,单次请求超时为 2 秒。这个探测有意直连,不读取自定义 Provider endpoint,也不替代后面使用完整配置的客户端探测。选择界面返回 __CANCELLED__ 时,启动函数返回取消错误。
源码文件:codex-rs/tui/src/startup_orchestration.rs
相关函数/类型:run_main_inner
// ...
// 默认模型作为 ConfigOverrides 输入,显式 -m 的优先级更高。
let model = if let Some(model) = &cli.model {
Some(model.clone())
} else if cli.oss {
// Use the provider from model_provider_override
model_provider_override
.as_ref()
.and_then(|provider_id| get_default_model_for_oss_provider(provider_id))
.map(std::borrow::ToOwned::to_owned)
} else {
None // No model specified, will use the default.
};
// ...这段代码中的 model 随后放进 ConfigOverrides 再构造正式 Config。因而 --oss 且没有 -m 时,入口已经按 Provider 选择了默认模型,不能把它理解成 ensure_oss_ready 最后才决定模型。Ollama 对应的默认值来自 codex-rs/ollama/src/lib.rs 的 DEFAULT_OSS_MODEL = "gpt-oss:20b"。
1.2 检查顺序
TUI 的 run_main_inner 与 exec 的启动路径最终调用共用的 ensure_oss_provider_ready。共用函数把本地产品的差异保留在各自 crate 内,并严格按下面的代码顺序等待。
源码文件:codex-rs/utils/oss/src/lib.rs
相关函数/类型:ensure_oss_provider_ready
pub async fn ensure_oss_provider_ready(
provider_id: &str,
config: &Config,
) -> Result<(), std::io::Error> {
match provider_id {
LMSTUDIO_OSS_PROVIDER_ID => {
codex_lmstudio::ensure_oss_ready(config)
.await
.map_err(|e| std::io::Error::other(format!("OSS setup failed: {e}")))?;
}
OLLAMA_OSS_PROVIDER_ID => {
// 一次构造后的同一客户端借给两个后续检查;版本失败会提前返回。
let client = codex_ollama::OllamaClient::try_from_oss_provider(config).await?;
codex_ollama::ensure_responses_supported(&client).await?;
codex_ollama::ensure_oss_ready(config, &client)
.await
.map_err(|e| std::io::Error::other(format!("OSS setup failed: {e}")))?;
}
_ => {
// Unknown provider, skip setup
}
}
Ok(())
}Ollama 分支只有一次客户端构造。后续两个函数借用同一实例,先检查 Responses 版本,再判断模型是否存在。前两步失败直接传播;模型准备失败则额外包上 OSS setup failed。顺序的作用是避免在已知不支持所需协议的服务上先进行大体积下载。
下图只画已经选定 ollama 的启动路径。蓝色入口负责等待,深蓝色代表配置,绿色代表本地服务准备;推理请求从准备返回之后的会话执行进入另一条链路。
图中的“准备返回”只代表这段预检逻辑允许启动继续。模型列表查询失败时也可能走到这里;后文会跟踪这一降级分支。图中 /v1/models 是兼容地址对应的探测接口,后面的版本、列表、拉取仍使用原生 /api/* 接口。
2. 地址与传输
2.1 内置地址
内置 Provider 定义在 codex-rs/model-provider-info/src/lib.rs 的 built_in_model_providers。Ollama 的默认端口为 11434,wire API 传入 WireApi::Responses。地址生成需要同时考虑完整 URL 与端口两个环境变量。
源码文件:codex-rs/model-provider-info/src/lib.rs
相关函数/类型:create_oss_provider
pub fn create_oss_provider(default_provider_port: u16, wire_api: WireApi) -> ModelProviderInfo {
// These CODEX_OSS_ environment variables are experimental: we may
// switch to reading values from config.toml instead.
let default_codex_oss_base_url = format!(
"http://localhost:{codex_oss_port}/v1",
codex_oss_port = std::env::var("CODEX_OSS_PORT")
.ok()
.filter(|value| !value.trim().is_empty())
.and_then(|value| value.parse::<u16>().ok())
.unwrap_or(default_provider_port)
);
// 完整地址优先于端口拼装;这里仅检查空白,不负责 URL 结构校验。
let codex_oss_base_url = std::env::var("CODEX_OSS_BASE_URL")
.ok()
.filter(|v| !v.trim().is_empty())
.unwrap_or(default_codex_oss_base_url);
create_oss_provider_with_base_url(&codex_oss_base_url, wire_api)
}完整的 CODEX_OSS_BASE_URL 优先;没有有效的非空值时,才使用 CODEX_OSS_PORT 或默认端口生成 http://localhost:<port>/v1。这两个环境变量在源码中标为实验性,而且是本地 Provider 共用的输入,不是只影响 Ollama 的设置。
注意 filter 只用 trim() 判断是否为空,没有把修剪后的字符串写回。端口随后对原字符串执行 parse::<u16>(),带首尾空格的数字不能视为已被规范化。最终完整地址同样不能依赖这里替你清理前后空白。
源码文件:codex-rs/model-provider-info/src/lib.rs
相关函数/类型:create_oss_provider_with_base_url
pub fn create_oss_provider_with_base_url(base_url: &str, wire_api: WireApi) -> ModelProviderInfo {
ModelProviderInfo {
name: "gpt-oss".into(),
base_url: Some(base_url.into()),
env_key: None,
env_key_instructions: None,
experimental_bearer_token: None,
auth: None,
aws: None,
wire_api,
query_params: None,
http_headers: None,
env_http_headers: None,
request_max_retries: None,
stream_max_retries: None,
stream_idle_timeout_ms: None,
websocket_connect_timeout_ms: None,
// 本地默认不要求 OpenAI 认证,也不启用 WebSocket;模型请求仍使用传入的 wire_api。
requires_openai_auth: false,
supports_websockets: false,
supports_standalone_web_search: false,
}
}默认对象没有 env key、Bearer token、AWS 认证、额外 Header 或 WebSocket 支持。wire_api 来自调用方,并不是这个函数根据 URL 猜出来的。requires_openai_auth: false 也不表示任何服务都无需认证;它只描述这份内置配置的默认策略。
一个容易误读的地方是 OllamaClient 构造函数里“读取 Config 以包含用户覆盖”的注释。当前目录合并代码对非 Bedrock Provider 实际采用下面的插入语义。
源码文件:codex-rs/model-provider-info/src/lib.rs
相关函数/类型:merge_configured_model_providers
// ...
// 这是非 Bedrock 分支;相同 key 已存在时保留内置值,只有新 key 才插入。
model_providers.entry(key).or_insert(provider);
// ...HashMap::entry(key).or_insert(provider) 只填充不存在的键。ollama 已在内置目录中,因此同名 [model_providers.ollama] 的 base_url 不能据此替换内置值。修改地址时应先核对环境变量生成路径;如果新建另一个 Provider ID,通用 OSS 准备函数的匹配分支又不会把这个新 ID 当作 ollama 自动管理。配置能够解析,与启动 helper 能够识别,是两项独立条件。
命令示例:
CODEX_OSS_BASE_URL=http://localhost:11434/v1 \
codex --oss --local-provider ollama -m gpt-oss:20b这条命令的环境变量只作用于该次进程启动。模型缺失时会触发真实下载;它用于说明实际入口参数,后面的只读诊断命令则不会拉取模型。
2.2 两组接口
源码文件:codex-rs/ollama/src/url.rs
相关函数/类型:is_openai_compatible_base_url / base_url_to_host_root
pub(crate) fn is_openai_compatible_base_url(base_url: &str) -> bool {
base_url.trim_end_matches('/').ends_with("/v1")
}
/// Convert a provider base_url into the native Ollama host root.
/// For example, "http://localhost:11434/v1" -> "http://localhost:11434".
pub fn base_url_to_host_root(base_url: &str) -> String {
// 字符串后缀规则不解析 query;代理前缀需要与最终原生路径一起检查。
let trimmed = base_url.trim_end_matches('/');
if trimmed.ends_with("/v1") {
trimmed
.trim_end_matches("/v1")
.trim_end_matches('/')
.to_string()
} else {
trimmed.to_string()
}
}兼容根地址是以 /v1 结尾的 Provider 地址;host root是删除该后缀后用于拼接原生接口的根。这里没有 URL parser,只有字符串后缀操作。
输入 base_url | 兼容模式 | 得到的 host_root | 需要注意的结果 |
|---|---|---|---|
http://host/v1/ | true | http://host | 尾斜杠先去掉 |
http://host/prefix/v1 | true | http://host/prefix | 原生接口也带 /prefix |
http://host/v1?tenant=x | false | 原字符串 | query 阻止 /v1 后缀匹配 |
http://host/v1/v1 | true | http://host | trim_end_matches 会重复移除匹配后缀 |
http://host | false | 原字符串 | 使用原生探测接口 |
第三行尤其适合排查反向代理问题:后续字符串拼接得到 .../v1?tenant=x/api/tags,其中 /api/tags 落进 query,而不是期望的路径。不要仅凭 URL “包含 v1” 判断客户端处于兼容模式。
下面的分流图展示同一份 Provider 地址如何服务不同请求。图中绿色节点属于 Ollama 管理客户端,模型推理的 ResponsesClient 走通用模型请求路径。
兼容探测成功不意味着服务具有 Ollama 管理接口。如果一个网关只代理 /v1/models 和 /v1/responses,版本接口可能被当作“未知版本”放行,但 /api/tags 返回非成功后会触发 /api/pull;仅代理 OpenAI 兼容接口不足以承诺这条自动准备流程可用。
2.3 客户端所有权
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient
// 持久成员只有传输池和地址信息;下载缓冲与进度不保存在客户端中。
pub struct OllamaClient {
client: RouteAwareClientPool,
host_root: String,
uses_openai_compat: bool,
}客户端长期持有的只有连接池、派生根地址和探测模式。这里没有模型目录缓存、活跃下载表、健康状态锁或取消 token。一次拉取的字节缓冲在返回的流内部,一次拉取的进度状态在独立 reporter 中。
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::try_from_oss_provider
pub async fn try_from_oss_provider(config: &Config) -> io::Result<Self> {
// Note that we must look up the provider from the Config to ensure that
// any overrides the user has in their config.toml are taken into
// account.
// 读取已经合并的目录;注释提到配置覆盖,实际允许范围仍由目录合并函数决定。
let provider = config
.model_providers
.get(OLLAMA_OSS_PROVIDER_ID)
.ok_or_else(|| {
io::Error::new(
io::ErrorKind::NotFound,
format!("Built-in provider {OLLAMA_OSS_PROVIDER_ID} not found",),
)
})?;
Self::try_from_provider(provider, config.http_client_factory()).await
}读取配置目录使传输工厂与地址由同一份正式 Config 决定。找不到内置键时返回 NotFound;客户端没有偷偷补回另一份默认配置。下一层要求 Provider 一定有 base_url,通过 expect 表达内部构造前提,而不是把缺失地址转换成普通连接错误。
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::try_from_provider
pub(crate) async fn try_from_provider(
provider: &ModelProviderInfo,
http_client_factory: HttpClientFactory,
) -> io::Result<Self> {
#![expect(clippy::expect_used)]
let base_url = provider
.base_url
.as_ref()
.expect("oss provider must have a base_url");
let uses_openai_compat = is_openai_compatible_base_url(base_url);
let host_root = base_url_to_host_root(base_url);
// 5 秒只限制建连;Other 路由保留配置的代理策略。
let client = RouteAwareClientPool::with_connect_timeout(
http_client_factory,
ClientRouteClass::Other,
OLLAMA_CONNECTION_TIMEOUT,
)
.with_legacy_custom_ca_fallback();
let client = Self {
client,
host_root,
uses_openai_compat,
};
// 探测成功后才把 owner 返回;失败不会生成一个可用客户端句柄。
client.probe_server().await?;
Ok(client)
}ClientRouteClass::Other 使请求参与配置的路由选择;它不等价于强制绕过系统代理。5 秒常量传给 with_connect_timeout,只限制连接建立,不是整个下载、每个流分片或模型准备阶段的总时限。此管理路径也没有读取 Provider 的 Responses 重试次数与 SSE 空闲超时。
with_legacy_custom_ca_fallback 保留旧行为:默认传输代理策略下,自定义 CA 构造失败可退回系统根证书;系统代理策略下仍传播构造错误。底层细节由 HTTP Client 路由与中间件说明。它与 TUI 前面固定 loopback 端口的直连探测属于不同阶段。
下面用对象图区分配置借用、客户端成员与一次拉取的临时状态。LineBuffer 由流拥有,因此没有画成 OllamaClient 的成员。
&mut dyn PullProgressReporter 在高层拉取期间被独占借用,报告器更新不需要另加锁。图中的依赖关系也说明,调用 pull_model_stream 并不会自动创建一个后台下载管理器;消费者必须继续轮询返回的流。
3. 就绪判定
3.1 可达性探测
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::probe_server
async fn probe_server(&self) -> io::Result<()> {
let url = if self.uses_openai_compat {
format!("{}/v1/models", self.host_root.trim_end_matches('/'))
} else {
format!("{}/api/tags", self.host_root.trim_end_matches('/'))
};
let resp = self
.client
.get(url)
.send()
.await
.map_err(|error| match error {
// 传输初始化错误保留原因;网络/HTTP 错误才归并成启动 Ollama 的提示。
RouteAwareRequestError::Route(error) => {
tracing::warn!(error = %error, "Failed to initialize Ollama HTTP transport");
io::Error::other(error)
}
error => {
tracing::warn!(error = ?error, "Failed to connect to Ollama server");
io::Error::other(OLLAMA_CONNECTION_ERROR)
}
})?;
if resp.status().is_success() {
Ok(())
} else {
tracing::warn!(
"Failed to probe server at {}: HTTP {}",
self.host_root,
resp.status()
);
Err(io::Error::other(OLLAMA_CONNECTION_ERROR))
}
}探测只检查 HTTP 成功状态,不读取模型列表的 JSON。传输路由初始化失败保留内部原因,其余请求错误和非成功状态统一映射为“没有运行的 Ollama 服务”提示。因此出现这句提示时,既可能是没有进程监听,也可能是路由已通但探测路径返回 401 或 404。先确认真实请求地址,比直接重装 Ollama 更能缩小范围。
3.2 版本门槛
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::fetch_version
pub async fn fetch_version(&self) -> io::Result<Option<Version>> {
let version_url = format!("{}/api/version", self.host_root.trim_end_matches('/'));
let resp = self
.client
.get(version_url)
.send()
.await
.map_err(io::Error::other)?;
// 未知版本与请求/JSON 失败具有不同的 Result 语义。
if !resp.status().is_success() {
return Ok(None);
}
let val = resp.json::<JsonValue>().await.map_err(io::Error::other)?;
let Some(version_str) = val.get("version").and_then(|v| v.as_str()).map(str::trim) else {
return Ok(None);
};
let normalized = version_str.trim_start_matches('v');
match Version::parse(normalized) {
Ok(version) => Ok(Some(version)),
Err(err) => {
tracing::warn!("Failed to parse Ollama version `{version_str}`: {err}");
Ok(None)
}
}
}返回类型 io::Result<Option<Version>> 表达三种情况:Some 是已解析版本,None 是版本不可用,Err 是本次查询失败。HTTP 非成功、缺少字符串字段以及 semver 解析失败都落入 None;网络错误和 JSON 解码错误落入 Err。所以“版本不可解析就放行”必须限于版本字符串解析,不能扩大到任意响应格式损坏。
源码文件:codex-rs/ollama/src/lib.rs
相关函数/类型:min_responses_version / supports_responses / ensure_responses_supported
fn min_responses_version() -> Version {
Version::new(0, 13, 4)
}
fn supports_responses(version: &Version) -> bool {
*version == Version::new(0, 0, 0) || *version >= min_responses_version()
}
/// Ensure the running Ollama server is new enough to support the Responses API.
///
/// Returns `Ok(())` when the version endpoint is missing or unparsable.
pub async fn ensure_responses_supported(client: &OllamaClient) -> std::io::Result<()> {
// None 采用宽松放行;Err 会被 ? 直接传回启动入口。
let Some(version) = client.fetch_version().await? else {
return Ok(());
};
if supports_responses(&version) {
return Ok(());
}
let min = min_responses_version();
Err(std::io::Error::other(format!(
"Ollama {version} is too old. Codex requires Ollama {min} or newer."
)))
}已知版本的门槛为 0.13.4,另允许精确的 0.0.0 开发版本。比较用的是 semver::Version 的实际相等与排序,而不是拆出三个数字自行比较。0.13.4-rc.1 排在正式版之前;0.0.0+dev 也不满足代码中的精确相等条件,不能把所有“零开头的开发构建”都视为特例。
/api/version 输入 | 查询结果 | 启动动作 |
|---|---|---|
{"version":"0.13.3"} | Some | 版本过旧,停止 |
{"version":"0.13.4"} | Some | 继续 |
{"version":"0.13.4-rc.1"} | Some | 低于正式门槛,停止 |
{"version":"0.0.0"} | Some | 开发版本特例,继续 |
{"version":"0.0.0+dev"} | Some | 未命中特例且低于门槛,停止 |
{"version":" v0.14.1 "} | Some | 修剪空白和 v 后继续 |
{"version":"unknown"} 或 {} | None | 继续,保留未知边界 |
| HTTP 404 | None | 继续 |
HTTP 200,响应体为 not-json | Err | 立即停止 |
这些结果可以由调用 ensure_responses_supported 的 mock 实验复现。放行未知版本是一项兼容策略,不能证明该服务真的支持 Responses;真正的协议错误仍会在首次推理时出现。
3.3 名称匹配
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::fetch_models
pub async fn fetch_models(&self) -> io::Result<Vec<String>> {
let tags_url = format!("{}/api/tags", self.host_root.trim_end_matches('/'));
let resp = self
.client
.get(tags_url)
.send()
.await
.map_err(io::Error::other)?;
// HTTP 非成功降为已知空列表,调用方可能因此发起拉取。
if !resp.status().is_success() {
return Ok(Vec::new());
}
// JSON 解码失败走 Err;它与上面的空列表会触发不同动作。
let val = resp.json::<JsonValue>().await.map_err(io::Error::other)?;
let names = val
.get("models")
.and_then(|m| m.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.get("name").and_then(|n| n.as_str()))
.map(str::to_string)
.collect::<Vec<_>>()
})
.unwrap_or_default();
Ok(names)
}模型来自原生 models[].name。这里过滤掉非字符串名称,缺少数组时给出空列表,也没有做去重、补 :latest、大小写归一化或别名解析。返回顺序保持服务数组中的顺序,但后续存在性检查只关心精确字符串相等。
源码文件:codex-rs/ollama/src/lib.rs
相关函数/类型:ensure_oss_ready
pub async fn ensure_oss_ready(config: &Config, client: &OllamaClient) -> std::io::Result<()> {
// Only download when the requested model is the default OSS model (or when -m is not provided).
let model = match config.model.as_ref() {
Some(model) => model,
None => DEFAULT_OSS_MODEL,
};
// If the model is not present locally, pull it.
match client.fetch_models().await {
Ok(models) => {
// 实际分支比较所有请求模型;没有“仅默认模型可下载”的限制。
if !models.iter().any(|m| m == model) {
let mut reporter = crate::CliProgressReporter::new();
client.pull_with_reporter(model, &mut reporter).await?;
}
}
// 查询失败仅记录警告;此处放行不代表模型存在。
Err(err) => {
// Not fatal; higher layers may still proceed and surface errors later.
tracing::warn!("Failed to query local models from Ollama: {}.", err);
}
}
Ok(())
}函数顶部原注释写着“只下载默认 OSS 模型”,实际条件却只是 !models.iter().any(|m| m == model)。对照 model 的来源可知,显式 -m lesson:1 同样会在未命中时触发拉取。阅读源码时,过时注释应当由分支条件与调用方纠正。
更重要的是 HTTP 状态与解码错误的分流:/api/tags 返回 503 被转为 Ok([]),因此会尝试下载;返回 200 但 JSON 损坏得到 Err,只记录 warning 并继续启动。ensure_oss_ready 并不承诺把所有“不知道是否存在”的情况都判为失败。
| 列表响应 | fetch_models | ensure_oss_ready 后续行为 |
|---|---|---|
有完全相同的 name | Ok(models) | 不拉取 |
只有 lesson:latest,请求 lesson | Ok(models) | 拉取 lesson |
| 非成功 HTTP 状态 | Ok([]) | 尝试拉取 |
成功 JSON 却只有 data[].id | Ok([]) | 尝试拉取 |
| 200 且 JSON 无法解码 | Err | warning 后放行 |
| 列表网络失败 | Err | warning 后放行 |
| 拉取明确失败 | 拉取函数 Err | 停止,外层补 OSS setup failed |
这张表可以作为断点选择依据:如果根本没有 /api/pull 请求,应先区分“列表已命中”与“列表查询失败被放行”,两者都会从这个函数返回 Ok(())。
4. 拉取协议
4.1 事件与观察者
NDJSON是一行一个 JSON 值的流格式。Ollama 拉取使用这种格式;模型推理 Responses 使用 SSE。二者的拆包规则和完成信号不同,不能共用对 response.completed 的理解。
源码文件:codex-rs/ollama/src/pull.rs
相关函数/类型:PullEvent / PullProgressReporter
#[derive(Debug, Clone)]
pub enum PullEvent {
/// A human-readable status message (e.g., "verifying", "writing").
Status(String),
/// Byte-level progress update for a specific layer digest.
// total/completed 可分别更新;digest 标识层而非本次网络块。
ChunkProgress {
digest: String,
total: Option<u64>,
completed: Option<u64>,
},
/// The pull finished successfully.
Success,
/// Error event with a message.
Error(String),
}
/// A simple observer for pull progress events. Implementations decide how to
/// render progress (CLI, TUI, logs, ...).
pub trait PullProgressReporter {
// 报告器也能失败;调用方通过 ? 将输出错误传播为准备失败。
fn on_event(&mut self, event: &PullEvent) -> io::Result<()>;
}PullEvent 是从服务响应中投影出来的事件,不是互斥的下载状态机:一条 JSON 可以产生 Status 与 ChunkProgress,也可以产生 Status 与 Success。ChunkProgress 中 total 与 completed 各自为 Option<u64>,允许服务分两条消息更新同一个 layer。
digest在此只作为层的键。客户端不验证 SHA-256,也不依据 digest 校验下载文件。模型文件由外部 Ollama 服务管理,Codex 只接收管理接口的进度;“有 digest 字段”不能外推成 Codex 实现了权重完整性验证。
4.2 请求与流
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::pull_model_stream
// ...
pub async fn pull_model_stream(
&self,
model: &str,
) -> io::Result<BoxStream<'static, PullEvent>> {
let url = format!("{}/api/pull", self.host_root.trim_end_matches('/'));
// 这里发送 HTTP;返回的 BoxStream 在随后被轮询时才继续解码响应体。
let resp = self
.client
.post(url)
.json(&serde_json::json!({"model": model, "stream": true}))
.send()
.await
.map_err(io::Error::other)?;
if !resp.status().is_success() {
return Err(io::Error::other(format!(
"failed to start pull: HTTP {}",
resp.status()
)));
}
// 字节流与缓冲由返回的异步流捕获,不再借用 self。
let mut stream = resp.bytes_stream();
let mut buf = LineBuffer::default();
let _pending: VecDeque<PullEvent> = VecDeque::new();
// ...高层传入的模型名被原样写进 {"model":...,"stream":true}。非成功 HTTP 状态在创建流之前就返回错误;HTTP 200 则只能说明服务接受了这一响应通道,后续行里仍可能携带 error。
BoxStream<'static, PullEvent> 表示返回流不再借用 &self。响应体流与 LineBuffer 被异步生成器拥有,调用者继续持有客户端不是流读取的借用前提。_pending 在当前函数中只初始化,没有任何入队或消费,不能依据这个变量推导出一个事件队列调度机制。
这也划清了内存归属:Codex 缓冲的是 NDJSON 进度字节,不是模型权重。HTTP chunk 可以同时包含多行,也可能仅包含一行的一部分,接下来必须先恢复行边界。
5. 拆行算法
5.1 扫描游标
源码文件:codex-rs/ollama/src/line_buffer.rs
相关函数/类型:LineBuffer
use bytes::BytesMut;
use memchr::memchr;
#[derive(Default)]
#[cfg_attr(test, derive(Debug, PartialEq, Eq))]
pub(crate) struct LineBuffer {
bytes: BytesMut,
/// Prefix already scanned and known not to contain a newline.
// 不变量:该前缀已经确认没有换行,后续扫描可以跳过它。
scanned_len: usize,
}
impl LineBuffer {
pub(crate) fn extend_from_slice(&mut self, bytes: &[u8]) {
self.bytes.extend_from_slice(bytes);
}
pub(crate) fn take_line(&mut self) -> Option<BytesMut> {
// 只扫描新增后缀;没有换行就推进游标,保留已有字节。
let Some(relative_index) = memchr(b'\n', &self.bytes[self.scanned_len..]) else {
self.scanned_len = self.bytes.len();
return None;
};
let newline_index = self.scanned_len + relative_index;
let line = self.bytes.split_to(newline_index + 1);
// split_to 改变了缓冲起点,必须重置相对游标。
self.scanned_len = 0;
Some(line)
}
}bytes 保存尚未取走的数据;scanned_len 记录其中已经确认没有换行的前缀长度。不变量为 0 <= scanned_len <= bytes.len(),且前缀 bytes[..scanned_len] 不含 \n。新增字节只会追加,不会改变已扫描前缀,因此没有换行时可以直接把游标推进到当前末尾。
如果每到一个小 chunk 都从缓冲开头重扫,长为 N、按单字节到达的一行要比较约 1 + 2 + ... + N 个字节。这里把已扫描区域排除后,每个字节只需为换行判断扫描一次;这是扫描工作的摊销线性结论,不包括内存分配、JSON 解析和后续事件处理成本。
找到换行后 split_to(newline_index + 1) 取走包括换行的前段。剩余缓冲的索引原点已经改变,所以 scanned_len 必须重置为 0。后面的剩余字节可能含有另一条完整行,调用方的内层 while 会继续取出。
下图表示扫描和消费关系;“等待更多字节”是 take_line 返回 None 后由外层网络流驱动的行为,没有隐藏的唤醒任务。
图中的消费者决定是否再次轮询。若上层收到成功后直接返回,图中剩余的拆行循环就不会继续执行。缓冲器没有最大行长限制:一个始终不发换行的对端仍可让它持续增长,游标优化只解决重复扫描,不能替代输入大小限制。
5.2 字节与字符
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::pull_model_stream
// ...
// Using an async stream adaptor backed by unfold-like manual loop.
let s = async_stream::stream! {
while let Some(chunk) = stream.next().await {
match chunk {
Ok(bytes) => {
buf.extend_from_slice(&bytes);
// 先按字节找完整行,再做 UTF-8 解码,允许多字节字符跨 HTTP chunk。
while let Some(line) = buf.take_line() {
if let Ok(text) = std::str::from_utf8(&line) {
let text = text.trim();
if text.is_empty() { continue; }
if let Ok(value) = serde_json::from_str::<JsonValue>(text) {
// yield 会暂停;消费者可以在继续检查 error/status 之前就结束读取。
for ev in pull_events_from_value(&value) { yield ev; }
if let Some(err_msg) = value.get("error").and_then(|e| e.as_str()) {
yield PullEvent::Error(err_msg.to_string());
return;
}
if let Some(status) = value.get("status").and_then(|s| s.as_str())
&& status == "success" { yield PullEvent::Success; return; }
}
}
}
}
// 网络读取错误在此仅结束流;高层负责把缺少 Success 判为失败。
Err(_) => {
// Connection error: end the stream.
return;
}
}
}
};
Ok(Box::pin(s))
}
// ...解码顺序是“网络字节 → 完整行 → UTF-8 → 去首尾空白 → JSON → 事件”。因此一个中文字符的三个 UTF-8 字节分散在不同 chunk 中仍可正确恢复;若在每个 chunk 刚到达时就转换为字符串,合法的跨块字符会被误判为损坏。
同时也要看到当前容错范围:空行、UTF-8 非法行和 JSON 非法行都被跳过,没有转换成 PullEvent::Error。如果后来出现完整的成功行,高层仍可成功;如果整个响应只有损坏行,高层最终因为没有消费到成功事件而失败。
EOF 时没有“把缓冲剩余部分当作最后一行”的分支。即使最后的字节恰好是完整的 {"status":"success"},缺少结尾换行也不会被解析。这是 NDJSON 帧边界与 JSON 语法边界的差别:JSON 内容完整,不等于流解析器认定帧已结束。
6. 终止语义
6.1 一行多个事件
源码文件:codex-rs/ollama/src/parser.rs
相关函数/类型:pull_events_from_value
pub(crate) fn pull_events_from_value(value: &JsonValue) -> Vec<PullEvent> {
let mut events = Vec::new();
if let Some(status) = value.get("status").and_then(|s| s.as_str()) {
// status 先生成文本事件;success 再额外生成终止事件。
events.push(PullEvent::Status(status.to_string()));
if status == "success" {
events.push(PullEvent::Success);
}
}
let digest = value
.get("digest")
.and_then(|d| d.as_str())
.unwrap_or("")
.to_string();
// 只接受非负整数;未知字段不会阻止识别其他事件。
let total = value.get("total").and_then(JsonValue::as_u64);
let completed = value.get("completed").and_then(JsonValue::as_u64);
if total.is_some() || completed.is_some() {
events.push(PullEvent::ChunkProgress {
digest,
total,
completed,
});
}
events
}纯解析函数先生成状态事件,再生成进度事件。它不识别 error 字段;错误提取在外层流中完成。因此 {"error":"missing-model"} 经过这个函数得到空向量,外层随后才产出 PullEvent::Error。
只有 total 或 completed 至少一个能按 as_u64 读取时,才会产生进度事件。负数、字符串和小数不被接受。缺少 digest 时使用空字符串,因此两条都没有 digest 的进度消息会在 reporter 里共享同一个键;这层没有依据模型名补全层身份。
将解析器与外层流放在一起,直接收集 pull_model_stream(...).collect::<Vec<_>>() 可得到下列结果。
| 一条完整 JSON 行 | 底层流产生的事件顺序 |
|---|---|
{"status":"verifying"} | Status("verifying") |
{"digest":"a","completed":42} | ChunkProgress(a, None, Some(42)) |
{"status":"success"} | Status("success") → Success → Success |
{"error":"missing-model"} | Error("missing-model") |
{"status":"success","error":"bad"} | Status("success") → Success → Error("bad") |
{"status":"complete"} | Status("complete"),随后 EOF |
第三行的重复来自两处代码:pull_events_from_value 已经为 success 加入一次 Success,外层在遍历解析结果后又检查 status 并产出一次。不要把它在讲解中“修正”为理想化的一次通知;自定义底层消费者确实会观察到重复。
6.2 高层提前返回
源码文件:codex-rs/ollama/src/client.rs
相关函数/类型:OllamaClient::pull_with_reporter
pub async fn pull_with_reporter(
&self,
model: &str,
reporter: &mut dyn PullProgressReporter,
) -> io::Result<()> {
reporter.on_event(&PullEvent::Status(format!("Pulling model {model}...")))?;
let mut stream = self.pull_model_stream(model).await?;
while let Some(event) = stream.next().await {
// 先通知报告器,再解释事件;报告器失败可以抢先结束调用。
reporter.on_event(&event)?;
match event {
// 收到第一次 Success 就返回,余下流会被丢弃。
PullEvent::Success => {
return Ok(());
}
PullEvent::Error(err) => {
// Empirically, ollama returns a 200 OK response even when
// the output stream includes an error message. Verify with:
//
// `curl -i http://localhost:11434/api/pull -d '{ "model": "foobarbaz" }'`
//
// As such, we have to check the event stream, not the
// HTTP response status, to determine whether to return Err.
return Err(io::Error::other(format!("Pull failed: {err}")));
}
PullEvent::ChunkProgress { .. } | PullEvent::Status(_) => {
continue;
}
}
}
Err(io::Error::other(
"Pull stream ended unexpectedly without success.",
))
}高层没有收集整个流,而是边读边处理。它先发出开始状态,再调用 HTTP 拉取;每次收到事件先交给 reporter,然后根据 Success 或 Error 决定是否返回。正常 EOF、网络读取错误引起的 EOF、没有换行的尾行、只有 complete 状态,最终都汇合到“未观察到成功”的错误。
yield 是理解重复成功为何不总出现在 UI 上的关键:生成器在交出事件时暂停。高层收到第一次 Success 后立即返回,返回时丢弃流,因此第二次成功所在的代码还没获得下一次轮询机会。底层产生什么,与特定消费者实际读到什么,必须分开描述。
下面从高层调用者的视角画出这一暂停点。PullProgressReporter 是被借用的观察者,图中没有后台事件转发 channel。
矛盾响应 {"status":"success","error":"bad"} 更能检验对顺序的理解。外层虽然把 error 检查写在它自己的最终 status 检查之前,但解析器早已 yield Success。高层会在读到后面的 Error 前返回成功。这个结果描述客户端面对矛盾输入的优先次序,不说明正常 Ollama 服务会发出这种响应。
6.3 输出失败
reporter.on_event(...)? 是有控制效果的调用,并非永远不失败的日志旁路。第一条“开始拉取”状态输出失败时,HTTP POST 尚未发出;已经收到服务事件后报告器失败,则返回该输出错误并丢弃正在消费的流。连 Success 都是先经过 reporter,若此时 stderr 写入失败,拉取函数仍返回错误。
可以给自定义 reporter 加入“第 N 次回调返回 io::Error”的注入点,观察 mock 服务接收到的请求数。N=1 时为 0 次 POST;N=2 时为 1 次 POST。这个反例说明,函数错误不仅来自 Ollama,也可能来自本地进度呈现。
底层响应流所有权随 drop 释放,但这层没有向 /api/pull 再发送取消 RPC,没有等待远端下载退出,也没有删除远端已取得的模型层。客户端停止读取与服务端是否保留、继续或清理下载是两个不同问题;后者需要 Ollama 服务端源码或运行实验,不能从 Codex 的 drop 推导。
7. 进度聚合
7.1 每层快照
源码文件:codex-rs/ollama/src/pull.rs
相关函数/类型:CliProgressReporter
// 这份可变进度只属于一次拉取;按 digest 保存最近值和相邻采样点。
pub struct CliProgressReporter {
printed_header: bool,
last_line_len: usize,
last_completed_sum: u64,
last_instant: std::time::Instant,
totals_by_digest: HashMap<String, (u64, u64)>,
}
impl Default for CliProgressReporter {
fn default() -> Self {
Self::new()
}
}
impl CliProgressReporter {
pub fn new() -> Self {
Self {
printed_header: false,
last_line_len: 0,
last_completed_sum: 0,
last_instant: std::time::Instant::now(),
totals_by_digest: HashMap::new(),
}
}
}totals_by_digest 的值为 (total, completed),保存各层最近已知值。printed_header 控制总量提示只打印一次;last_line_len 用于覆盖终端上一行;last_completed_sum 与 last_instant 是速度采样点。这里没有跨拉取共享状态,每次 ensure_oss_ready 创建新的 reporter,状态从零开始。
源码文件:codex-rs/ollama/src/pull.rs
相关函数/类型:CliProgressReporter::on_event
// ...
PullEvent::ChunkProgress {
digest,
total,
completed,
} => {
// 赋值替换快照,重复上报不应当重复累计。
if let Some(t) = *total {
self.totals_by_digest
.entry(digest.clone())
.or_insert((0, 0))
.0 = t;
}
if let Some(c) = *completed {
self.totals_by_digest
.entry(digest.clone())
.or_insert((0, 0))
.1 = c;
}
let (sum_total, sum_completed) = self
.totals_by_digest
.values()
.fold((0u64, 0u64), |acc, (t, c)| (acc.0 + *t, acc.1 + *c));
if sum_total > 0 {
if !self.printed_header {
let gb = (sum_total as f64) / (1024.0 * 1024.0 * 1024.0);
let header = format!("Downloading model: total {gb:.2} GB\n");
out.write_all(b"\r\x1b[2K")?;
out.write_all(header.as_bytes())?;
self.printed_header = true;
}
let now = std::time::Instant::now();
let dt = now
.duration_since(self.last_instant)
.as_secs_f64()
.max(0.001);
// 下降时钳制为零,避免 u64 下溢;速度依据上一次事件间隔。
let dbytes = sum_completed.saturating_sub(self.last_completed_sum) as f64;
let speed_mb_s = dbytes / (1024.0 * 1024.0) / dt;
self.last_completed_sum = sum_completed;
self.last_instant = now;
let done_gb = (sum_completed as f64) / (1024.0 * 1024.0 * 1024.0);
let total_gb = (sum_total as f64) / (1024.0 * 1024.0 * 1024.0);
let pct = (sum_completed as f64) * 100.0 / (sum_total as f64);
let text =
format!("{done_gb:.2}/{total_gb:.2} GB ({pct:.1}%) {speed_mb_s:.1} MB/s");
let pad = self.last_line_len.saturating_sub(text.len());
let line = format!("\r{text}{}", " ".repeat(pad));
self.last_line_len = text.len();
out.write_all(line.as_bytes())?;
out.flush()
} else {
Ok(())
}
}
// ...有 Some(total) 就替换这一层的总量,有 Some(completed) 就替换这一层的完成量,缺失字段保留之前的值。它们不是“本次又下载了多少”的增量。若把每条 completed 直接加到全局计数,同一快照重发时会重复计算,进度很快超过真实值。
下面以 MiB 为输入单位演算,实际事件字段仍然是字节整数。
| 到达事件 | A 的状态 | B 的状态 | 总进度 |
|---|---|---|---|
| A: total=100,completed=50 | 100 / 50 | 未出现 | 50 / 100 = 50% |
| B: total=300,completed=0 | 100 / 50 | 300 / 0 | 50 / 400 = 12.5% |
| A: completed=60 | 100 / 60 | 300 / 0 | 60 / 400 = 15% |
| A: completed=60 再发一次 | 100 / 60 | 300 / 0 | 仍为 15% |
| B: completed=100 | 100 / 60 | 300 / 100 | 160 / 400 = 40% |
新层出现使分母变大,所以百分比可以下降,不能据此判定下载回滚。初次 Downloading model: total ... 提示只打印一次,其总量可能只是当时已经见到的层总量;后续实时行才会反映新发现的层。
速度计算用相邻进度事件之间的完成量差除以时间差,时间下限为 0.001 秒,字节差用 saturating_sub 防止计数下降时无符号下溢。它是事件间采样速率,不是整次下载的平均速度,更不是网络链路测速。显示标签写 GB/MB,但除数是 1024 的幂;需要精确换算时应按 GiB/MiB 理解。
代码没有把百分比钳制在 100%,也没有验证 completed <= total。如果服务提供不一致的计数,显示可能超过 100%;客户端当前优先按原始快照展示,而不是自行修补进度数据。
7.2 终端呈现
源码文件:codex-rs/ollama/src/pull.rs
相关函数/类型:CliProgressReporter::on_event
// ...
PullEvent::Status(status) => {
// Avoid noisy manifest messages; otherwise show status inline.
if status.eq_ignore_ascii_case("pulling manifest") {
return Ok(());
}
// 覆盖较短文字后要补空格,否则旧进度尾部仍留在终端。
let pad = self.last_line_len.saturating_sub(status.len());
let line = format!("\r{status}{}", " ".repeat(pad));
self.last_line_len = status.len();
out.write_all(line.as_bytes())?;
out.flush()
}
// ...状态消息通过回车覆盖当前行。新状态比旧文本短时补空格,避免残留上一条进度尾部;pulling manifest 按忽略大小写的比较跳过。这里使用字符串字节长度计算 padding,因此它不是通用的 Unicode 终端显示宽度布局器。
成功事件在 CliProgressReporter::on_event 的 Success 分支输出换行;错误事件不重复打印,由上层返回错误统一呈现。这也解释了为什么底层 Error 不一定在进度行里直接留下错误文本。
源码文件:codex-rs/ollama/src/pull.rs
相关函数/类型:TuiProgressReporter
// 导出类型的名字不代表独立 TUI widget;当前委托给 CLI 输出。
/// For now the TUI reporter delegates to the CLI reporter. This keeps UI and
/// CLI behavior aligned until a dedicated TUI integration is implemented.
#[derive(Default)]
pub struct TuiProgressReporter(CliProgressReporter);
impl PullProgressReporter for TuiProgressReporter {
fn on_event(&mut self, event: &PullEvent) -> io::Result<()> {
self.0.on_event(event)
}
}TuiProgressReporter 当前是 CLI reporter 的包装,而 ensure_oss_ready 实际直接构造 CliProgressReporter。因此不能根据公开导出名画出一条独立的 TUI widget 事件队列。若要开发真正的 TUI 进度视图,接入点是 PullProgressReporter,同时必须保留回调失败会影响主流程的契约。
8. 交接与中断
8.1 进入模型请求
准备结束后,启动流程继续创建或连接会话;之后一次 Turn,即一次用户输入驱动的执行周期,会通过通用模型客户端生成 Responses 请求。codex-ollama 不持有 Turn,也不把拉取进度塞进 ResponseEvent。这两条事件链在类型和消费者上都是分开的。
沿源码继续走时,可用下表确定交接位置,而不需要在本篇重复整个模型客户端。
| 阶段 | 真实入口/消费者 | 传递内容 |
|---|---|---|
| 本地准备返回 | utils/oss/src/lib.rs::ensure_oss_provider_ready | Result<()>,没有模型句柄 |
| 启动与会话装配 | tui/src/startup_orchestration.rs::run_main_inner;exec/src/lib.rs 的启动流程 | 已选 Provider 和模型的 Config |
| 普通模型请求 | core/src/client.rs 的 ModelClient/Turn session | 模型请求与 Provider setup |
| HTTP Responses | codex-api/src/endpoint/responses.rs::ResponsesClient::stream_request | 编码后的 ResponsesApiRequest |
| 流式响应 | 同文件 stream_encoded → spawn_response_stream | POST responses,Accept: text/event-stream |
表中路径均相对于 codex-rs/。最后一个接口路径由 Provider 的基础地址拼接,内置 Ollama 对应 /v1/responses。当前配置的 WireApi::Responses 不会因 Ollama 的旧版本自动切回 Chat Completions;旧 wire_api = "chat" 和 ollama-chat 的移除错误在 model-provider-info/src/lib.rs 中明确声明。
列表含有模型名称,也不能证明模型能够理解工具 schema、多模态输入或返回所需事件。遇到模型输出协议问题,应继续读模型请求构造和 Responses 流式事件解析;若失败发生在 /api/pull,则留在本篇的管理接口链路排查。
8.2 取消边界
TUI 在执行本地准备之前恢复终端,再调用共用 helper。
源码文件:codex-rs/tui/src/startup_orchestration.rs
相关函数/类型:run_main_inner
// ...
startup_draft
.tui_mut()
.with_restored(|| async {
// Provider setup may print progress or block in an external downloader.
// Restore ordinary signal handling so Ctrl+C can interrupt that process.
// 启动准备发生在恢复后的终端环境,让 Ctrl+C 交回普通信号处理。
crossterm::terminal::disable_raw_mode()?;
ensure_oss_provider_ready(provider_id, &config).await
})
.await?;
// ...这段等待位于启动过程,普通信号处理允许用户用 Ctrl+C 中断准备。它没有进入 Core Turn 的 Interrupt/取消 token 管理,也没有把下载注册成 Session 后台任务。上述 Rust 管理模块没有按操作系统选择另一套拉取算法;平台差异主要来自终端信号和底层传输,而不是一个隐藏的 Windows 下载分支。
正常路径由消费者读到 Success 返回;失败路径由 HTTP 建立错误、事件错误、报告器错误或未成功 EOF 返回;外部取消通过停止等待与进程/运行时生命周期结束局部流。三者的共同点是这里没有远端回滚协议。需要判断远端是否还在下载时,应检查服务自身状态,不能只看 Codex 已经退出。
9. 边界复现
9.1 游标断言
先从最小的无网络测试练习不变量,再读 HTTP 测试。以下上游测试分三次追加数据,并直接比较内部游标。
源码文件:codex-rs/ollama/src/line_buffer_tests.rs
相关函数/类型:searches_only_new_bytes_after_partial_line
// 测试直接观察缓冲内容和扫描游标,验证分片间的不变量。
fn searches_only_new_bytes_after_partial_line() {
let mut buffer = LineBuffer::default();
buffer.extend_from_slice(b"partial");
assert_eq!(buffer.take_line(), None);
assert_eq!(
buffer,
LineBuffer {
bytes: BytesMut::from(&b"partial"[..]),
scanned_len: 7,
}
);
buffer.extend_from_slice(b" line");
assert_eq!(buffer.take_line(), None);
assert_eq!(
buffer,
LineBuffer {
bytes: BytesMut::from(&b"partial line"[..]),
scanned_len: 12,
}
);
buffer.extend_from_slice(b"\nnext");
assert_eq!(
buffer.take_line(),
Some(BytesMut::from(&b"partial line\n"[..]))
);
assert_eq!(
buffer,
LineBuffer {
bytes: BytesMut::from(&b"next"[..]),
scanned_len: 0,
}
);
}第一次只有 partial,断言游标为 7;追加 line 后为 12;追加换行和 next 后,取出的行包含换行,剩余内容为 next 且游标回到 0。它验证的是“已扫描前缀”和“拆行后重置”的正确性,不证明任意网络故障下都能恢复,也不证明没有内存上限问题。
进一步把 {"status":"正在下载"}\n 的 UTF-8 字节逐字节喂给原始 LineBuffer,断言取回的行能解码且拼接后与输入一致,可以覆盖字符跨 chunk 的场景。这个实验的对象仍是字节缓冲,不应把正确拆行直接等同于完整拉取成功。
9.2 请求次数
下面保留 version_check_reuses_existing_ollama_client 的 mock 设置、调用和关键断言。测试采用不带 /v1 的原生根地址。
源码文件:codex-rs/ollama/src/lib.rs
相关函数/类型:version_check_reuses_existing_ollama_client
// ...
let server = wiremock::MockServer::start().await;
wiremock::Mock::given(wiremock::matchers::method("GET"))
.and(wiremock::matchers::path("/api/tags"))
.respond_with(
wiremock::ResponseTemplate::new(200)
.set_body_json(serde_json::json!({"models": [{"name": "gpt-oss:20b"}]})),
)
// 原生 URL 构造探测一次,后续 fetch_models 再查一次。
.expect(2)
.mount(&server)
.await;
wiremock::Mock::given(wiremock::matchers::method("GET"))
.and(wiremock::matchers::path("/api/version"))
.respond_with(
wiremock::ResponseTemplate::new(200)
.set_body_json(serde_json::json!({"version": "0.14.1"})),
)
// 版本检查只请求 version;不另建客户端重复探测。
.expect(1)
.mount(&server)
.await;
let provider = create_oss_provider_with_base_url(&server.uri(), WireApi::Responses);
let client = OllamaClient::try_from_provider(
&provider,
HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault),
)
.await
.expect("create Ollama client");
ensure_responses_supported(&client)
.await
.expect("version check should reuse the existing client");
assert_eq!(
client.fetch_models().await.expect("fetch models"),
vec!["gpt-oss:20b"]
);
server.verify().await;
// .../api/tags 期望 2 次:构造时原生探测一次,显式 fetch_models 一次;/api/version 只需 1 次。server.verify() 将请求次数变成可失败的断言,能够抓到“版本检查偷偷重建客户端并重复探测”的改动。它没有覆盖自动下载,因为测试没有调用 ensure_oss_ready。
另一项 client::tests::test_pull_model_stream_parses_large_json_lines 在首行加入 128 KiB 的 padding,第二行使用 status: "complete",断言底层收集到两个 Status。这个测试验证大行解析与忽略额外字段;第二条没有使用 success,所以它不能证明高层拉取成功。把同样的流交给 pull_with_reporter,就会在 EOF 返回错误。
在源码仓库根目录执行:
cd codex-rs
just test -p codex-ollama --locked --lib该 crate 的 16 项测试包含版本门槛、URL 派生、行游标、模型查询、探测、代理策略和自定义 CA 分支。网络测试有 CODEX_SANDBOX_NETWORK_DISABLED 环境变量 guard;判断是否真正执行了 mock 请求时,需要核对该 guard,不能仅看外层运行器显示 passed。测试使用本地 mock 服务,不下载模型。
9.3 异常输入
为检查上游测试没有直接覆盖的消费者差异,可把 /api/pull 的 mock 响应替换成下列原始字节,再分别调用底层流与高层 helper。表中 \n 代表一个真正的换行字节,而不是两个文本字符。
| mock 字节输入 | 底层收集结果 | 高层结果 |
|---|---|---|
{"status":"success"}\n | 状态、成功、成功 | 成功,只消费第一次成功 |
{"status":"success"} | 无事件 | 未成功 EOF 错误 |
{"status":"complete"}\n | 一个状态事件 | 未成功 EOF 错误 |
{"error":"missing-model"}\n,HTTP 200 | 一个错误事件 | Pull failed: missing-model |
not-json\n{"status":"success"}\n | 跳过首行,随后成功事件 | 成功 |
0xff 加换行,再加成功行 | 跳过非法 UTF-8 行 | 成功 |
{"status":"success","error":"bad"}\n | 状态、成功、错误 | 提前成功 |
这些输入分别改变帧边界、终止拼写、错误载体、容错行和事件优先顺序。它们验证客户端对给定字节的行为,不证明真实服务一定产生矛盾或损坏响应,也不能证明真实模型推理能力。
地址和列表问题可以先通过只读请求定位。将根地址替换为实际服务,按阶段观察 HTTP 状态、JSON 外形与模型名称。
curl -i --max-time 5 http://localhost:11434/v1/models
curl -i --max-time 5 http://localhost:11434/api/version
curl -i --max-time 5 http://localhost:11434/api/tags这里 curl 的 5 秒是诊断命令的总时限,与源码里的 5 秒建连超时含义不同。第一条成功而第三条失败时,优先检查代理是否覆盖原生接口;第三条成功但自动拉取仍发生时,逐字符比较 models[].name 与请求模型。
最后用两个具体任务检验是否能继续独立读代码:给定 --oss --local-provider ollama -m lesson:1,从 SharedCliOptions 追踪到实际 POST 的 JSON,并说明模型列表 503 与损坏 JSON 为何改变是否下载;再给定“进度显示完成却返回未成功 EOF”,指出需要检查的换行字节、complete/success 值和 pull_with_reporter 的返回位置。
如果要修改实现,应该能够把改动定位到清晰边界:URL parser 属于地址派生,最后一行补刷属于 LineBuffer 与流的 EOF 契约,去重成功通知属于 parser 与外层终止检查的职责划分,报告器不影响下载则需要改变 on_event 错误传播策略。每一种修改都需要保留对应的反例,而不能只重跑“正常下载成功”。
另一种本地服务采用不同的完成边界:LM Studio 本地 Provider通过外部 lms 进程下载,并把预载交给分离的后台任务。可以对照两篇的返回条件,判断哪些本地准备保证由共同入口提供,哪些只属于特定 Provider。
