trusty-common 0.52.3

Shared utilities and provider-agnostic streaming chat (ChatProvider, OllamaProvider, OpenRouter, tool-use) for trusty-* projects
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
//! Liveness probe for a local OpenAI-compatible model server (#4490).
//!
//! Why: two call sites answer the same question — "is a local model server
//! actually reachable right now?" — and a wrong answer costs money and latency
//! silently instead of failing loudly. `chat::auto_detect_local_provider` has
//! probed since the `ChatProvider` era, while `inference::providers::local`
//! built an adapter for a fixed base URL unconditionally, so a consumer moving
//! from `ChatProvider` to `InferenceAdapter` (#4427) would have fallen through
//! to OpenRouter on every turn even with a local model running. Copying the
//! probe across would leave two independent implementations of one capability,
//! which CLAUDE.md's common-entry-point rule forbids, so the probe lives here
//! and both callers route through it — the timeout included, as one named
//! constant rather than a literal repeated per call site.
//!
//! What: [`LOCAL_PROBE_TIMEOUT`](crate::local_probe::LOCAL_PROBE_TIMEOUT) (the
//! shared budget, applied to both connect and whole-request),
//! [`models_url`](crate::local_probe::models_url) (the `{host}/v1/models`
//! derivation),
//! [`probe_models_endpoint`](crate::local_probe::probe_models_endpoint) (the GET
//! plus status check), [`probe_local`](crate::local_probe::probe_local) (the
//! two composed), [`list_models`](crate::local_probe::list_models) (the same
//! request, reading back the served model ids), and
//! [`local_host`](crate::local_probe::local_host) (the
//! [`LOCAL_HOST_ENV`](crate::local_probe::LOCAL_HOST_ENV) resolution every
//! caller shares). Failure is a typed
//! [`LocalProbeError`](crate::local_probe::LocalProbeError) that always names the
//! endpoint it dialled, so a caller can report which address was dead rather than
//! surfacing a bare transport timeout.
//!
//! Paths here are crate-absolute on purpose: this module carries docs both here
//! and on the `pub mod local_probe;` declaration in `lib.rs`, rustdoc merges the
//! two, and link resolution takes its scope from the FIRST fragment — the
//! `lib.rs` one — so a bare `LOCAL_PROBE_TIMEOUT` resolves against the crate root
//! and is not found (#6027).
//!
//! Scope: this deliberately does NOT route through
//! [`crate::http_client::loopback_client_builder`]. That entry point disables
//! proxies for LOOPBACK daemon targets, and a local model server is pointable at
//! a non-loopback host (`OLLAMA_HOST=http://192.168.1.50:11434`) where the
//! operator's proxy must still apply — the same carve-out `http_client`'s own
//! docs state for inference providers. [`crate::health_probe::probe_health`] is
//! the loopback-daemon counterpart; it answers `bool`, which cannot carry the
//! endpoint a caller has to report.
//!
//! Test: inline `tests` — `models_url_appends_v1_when_absent`,
//! `models_url_does_not_double_an_existing_v1_suffix`,
//! `probe_reports_unreachable_naming_the_endpoint`,
//! `probe_reports_non_success_status`, `probe_accepts_a_live_endpoint`,
//! `probe_timeout_is_one_second`, `list_models_returns_the_served_ids`,
//! `list_models_reports_an_unreadable_body`,
//! `list_models_is_empty_when_the_server_serves_none`,
//! `local_host_reads_the_env_override`, `local_host_defaults_when_unset`.

use std::time::Duration;

/// The budget one liveness probe gets, for connect AND for the whole request.
///
/// Why: the probe runs on a startup path in front of a fallback decision, so it
/// must never be the thing that makes a command feel hung — a local server that
/// is not up has to be ruled out in about the time a human would wait. One
/// second is what `chat::auto_detect_local_provider` has used since it shipped;
/// naming it once here is what keeps the two callers from drifting apart.
/// What: one second, passed to both `connect_timeout` and `timeout`.
/// Test: `probe_timeout_is_one_second`.
pub const LOCAL_PROBE_TIMEOUT: Duration = Duration::from_secs(1);

/// The OpenAI-compatible path every local server implements for a cheap,
/// side-effect-free liveness check.
pub const LOCAL_MODELS_PATH: &str = "/v1/models";

/// Env var naming the local model server's BARE host (no `/v1` suffix).
///
/// Why: four modules across three crates read this variable to answer the same
/// question — which host is the local model server on. Named `OLLAMA_HOST`
/// rather than a `TRUSTY_*` spelling because `trusty-agents`' legacy adapter
/// already read it, so one export configures every path.
/// What: the variable name only; [`local_host`] applies it.
/// Test: `local_host_reads_the_env_override`.
pub const LOCAL_HOST_ENV: &str = "OLLAMA_HOST";

/// The host a local model server listens on when nothing overrides it.
pub const DEFAULT_LOCAL_HOST: &str = "http://localhost:11434";

/// Resolve the local model server's bare host.
///
/// Why (#4490): `repl::ollama::ollama_host`, `llm::adapter::ollama_host`,
/// `api::server::models::ollama_host` and `LocalConfig::from_env` were four
/// independent copies of this one line, so a host spelled one way in the REPL
/// and another in the model catalog could disagree about which server was being
/// probed. One resolver is what keeps the probe and the request dialling the
/// same machine.
/// What: [`LOCAL_HOST_ENV`] when set and non-blank (trailing slashes trimmed),
/// otherwise [`DEFAULT_LOCAL_HOST`]. The result carries NO `/v1` suffix —
/// [`models_url`] and `LocalConfig::from_env` each append what they need.
/// Test: `local_host_reads_the_env_override`, `local_host_defaults_when_unset`.
pub fn local_host() -> String {
    std::env::var(LOCAL_HOST_ENV)
        .ok()
        .filter(|v| !v.trim().is_empty())
        .map(|host| host.trim_end_matches('/').to_string())
        .unwrap_or_else(|| DEFAULT_LOCAL_HOST.to_string())
}

/// Why the local model server could not be confirmed live.
///
/// Why: the caller's job after a failed probe is to tell an operator WHICH
/// address was dead — "connection refused" alone sends them looking at the wrong
/// host when `OLLAMA_HOST` points somewhere unexpected. Every variant therefore
/// carries the endpoint, and [`Self::endpoint`] reads it back without a match.
/// What: `ClientBuild` (the HTTP client could not be constructed at all),
/// `Unreachable` (no response inside [`LOCAL_PROBE_TIMEOUT`] — refused, timed
/// out, DNS, TLS), `Status` (the server answered, but not 2xx), and `Body` (a
/// 2xx whose payload [`list_models`] could not read). No variant carries a
/// credential: the probe sends no `Authorization` header.
/// Test: `probe_reports_unreachable_naming_the_endpoint`,
/// `probe_reports_non_success_status`, `list_models_reports_an_unreadable_body`.
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum LocalProbeError {
    /// The `reqwest` client could not be built (never observed in practice).
    ClientBuild {
        /// The endpoint the probe would have dialled.
        endpoint: String,
        /// The stringified builder error.
        cause: String,
    },
    /// No response arrived within [`LOCAL_PROBE_TIMEOUT`].
    Unreachable {
        /// The endpoint that did not answer.
        endpoint: String,
        /// The stringified transport error.
        cause: String,
    },
    /// The server answered with a non-2xx status.
    Status {
        /// The endpoint that answered.
        endpoint: String,
        /// The HTTP status it returned.
        status: u16,
    },
    /// The server answered 2xx, but the body was not the expected JSON shape
    /// (#4490 — only [`list_models`] can produce this; a liveness probe never
    /// reads the body).
    Body {
        /// The endpoint that answered.
        endpoint: String,
        /// The stringified decode error.
        cause: String,
    },
}

impl LocalProbeError {
    /// The endpoint this probe dialled.
    ///
    /// Why: every caller reports it, and none of them should have to match on
    /// the variant to get at it.
    /// What: the `endpoint` field of whichever variant this is.
    /// Test: `probe_reports_unreachable_naming_the_endpoint`.
    pub fn endpoint(&self) -> &str {
        match self {
            Self::ClientBuild { endpoint, .. }
            | Self::Unreachable { endpoint, .. }
            | Self::Status { endpoint, .. }
            | Self::Body { endpoint, .. } => endpoint,
        }
    }
}

impl std::fmt::Display for LocalProbeError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            Self::ClientBuild { endpoint, cause } => write!(
                f,
                "could not build a probe client for local model server {endpoint}: {cause}"
            ),
            Self::Unreachable { endpoint, cause } => write!(
                f,
                "local model server not reachable at {endpoint} within {}s: {cause}",
                LOCAL_PROBE_TIMEOUT.as_secs()
            ),
            Self::Status { endpoint, status } => write!(
                f,
                "local model server at {endpoint} answered HTTP {status}, not a success status"
            ),
            Self::Body { endpoint, cause } => write!(
                f,
                "local model server at {endpoint} answered with an unreadable model list: {cause}"
            ),
        }
    }
}

impl std::error::Error for LocalProbeError {}

/// Derive the `/v1/models` probe URL from a local server's base URL.
///
/// Why: the two callers spell their base URL differently and neither should have
/// to know the other's convention. `chat::LocalModelConfig::base_url` is a bare
/// host (`http://localhost:11434`); `inference::providers::local::LOCAL_BASE_URL`
/// already carries the `/v1` suffix the OpenAI dialect needs. One derivation that
/// accepts both is what lets a single probe serve both.
/// What: trims trailing slashes, drops a trailing `/v1` if present, and appends
/// [`LOCAL_MODELS_PATH`]. A bare host is therefore unchanged from what
/// `chat::auto_detect_local_provider` has always produced.
/// Test: `models_url_appends_v1_when_absent`,
/// `models_url_does_not_double_an_existing_v1_suffix`.
pub fn models_url(base_url: &str) -> String {
    let base = base_url.trim_end_matches('/');
    let host = base.strip_suffix("/v1").unwrap_or(base);
    format!("{host}{LOCAL_MODELS_PATH}")
}

/// GET an already-derived models endpoint and report whether it is live.
///
/// Why: the single implementation of the probe itself, so the timeout, the
/// success criterion, and the error text cannot drift between `chat::` and
/// `inference::`.
/// What: builds a client bounded by [`LOCAL_PROBE_TIMEOUT`] on connect and on
/// the whole request, GETs `url`, and returns `Ok(())` for any 2xx. Anything
/// else — a build failure, a transport failure, a timeout, a non-2xx status — is
/// a [`LocalProbeError`] naming `url`. Sends no credential.
/// Test: `probe_reports_unreachable_naming_the_endpoint`,
/// `probe_reports_non_success_status`, `probe_accepts_a_live_endpoint`.
pub async fn probe_models_endpoint(url: &str) -> Result<(), LocalProbeError> {
    get_success(url).await.map(|_| ())
}

/// GET `url` inside the shared budget and return the 2xx response.
///
/// Why: [`probe_models_endpoint`] and [`list_models`] issue the IDENTICAL
/// request and differ only in whether they read the body, so the client, the
/// timeout, and the success criterion live here once.
/// What: builds a client bounded by [`LOCAL_PROBE_TIMEOUT`] on connect and on
/// the whole request, GETs `url`, and returns the response for any 2xx. Sends no
/// credential.
/// Test: `probe_reports_unreachable_naming_the_endpoint`,
/// `probe_reports_non_success_status`, `probe_accepts_a_live_endpoint`.
async fn get_success(url: &str) -> Result<reqwest::Response, LocalProbeError> {
    // #4490: proxies are deliberately left enabled — see the module's Scope note.
    let client = reqwest::Client::builder()
        .connect_timeout(LOCAL_PROBE_TIMEOUT)
        .timeout(LOCAL_PROBE_TIMEOUT)
        .build()
        .map_err(|e| LocalProbeError::ClientBuild {
            endpoint: url.to_string(),
            cause: e.to_string(),
        })?;

    match client.get(url).send().await {
        Ok(resp) if resp.status().is_success() => Ok(resp),
        Ok(resp) => Err(LocalProbeError::Status {
            endpoint: url.to_string(),
            status: resp.status().as_u16(),
        }),
        Err(e) => Err(LocalProbeError::Unreachable {
            endpoint: url.to_string(),
            cause: e.to_string(),
        }),
    }
}

/// Probe a local model server given its base URL.
///
/// Why: the form both callers actually want — they hold a base URL, not a models
/// endpoint.
/// What: [`models_url`] followed by [`probe_models_endpoint`].
/// Test: `probe_accepts_a_live_endpoint`,
/// `probe_reports_unreachable_naming_the_endpoint`.
pub async fn probe_local(base_url: &str) -> Result<(), LocalProbeError> {
    probe_models_endpoint(&models_url(base_url)).await
}

/// Probe a local model server AND read back the model ids it serves.
///
/// Why (#4490): `trusty-agents`' `/provider local` needs more than liveness — it
/// prints the pulled models so the user can pick one with `/model`. It answered
/// that with its own `GET {host}/api/tags` and its own 2s timeout, a second
/// independent implementation of this capability. Extending the shared probe to
/// return the list is what lets that call site delegate instead: liveness and
/// the catalog come from ONE request, one timeout, and one typed error.
/// What: [`get_success`] against [`models_url`], then reads the OpenAI-dialect
/// `{"data":[{"id":"…"}]}` body and collects the ids in server order. Ollama's
/// `/v1/models` shim reports the same names its native `/api/tags` does, so a
/// caller loses nothing by moving; LM Studio and vLLM serve the same shape,
/// which `/api/tags` never covered. An unreadable body is
/// [`LocalProbeError::Body`]; a missing or non-array `data` is an empty list,
/// not an error — a server with no models pulled is live.
/// Test: `list_models_returns_the_served_ids`,
/// `list_models_reports_an_unreadable_body`,
/// `list_models_is_empty_when_the_server_serves_none`.
pub async fn list_models(base_url: &str) -> Result<Vec<String>, LocalProbeError> {
    let url = models_url(base_url);
    let resp = get_success(&url).await?;
    let body: serde_json::Value = resp.json().await.map_err(|e| LocalProbeError::Body {
        endpoint: url.clone(),
        cause: e.to_string(),
    })?;
    Ok(body
        .get("data")
        .and_then(|v| v.as_array())
        .map(|arr| {
            arr.iter()
                .filter_map(|m| m.get("id").and_then(|id| id.as_str()).map(str::to_string))
                .collect()
        })
        .unwrap_or_default())
}

// ── Tests ────────────────────────────────────────────────────────────────────

#[cfg(test)]
mod tests {
    use super::*;

    /// A loopback stub answering one canned status per connection, forever.
    ///
    /// Modelled on `http_client::tests::stub_server` — a real listener, so the
    /// probe under test drives the real transport.
    async fn stub_server(response: impl Into<String>) -> String {
        let response: std::sync::Arc<str> = response.into().into();
        let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
            .await
            .expect("bind loopback stub");
        let addr = listener.local_addr().expect("stub addr").to_string();
        tokio::spawn(async move {
            while let Ok((mut stream, _)) = listener.accept().await {
                let response = std::sync::Arc::clone(&response);
                tokio::spawn(async move {
                    use tokio::io::AsyncWriteExt;
                    let _ = stream.write_all(response.as_bytes()).await;
                });
            }
        });
        addr
    }

    /// A 200 response carrying `body` as JSON, correctly framed.
    fn json_ok(body: &str) -> String {
        format!(
            "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{body}",
            body.len()
        )
    }

    /// An address with nothing listening: bind port 0, read it, release it.
    fn dead_addr() -> String {
        let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind to free a port");
        let addr = listener.local_addr().expect("dead addr").to_string();
        drop(listener);
        addr
    }

    /// Why: a bare host is what `chat::LocalModelConfig::base_url` holds, and
    /// the URL this produces for it must stay byte-identical to the one
    /// `chat::auto_detect_local_provider` has always built.
    /// Test: this test.
    #[test]
    fn models_url_appends_v1_when_absent() {
        assert_eq!(
            models_url("http://localhost:11434"),
            "http://localhost:11434/v1/models"
        );
        assert_eq!(
            models_url("http://localhost:11434/"),
            "http://localhost:11434/v1/models"
        );
    }

    /// Why: `inference::providers::local::LOCAL_BASE_URL` already ends in `/v1`,
    /// so a naive append would dial `/v1/v1/models` and report every live server
    /// as dead.
    /// Test: this test.
    #[test]
    fn models_url_does_not_double_an_existing_v1_suffix() {
        assert_eq!(
            models_url("http://localhost:11434/v1"),
            "http://localhost:11434/v1/models"
        );
        assert_eq!(
            models_url("http://localhost:11434/v1/"),
            "http://localhost:11434/v1/models"
        );
    }

    /// Why: the whole point of the typed error is that a caller can name the
    /// address that was dead; a bare "connection refused" sends an operator to
    /// the wrong host.
    /// Test: this test.
    #[tokio::test]
    async fn probe_reports_unreachable_naming_the_endpoint() {
        let addr = dead_addr();
        let base = format!("http://{addr}");
        let err = probe_local(&base).await.expect_err("closed port must fail");
        assert!(
            matches!(err, LocalProbeError::Unreachable { .. }),
            "expected Unreachable, got {err:?}"
        );
        assert_eq!(err.endpoint(), format!("{base}/v1/models"));
        assert!(err.to_string().contains(&addr), "{err}");
    }

    /// Why: a server that answers but does not serve `/v1/models` is not a
    /// usable local model server, and must be rejected rather than treated as
    /// live because a socket accepted the connection.
    /// Test: this test.
    #[tokio::test]
    async fn probe_reports_non_success_status() {
        let addr = stub_server("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\n\r\n").await;
        let err = probe_local(&format!("http://{addr}"))
            .await
            .expect_err("404 must fail");
        assert_eq!(
            err,
            LocalProbeError::Status {
                endpoint: format!("http://{addr}/v1/models"),
                status: 404,
            }
        );
    }

    /// Why: the positive half — a reachable server must pass, or the probe would
    /// disable local inference entirely.
    /// Test: this test.
    #[tokio::test]
    async fn probe_accepts_a_live_endpoint() {
        let addr = stub_server("HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok").await;
        probe_local(&format!("http://{addr}"))
            .await
            .expect("live endpoint must probe clean");
    }

    /// Why (#4490): the budget is a contract, not an implementation detail —
    /// both callers depend on a failed probe costing about a second, and this is
    /// the single definition they share.
    /// Test: this test.
    #[test]
    fn probe_timeout_is_one_second() {
        assert_eq!(LOCAL_PROBE_TIMEOUT, Duration::from_secs(1));
    }

    /// Why (#4490): the REPL's `/provider local` prints this list, so the shared
    /// probe has to hand back the same names the bespoke `/api/tags` call did —
    /// otherwise the delegation silently empties the model picker.
    /// Test: this test.
    #[tokio::test]
    async fn list_models_returns_the_served_ids() {
        let addr = stub_server(json_ok(
            r#"{"object":"list","data":[{"id":"qwen3:30b"},{"id":"llama3.1:8b"}]}"#,
        ))
        .await;
        let models = list_models(&format!("http://{addr}"))
            .await
            .expect("live endpoint must list");
        assert_eq!(models, vec!["qwen3:30b", "llama3.1:8b"]);
    }

    /// Why: a 2xx with a body that is not JSON is a misconfigured endpoint (a
    /// proxy login page, say), and must name the endpoint rather than surface as
    /// an empty model list that reads like "no models pulled".
    /// Test: this test.
    #[tokio::test]
    async fn list_models_reports_an_unreadable_body() {
        let addr = stub_server(
            "HTTP/1.1 200 OK\r\nContent-Type: text/html\r\nContent-Length: 5\r\n\r\nnope!",
        )
        .await;
        let base = format!("http://{addr}");
        let err = list_models(&base).await.expect_err("non-JSON must fail");
        assert!(
            matches!(err, LocalProbeError::Body { .. }),
            "expected Body, got {err:?}"
        );
        assert_eq!(err.endpoint(), format!("{base}/v1/models"));
    }

    /// Why: a running server with nothing pulled is LIVE — reporting that as an
    /// error would make `/provider local` refuse a server the operator can fix
    /// with one `ollama pull`.
    /// Test: this test.
    #[tokio::test]
    async fn list_models_is_empty_when_the_server_serves_none() {
        let addr = stub_server(json_ok(r#"{"object":"list","data":[]}"#)).await;
        assert!(
            list_models(&format!("http://{addr}"))
                .await
                .expect("empty catalog is still live")
                .is_empty()
        );
    }

    /// Why (#4490): four call sites across three crates read this variable; if
    /// the shared resolver ignored it, a remote Ollama host would be probed at
    /// localhost and reported dead.
    /// Test: this test.
    #[test]
    #[serial_test::serial]
    fn local_host_reads_the_env_override() {
        // SAFETY: guarded by `#[serial]`; no other thread reads the env here.
        unsafe { std::env::set_var(LOCAL_HOST_ENV, "http://192.168.1.50:11434/") };
        assert_eq!(local_host(), "http://192.168.1.50:11434");
        unsafe { std::env::remove_var(LOCAL_HOST_ENV) };
    }

    /// Why: the default is what every caller gets on a bare install, so it must
    /// survive an unset AND a blank override.
    /// Test: this test.
    #[test]
    #[serial_test::serial]
    fn local_host_defaults_when_unset() {
        // SAFETY: guarded by `#[serial]`; no other thread reads the env here.
        unsafe { std::env::remove_var(LOCAL_HOST_ENV) };
        assert_eq!(local_host(), DEFAULT_LOCAL_HOST);
        unsafe { std::env::set_var(LOCAL_HOST_ENV, "   ") };
        assert_eq!(local_host(), DEFAULT_LOCAL_HOST);
        unsafe { std::env::remove_var(LOCAL_HOST_ENV) };
    }
}