Skip to main content

qcode/provider/
relay.rs

1//! The relay that lets a harness inside a container talk to a provider without the provider's
2//! key ever entering that container.
3//!
4//! A profile container has no network of its own, so a harness inside it cannot reach a provider
5//! directly. [`assets/provider/qcode-relay.mjs`](../../../assets/provider/qcode-relay.mjs) is a
6//! small HTTP server the container runs on its loopback interface; a harness is pointed at it the
7//! way it would be pointed at the provider itself. The script forwards every request over a unix
8//! socket to this module, which asks the real provider with the tab's key added, and streams the
9//! answer back.
10//!
11//! This is the same shape as [`crate::bridge`]: a socket in the workspace's `Containers/MCP/`
12//! folder ([`crate::base::paths::MCP_DIR`] in the container), a script written beside it, and a
13//! token every harness tab already carries in [`crate::bridge::TOKEN_VARIABLE`] so QCode knows
14//! which tab is asking and therefore which provider entry to use. No second token is minted: a
15//! tab has one identity, and the bridge already gives it one nobody outside this process can
16//! guess.
17//!
18//! # Framing
19//!
20//! One connection to the socket carries one request and its answer. The container writes one
21//! line of JSON — the token, the method, the path, the headers it read from the harness minus
22//! `authorization`, `x-api-key` and `api-key` — then the request body as raw bytes, ending the write side of
23//! the connection when the body is done (a harness's request is small, so this is read whole
24//! before it goes on). QCode reads that line, decides whether to carry the request at all, and
25//! if so writes back its own line of JSON — the status and the headers the provider answered
26//! with — followed by the answer's body as raw bytes, copied across as they arrive rather than
27//! collected first. A line of JSON for what fits on one line and needs to be decided before
28//! anything else, raw bytes for a body that must never wait to be whole: the same trade [`crate::
29//! bridge::protocol`] makes, and a harness reading an answer as server-sent events depends on it
30//! exactly the way a tab's own agent depends on the bridge answering promptly.
31//!
32//! QCode never has to speak HTTP to do this: the container speaks HTTP to the harness and to
33//! nothing else, and everything on this side of the socket is the small framing above.
34//!
35//! # What is allowed
36//!
37//! A connection may send a message (`POST`, at the path [`crate::provider::Wire::messages_path`]
38//! names for either shape) or list the provider's models (`GET /v1/models`, which both shapes
39//! serve at the same path). Nothing else is carried: the relay forwards one conversation, not
40//! whatever address a container asks for.
41//!
42//! # Which model, and which one after that
43//!
44//! The harness names the model it wants in the body it sends, and that name is not QCode's
45//! business to obey: the person chose a model when they made the profile, not the harness when it
46//! started a background task. So the `model` in a request that has one is set to the step's model
47//! before the request goes out, and a body that is not a JSON object naming a model — a model
48//! listing, a form-encoded body — goes on exactly as it was written. That pinning is the one
49//! assurance that no request of a person's tab ever reaches a model they did not choose, which on
50//! a provider that bills per model is money spent on somebody else's idea of a small task.
51//!
52//! A profile that chose one model has one step and nothing to fall back to. A profile that chose
53//! a lineup has its steps in the order the person wrote them, and a step is given up on only on
54//! the answer's head: the provider could not be reached, `402`, `403`, `408`, `429` or any `5xx`,
55//! and `400` and `404` only when what they say names the model that was asked — those two are the
56//! same complaint for every model unless they are about one. A `401` is never a reason: there is
57//! one key, and the next step would be refused the same way. Nothing is decided on a `2xx`: once a
58//! good answer has started going to the harness, the model for this request is settled and a
59//! stream that breaks halfway is passed on as it broke.
60//!
61//! The last step's answer is the harness's whatever it is, body intact, so a tab shows the
62//! provider's own error instead of one QCode made up — unless every step of the request failed for
63//! good while an earlier one failed only for now. A `408`, a `429`, a `5xx` or a provider that
64//! could not be reached tells the person to wait; a model that does not exist, or a key the
65//! provider will not take, tells them to change the lineup. When one request produced both, the
66//! reason to wait is the true story of it: a person told that the last model of their lineup does
67//! not exist goes looking for a lineup that is not broken, while a minute of waiting would have
68//! answered. So the first passing answer of the request is shown instead, with the status and the
69//! body the provider sent, and a lasting answer that was never a failure of the lineup — a `2xx`, a
70//! `401`, a `400` about the request itself — goes on exactly as it always did.
71//!
72//! A step that was given up on is not asked again for a minute, so a busy free model stops costing
73//! a round trip per turn; that order is kept in this listener, behind one lock its connections
74//! share.
75//!
76//! Which of the two message paths a request takes is the harness's to say, not the provider
77//! entry's. Nothing translates between the shapes, so the only shape that can work is the one
78//! the harness speaks — Claude Code the Anthropic one, opencode the OpenAI one — and every kind
79//! of provider QCode knows serves both. Holding a tab to its provider's recorded shape would only
80//! refuse the one request its harness can make. The shape does decide where the request goes,
81//! since one service answers the two shapes under different roots ([`crate::provider::
82//! ProviderKind::api_root`]), and the key goes in the header the entry's kind reads
83//! ([`crate::provider::ProviderKind::key_header`]).
84//!
85//! # What is being served
86//!
87//! The relay is the only side of a request that can say whether one is on its way to a model: the
88//! script in the container hands it over and gets the answer back, so nothing between the harness
89//! and the provider is left holding a moment in time. [`Activity`] is what a screen reads for
90//! that — per token, how many of its requests are being served and when the last one finished.
91
92use std::collections::HashMap;
93use std::io::{self, Read};
94use std::path::{Path, PathBuf};
95use std::sync::{Arc, Mutex};
96use std::time::{Duration, Instant};
97
98use serde_json::{Map, Value, json};
99
100use super::ProviderEntry;
101use super::ask::{Method, Secret};
102
103/// Where one request goes and what it may be asked with: the provider entry to ask, the models
104/// to try for it in order, and the lineup's name when the profile chose one.
105///
106/// A profile that chose a single model sends one step. A profile that chose a lineup sends its
107/// steps in the order the person wrote them, and the relay walks down that order whenever a step
108/// answers with a reason the next one may not have.
109#[derive(Debug, Clone)]
110pub struct Route {
111    /// The provider to ask, as `providers.toml` holds it and as [`Listener::open`]'s `resolve`
112    /// read it for this very request.
113    pub entry: ProviderEntry,
114    /// The models to ask, in the order they are tried. At least one: a request with no step to
115    /// ask has nowhere to go, and a `resolve` that says so is better than a relay that decides.
116    pub models: Vec<String>,
117    /// The lineup's name, when the profile chose one rather than a single model. It is carried
118    /// for the relay's own reports only: which lineup a tab runs on is the person's own choice,
119    /// and the relay only has to repeat it in what it says.
120    pub lineup: Option<String>,
121}
122
123/// The name of the relay's script in `Containers/MCP/`, beside the bridge's.
124pub const SCRIPT_NAME: &str = "qcode-relay.mjs";
125
126/// The name of the relay's socket in `Containers/MCP/`, beside the bridge's `bridge.sock`.
127pub const SOCKET_NAME: &str = "relay.sock";
128
129/// The port the relay listens on inside the container, on its loopback interface alone. Chosen
130/// for being unlikely to be anything else a harness image already runs; the script and the tests
131/// here are the only two places this number may appear.
132pub const PORT: u16 = 41417;
133
134/// The relay's script, carried inside the binary and written into each open workspace's
135/// `Containers/MCP/` folder, the way [`crate::bridge::SCRIPT`] is.
136pub const SCRIPT: &str = include_str!("../../assets/provider/qcode-relay.mjs");
137
138/// Where the relay's script is inside a profile container, in the folder the workspace's
139/// `Containers/MCP/` is mounted at.
140#[must_use]
141pub fn script_in_container() -> String {
142    format!("{}/{SCRIPT_NAME}", crate::base::paths::MCP_DIR)
143}
144
145/// The program a tab of a provider profile runs: the relay, with the harness after it.
146///
147/// The relay starts the harness itself, once its server is listening, and lives exactly as long
148/// as the harness does. Started beside it instead, the harness could ask before the server was
149/// up and take the refusal for a provider that does not work.
150#[must_use]
151pub fn wrapping(harness: &[String]) -> Vec<String> {
152    let mut command = vec!["node".to_owned(), script_in_container()];
153    command.extend(harness.iter().cloned());
154    command
155}
156
157/// How long a connection may take to send its head line and body, and how long QCode waits for
158/// the provider to answer before giving up on it. Generous, because the provider may have to
159/// load a model first; finite, because a container waiting forever is a container stuck.
160const PATIENCE: Duration = Duration::from_secs(120);
161
162/// The longest a head line may be. Wide enough for a generous set of headers, narrow enough that
163/// a line this long is never mistaken for a request that meant to carry a body on it.
164const MOST_HEAD: usize = 64 * 1024;
165
166/// The most connections served at once, mirroring [`crate::bridge::socket`]'s limit: one per tab
167/// asking at a time, and the rest are something else knocking.
168const MOST_CONNECTIONS: usize = 16;
169
170/// The most a request's body may be, and the reason it is read whole rather than handed on as it
171/// arrives: the model in it is set to the one the profile chose before anything goes out, and
172/// every step of a lineup must be asked the very same bytes. A harness's own request is a
173/// conversation and far smaller, so the bound is against a body that never ends rather than
174/// against a real one.
175const MOST_BODY: usize = 64 * 1024 * 1024;
176
177/// How much of an answer's body is read before its head is decided on: enough of a `400` or a `404`
178/// to see whether it names the model that was asked, and as much of a passing answer as is kept in
179/// case that one is the answer the harness is shown. Wide enough for the sentence any provider has
180/// been seen to put in one, and a bound rather than a policy: what is not in the first 64 KiB of a
181/// complaint is not in the complaint.
182const MOST_SAID: usize = 64 * 1024;
183
184/// How long a step that answered with a reason to fall back is not asked again. Long enough that
185/// a busy free model stops costing a request per turn, short enough that a model which is back is
186/// used without the person thinking about it.
187const COOLING: Duration = Duration::from_secs(60);
188
189/// The headers a request or an answer never carries across the relay because the framing on
190/// either side already speaks for them, or because letting them through would carry someone
191/// else's idea of authentication into a request QCode is about to add its own key to.
192///
193/// `api-key` is among them because it is the header Xiaomi reads a key from: a harness that was
194/// given one of its own must not have it reach a provider beside the key QCode adds.
195const STRIPPED: [&str; 7] =
196    ["authorization", "x-api-key", "api-key", "host", "content-length", "transfer-encoding", "connection"];
197
198/// The headers a replayed answer does not carry. Its body was read before the step was given up on
199/// and may have been cut at [`MOST_SAID`], so neither the length the provider measured nor a
200/// framing of its own goes with it: the script in the container writes the head and the bytes as
201/// they come, and a length that does not match is an answer the harness reads short.
202const REPLAYED_WITHOUT: [&str; 2] = ["content-length", "transfer-encoding"];
203
204/// One request read off the socket, before it is decided whether to carry it anywhere.
205#[derive(Debug, Clone, PartialEq, Eq)]
206struct Incoming {
207    /// The token of the tab asking, as the relay script found it.
208    token: String,
209    /// `GET` or `POST`; anything else is refused before a provider is even looked up.
210    method: String,
211    /// The path asked for, exactly as the harness sent it, query string included.
212    path: String,
213    /// The headers the harness sent, already without its own `authorization` or `x-api-key`.
214    headers: Vec<(String, String)>,
215}
216
217/// A head line that was not one.
218#[derive(Debug, Clone, Copy, PartialEq, Eq)]
219struct Malformed;
220
221/// Reads one head line's worth of JSON.
222fn parse_incoming(line: &str) -> Result<Incoming, Malformed> {
223    let value: Value = serde_json::from_str(line).map_err(|_| Malformed)?;
224    let object = value.as_object().ok_or(Malformed)?;
225    let text = |key: &str| object.get(key).and_then(Value::as_str).map(str::to_owned).ok_or(Malformed);
226    let token = text("token")?;
227    let method = text("method")?;
228    let path = text("path")?;
229    let headers = match object.get("headers") {
230        None => Vec::new(),
231        Some(value) => {
232            let map = value.as_object().ok_or(Malformed)?;
233            map.iter()
234                .map(|(name, value)| value.as_str().map(|value| (name.to_lowercase(), value.to_owned())))
235                .collect::<Option<Vec<_>>>()
236                .ok_or(Malformed)?
237        }
238    };
239    Ok(Incoming { token, method, path, headers })
240}
241
242/// Where Codex sends its conversation: OpenAI's Responses shape, the only one Codex still speaks
243/// to a provider of one's own (its 0.156.1 program answers `wire_api = "chat"` with "is no longer
244/// supported"). It is a conversation in the OpenAI shape like the completion path, so it goes to
245/// the same root.
246pub const RESPONSES_PATH: &str = "/v1/responses";
247
248/// The path without its query string, which is what a rule about what may be asked checks
249/// against; the query string still travels with the request that is actually sent.
250fn path_only(path: &str) -> &str {
251    path.split('?').next().unwrap_or(path)
252}
253
254/// Whether `entry` may be asked for `method` `path` through the relay. See the module's doc
255/// comment for why these two and nothing else.
256fn allowed(method: &str, path: &str) -> bool {
257    let path = path_only(path);
258    match method {
259        "POST" => super::Wire::ALL.iter().any(|wire| path == wire.messages_path()) || path == RESPONSES_PATH,
260        "GET" => path == "/v1/models",
261        _ => false,
262    }
263}
264
265/// A line of JSON answering the container, with `status` and the headers the provider (or the
266/// relay itself, for a refusal) answered with.
267fn head_line(status: u16, headers: &[(String, String)]) -> String {
268    let mut object = Map::new();
269    object.insert("status".to_owned(), Value::from(status));
270    let headers: Map<String, Value> = headers.iter().map(|(name, value)| (name.clone(), json!(value))).collect();
271    object.insert("headers".to_owned(), Value::Object(headers));
272    let mut line = Value::Object(object).to_string();
273    line.push('\n');
274    line
275}
276
277/// A short refusal, written back as a whole answer: the head line, a small JSON body naming why,
278/// nothing else. Never the request that was refused and never a key: refusing a connection is not
279/// a place either could belong.
280fn refusal_body(why: &str) -> Vec<u8> {
281    json!({ "error": why }).to_string().into_bytes()
282}
283
284/// `count` as the bound [`Read::take`] takes, for a constant written as a `usize`.
285fn bound(count: usize) -> u64 {
286    u64::try_from(count).unwrap_or(u64::MAX)
287}
288
289/// `body` with its `model` set to `model`, or `body` itself unchanged when it is not a JSON
290/// object naming one: a body QCode cannot read the model out of is a body QCode does not touch,
291/// which is what a `GET /v1/models` and anything not written in the shapes QCode knows arrive as.
292///
293/// The rest of the body goes on exactly as the harness wrote it. Only the model is QCode's
294/// business, and the model is the person's choice rather than the harness's.
295fn pinned(body: Vec<u8>, model: &str) -> Vec<u8> {
296    let Ok(Value::Object(mut fields)) = serde_json::from_slice::<Value>(&body) else { return body };
297    if !fields.contains_key("model") {
298        return body;
299    }
300    fields.insert("model".to_owned(), Value::from(model));
301    serde_json::to_vec(&Value::Object(fields)).unwrap_or(body)
302}
303
304/// One request sent on to a provider, with its key already decided.
305pub struct UpstreamAsk {
306    /// `GET` for a model listing, `POST` for a message.
307    pub method: Method,
308    /// Where it goes, the provider's base with the requested path (and any query) after it.
309    pub url: String,
310    /// The headers the harness sent, minus the ones the relay never carries across.
311    pub headers: Vec<(String, String)>,
312    /// The provider's key, for a provider that needs one.
313    pub secret: Option<Secret>,
314    /// The request's body, read as it is sent rather than collected first.
315    pub body: Box<dyn Read + Send>,
316}
317
318/// What a provider answered, read incrementally so an answer streamed as server-sent events is
319/// never held back waiting for the rest of it.
320pub struct UpstreamAnswer {
321    /// The status the provider answered with.
322    pub status: u16,
323    /// The headers to carry back, minus the ones framing already speaks for.
324    pub headers: Vec<(String, String)>,
325    /// The body, read as it arrives.
326    pub body: Box<dyn Read + Send>,
327}
328
329/// Why a provider could not be asked.
330#[derive(Debug, Clone, PartialEq, Eq)]
331pub struct UpstreamError {
332    /// Where the request was going. Never the key: a key is a header value, never part of an
333    /// address, so naming where a request went never names what it carried.
334    pub url: String,
335    /// What went wrong, in the transport's own words.
336    pub reason: String,
337}
338
339/// How a request is carried to the provider.
340///
341/// [`network`](Upstream::network) is the real one. [`new`](Upstream::new) takes anything else,
342/// which is what every test here uses, the same seam [`super::ask::Web`] gives the rest of this
343/// crate: no test in this module reaches the network, whatever it asks for.
344#[derive(Clone)]
345pub struct Upstream(Arc<dyn Fn(UpstreamAsk) -> Result<UpstreamAnswer, UpstreamError> + Send + Sync>);
346
347impl std::fmt::Debug for Upstream {
348    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
349        formatter.write_str("Upstream")
350    }
351}
352
353impl Upstream {
354    /// The real one, which goes out over the network and streams the answer back rather than
355    /// reading it whole.
356    #[must_use]
357    pub fn network() -> Self {
358        Self::new(|ask| {
359            let config =
360                ureq::Agent::config_builder().timeout_global(Some(PATIENCE)).http_status_as_error(false).build();
361            let agent = ureq::Agent::new_with_config(config);
362            let mut builder = match ask.method {
363                Method::Get => agent.get(&ask.url).force_send_body(),
364                Method::Post => agent.post(&ask.url),
365            };
366            for (name, value) in &ask.headers {
367                builder = builder.header(name, value);
368            }
369            if let Some(secret) = &ask.secret {
370                builder = builder.header(&secret.header, &format!("{}{}", secret.prefix, secret.key.expose()));
371            }
372            let sent = match ask.method {
373                Method::Get => builder.send_empty(),
374                Method::Post => builder.send(ureq::SendBody::from_owned_reader(ask.body)),
375            };
376            let answer = sent.map_err(|error| UpstreamError { url: ask.url.clone(), reason: error.to_string() })?;
377            let status = answer.status().as_u16();
378            let headers = answer
379                .headers()
380                .iter()
381                .filter_map(|(name, value)| {
382                    let name = name.as_str().to_lowercase();
383                    (!STRIPPED.contains(&name.as_str())).then(|| Some((name, value.to_str().ok()?.to_owned())))?
384                })
385                .collect();
386            let body = answer.into_body().into_reader();
387            Ok(UpstreamAnswer { status, headers, body: Box::new(body) })
388        })
389    }
390
391    /// An upstream that carries requests the way `carry` says.
392    #[must_use]
393    pub fn new(carry: impl Fn(UpstreamAsk) -> Result<UpstreamAnswer, UpstreamError> + Send + Sync + 'static) -> Self {
394        Self(Arc::new(carry))
395    }
396
397    /// Asks `carry` to send `ask`.
398    pub(super) fn call(&self, ask: UpstreamAsk) -> Result<UpstreamAnswer, UpstreamError> {
399        (self.0)(ask)
400    }
401}
402
403/// Something the relay found out, in a shape a screen can show: a request refused, one carried
404/// through, or a provider that could not be reached. Never printed on its own; the screen that
405/// opened the workspace decides what the person sees and when.
406#[derive(Debug, Clone, PartialEq, Eq)]
407pub enum Event {
408    /// A connection carried no token, or one no open tab answers to.
409    UnknownTab,
410    /// The head line the relay script sent was not one QCode reads.
411    Malformed,
412    /// A connection asked for a method or a path the relay does not carry.
413    NotAllowed {
414        /// The method that was asked for.
415        method: String,
416        /// The path that was asked for.
417        path: String,
418    },
419    /// The provider could not be reached, or refused the transport itself.
420    Unreachable {
421        /// The tag of the provider that was asked.
422        tag: String,
423        /// What the transport said.
424        reason: String,
425    },
426    /// A request was carried to the provider and its answer's head came back.
427    Forwarded {
428        /// The tag of the provider that answered.
429        tag: String,
430        /// The path that was asked for.
431        path: String,
432        /// The status the provider answered with.
433        status: u16,
434    },
435    /// One model of a lineup was left for the next: a person who chose an order of models should
436    /// see which one was passed over and which one is being asked instead, since only QCode knows
437    /// that both happened.
438    FellBack {
439        /// The token of the tab whose request this was, which is the only thing that says whose
440        /// conversation fell back: two tabs of one workspace can be on different providers, and a
441        /// screen writes this under the tab it belongs to rather than wherever it happens to be.
442        token: String,
443        /// The tag of the provider that was asked.
444        tag: String,
445        /// The lineup's name, when the profile chose one rather than a single model.
446        lineup: Option<String>,
447        /// The model that was asked and passed over.
448        from: String,
449        /// The model that is asked instead.
450        to: String,
451        /// The status the first step answered with, or `None` when it could not be reached at all.
452        status: Option<u16>,
453    },
454}
455
456/// What one token's requests are doing at this moment: how many of them the relay is serving, and
457/// when the last of them finished.
458#[derive(Debug, Clone, Copy, Default)]
459struct Asking {
460    /// How many requests of this token are being served right now.
461    serving: usize,
462    /// When the last request of this token finished, however it ended.
463    last_finished: Option<Instant>,
464}
465
466/// What the relay is doing for each of a workspace's tabs, read from anywhere and cheaply.
467///
468/// The screen has to be able to tell a tab that is waiting on a model from one that has been
469/// quiet for a while, and the relay is the only side of a request that knows: the container's
470/// script hands it over and gets the answer back, so nothing between the harness and the provider
471/// is left holding a moment in time. A clone is a second pointer at the same counts, so a screen
472/// and the relay share one truth rather than two copies that can drift apart.
473#[derive(Debug, Clone, Default)]
474pub struct Activity(Arc<Mutex<HashMap<String, Asking>>>);
475
476impl Activity {
477    /// Whether a request of `token` is on its way to a provider right now.
478    ///
479    /// A lock that cannot be taken says it is not. A screen reading this is answering a person who
480    /// is waiting, and waiting on a count that a poisoned lock would hold forever helps nobody.
481    #[must_use]
482    pub fn busy(&self, token: &str) -> bool {
483        self.0.lock().is_ok_and(|asking| asking.get(token).is_some_and(|asking| asking.serving > 0))
484    }
485
486    /// When the last request of `token` finished, or `None` when it has not finished one: a token
487    /// no tab has, and a token whose every request was refused before it reached a route, both
488    /// have no time.
489    #[must_use]
490    pub fn last_finished(&self, token: &str) -> Option<Instant> {
491        self.0.lock().ok()?.get(token)?.last_finished
492    }
493}
494
495/// One request being served, counted in for as long as it lasts and out again when it ends.
496///
497/// The count is taken out in [`Drop`], so it goes however the serving returns: an answer written,
498/// a refusal, an early return, or a panic on the way. A request left counted in would keep its tab
499/// busy for good, which is the one thing a screen reading [`Activity`] must never be told.
500struct Counting {
501    asking: Arc<Mutex<HashMap<String, Asking>>>,
502    token: String,
503}
504
505impl Counting {
506    /// Counts one request of `token` in, from now until what this returns is dropped.
507    fn serving(activity: &Activity, token: &str) -> Self {
508        let counting = Self { asking: Arc::clone(&activity.0), token: token.to_owned() };
509        // A lock that cannot be taken is not a reason to refuse the request: the request is served
510        // either way, and the count is then whatever the lock holds.
511        if let Ok(mut asking) = counting.asking.lock() {
512            asking.entry(counting.token.clone()).or_default().serving += 1;
513        }
514        counting
515    }
516}
517
518impl Drop for Counting {
519    fn drop(&mut self) {
520        let Ok(mut asking) = self.asking.lock() else { return };
521        let Some(asked) = asking.get_mut(&self.token) else { return };
522        asked.serving = asked.serving.saturating_sub(1);
523        asked.last_finished = Some(Instant::now());
524    }
525}
526
527/// The workspace's provider relay, listened on for as long as this value lives.
528#[derive(Debug)]
529pub struct Listener {
530    socket: PathBuf,
531    activity: Activity,
532    #[cfg(unix)]
533    open: Arc<std::sync::atomic::AtomicBool>,
534}
535
536impl Listener {
537    /// Listens in `folder`, the workspace's `Containers/MCP/`, beside the bridge's socket: makes
538    /// the folder if it is not there already, writes the relay script, and starts waiting for
539    /// connections. `resolve` finds the route a token belongs to — the provider to ask and the
540    /// models to ask it with; `upstream` carries the request the route describes; `report` is
541    /// told what happened to each connection, in a shape a screen can show.
542    ///
543    /// A socket left behind by a QCode that ended without clearing it is replaced; one another
544    /// QCode still answers on is left to it.
545    ///
546    /// # Errors
547    ///
548    /// When the folder or the script cannot be written, another QCode answers on the socket
549    /// ([`io::ErrorKind::AddrInUse`]), the socket cannot be made, or the system has none.
550    pub fn open(
551        folder: &Path,
552        resolve: impl Fn(&str) -> Option<Route> + Send + Sync + 'static,
553        upstream: Upstream,
554        report: impl Fn(Event) + Send + Sync + 'static,
555    ) -> io::Result<Self> {
556        #[cfg(unix)]
557        {
558            unix::open(folder, resolve, upstream, report)
559        }
560        #[cfg(not(unix))]
561        {
562            let _ = (folder, resolve, upstream, report);
563            Err(io::Error::new(io::ErrorKind::Unsupported, "this system has no unix sockets"))
564        }
565    }
566
567    /// Where the socket is.
568    #[must_use]
569    pub fn socket(&self) -> &Path {
570        &self.socket
571    }
572
573    /// What this relay is doing for each of the workspace's tokens, as a handle the screen keeps.
574    ///
575    /// Cheap to clone and safe to keep: a clone held after this listener is closed still answers,
576    /// with the last thing that was true.
577    #[must_use]
578    pub fn activity(&self) -> Activity {
579        self.activity.clone()
580    }
581}
582
583/// Writes the relay's script into `folder` unless it is there already as this QCode carries it.
584fn write_script(folder: &Path) -> io::Result<()> {
585    let path = folder.join(SCRIPT_NAME);
586    if std::fs::read(&path).is_ok_and(|written| written == SCRIPT.as_bytes()) {
587        return Ok(());
588    }
589    qframe::storage::atomic_write(&path, SCRIPT.as_bytes())
590}
591
592#[cfg(unix)]
593mod unix {
594    use std::io::{self, BufRead, BufReader, Read, Write};
595    use std::os::unix::fs::PermissionsExt;
596    use std::os::unix::net::{UnixListener, UnixStream};
597    use std::path::{Path, PathBuf};
598    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
599    use std::sync::{Arc, Mutex};
600    use std::time::Instant;
601
602    use super::super::Wire;
603    use super::super::ask::{Method, secret_of};
604    use super::{
605        Activity, Cooling, Counting, Event, Listener, MOST_BODY, MOST_CONNECTIONS, MOST_HEAD, MOST_SAID, PATIENCE,
606        Passing, Route, SOCKET_NAME, STRIPPED, Upstream, UpstreamAsk, allowed, bound, cool, head_line, lasting,
607        order_now, parse_incoming, passing, pinned, refusal_body, write_script,
608    };
609
610    /// The longest socket path the system takes, less the byte its terminator needs; the same
611    /// limit [`crate::bridge::socket`] works around, for the same reason.
612    const MOST_PATH: usize = 107;
613
614    pub(super) fn open(
615        folder: &Path,
616        resolve: impl Fn(&str) -> Option<Route> + Send + Sync + 'static,
617        upstream: Upstream,
618        report: impl Fn(Event) + Send + Sync + 'static,
619    ) -> io::Result<Listener> {
620        std::fs::create_dir_all(folder)?;
621        std::fs::set_permissions(folder, std::fs::Permissions::from_mode(0o700))?;
622        write_script(folder)?;
623        let socket = folder.join(SOCKET_NAME);
624        if socket.exists() {
625            if reach(&socket, |path| UnixStream::connect(path)).is_ok() {
626                return Err(io::Error::new(io::ErrorKind::AddrInUse, socket.display().to_string()));
627            }
628            std::fs::remove_file(&socket)?;
629        }
630        let listener = reach(&socket, |path| UnixListener::bind(path))?;
631        let open = Arc::new(AtomicBool::new(true));
632        let still_open = Arc::clone(&open);
633        let resolve = Arc::new(resolve);
634        let report = Arc::new(report);
635        // Which models of this workspace have refused, and until when each is not asked again.
636        // Held by the thread that waits for connections and shared with the ones that serve them,
637        // so it lives exactly as long as the relay listens and every connection of this relay
638        // shares it. See [`ordered`].
639        let cooling = Arc::new(Mutex::new(Cooling::new()));
640        let cooling_in_threads = Arc::clone(&cooling);
641        // Held by every thread that serves a connection and by whoever asked to watch, so a screen
642        // and the relay read the same counts.
643        let activity = Activity::default();
644        let watching = activity.clone();
645        std::thread::Builder::new().name("qcode-relay".to_owned()).spawn(move || {
646            accept(&listener, &still_open, &resolve, &upstream, &cooling_in_threads, &report, &watching);
647        })?;
648        Ok(Listener { socket, activity, open })
649    }
650
651    /// Runs `with` on the socket path, or, when the path is too long for a socket address, on the
652    /// same socket named through the folder's open handle in `/proc/self/fd`.
653    fn reach<T>(socket: &Path, with: impl Fn(&Path) -> io::Result<T>) -> io::Result<T> {
654        if socket.as_os_str().len() <= MOST_PATH {
655            return with(socket);
656        }
657        let folder = socket.parent().ok_or_else(|| io::Error::from(io::ErrorKind::InvalidInput))?;
658        let handle = std::fs::File::open(folder)?;
659        let short = PathBuf::from(format!("/proc/self/fd/{}/{SOCKET_NAME}", std::os::fd::AsRawFd::as_raw_fd(&handle)));
660        let result = with(&short);
661        drop(handle);
662        result
663    }
664
665    /// Waits for connections until the listener is closed, serving each on a thread of its own.
666    fn accept(
667        listener: &UnixListener,
668        open: &AtomicBool,
669        resolve: &Arc<impl Fn(&str) -> Option<Route> + Send + Sync + 'static>,
670        upstream: &Upstream,
671        cooling: &Arc<Mutex<Cooling>>,
672        report: &Arc<impl Fn(Event) + Send + Sync + 'static>,
673        activity: &Activity,
674    ) {
675        let serving = Arc::new(AtomicUsize::new(0));
676        for stream in listener.incoming() {
677            if !open.load(Ordering::SeqCst) {
678                break;
679            }
680            let Ok(stream) = stream else { continue };
681            if serving.load(Ordering::SeqCst) >= MOST_CONNECTIONS {
682                continue;
683            }
684            serving.fetch_add(1, Ordering::SeqCst);
685            let (resolve, upstream, cooling, report, done, activity) = (
686                Arc::clone(resolve),
687                upstream.clone(),
688                Arc::clone(cooling),
689                Arc::clone(report),
690                Arc::clone(&serving),
691                activity.clone(),
692            );
693            let spawned = std::thread::Builder::new().name("qcode-relay-call".to_owned()).spawn(move || {
694                serve(stream, resolve.as_ref(), &upstream, &cooling, report.as_ref(), &activity);
695                done.fetch_sub(1, Ordering::SeqCst);
696            });
697            if spawned.is_err() {
698                serving.fetch_sub(1, Ordering::SeqCst);
699            }
700        }
701    }
702
703    /// Reads one request off `stream`, decides whether it may be carried, and either refuses it
704    /// or forwards it, walking down the route's steps until one of them answers, and streams that
705    /// answer back.
706    fn serve(
707        stream: UnixStream,
708        resolve: &(impl Fn(&str) -> Option<Route> + ?Sized),
709        upstream: &Upstream,
710        cooling: &Arc<Mutex<super::Cooling>>,
711        report: &(impl Fn(Event) + ?Sized),
712        activity: &Activity,
713    ) {
714        if stream.set_read_timeout(Some(PATIENCE)).is_err() {
715            return;
716        }
717        let Ok(writing) = stream.try_clone() else { return };
718        let mut writing = writing;
719        let mut reading = BufReader::new(stream);
720
721        let mut line = String::new();
722        let limit = u64::try_from(MOST_HEAD).unwrap_or(u64::MAX) + 1;
723        let read = (&mut reading).take(limit).read_line(&mut line);
724        let Ok(incoming) = (match read {
725            Ok(_) if line.ends_with('\n') => parse_incoming(line.trim_end_matches(['\n', '\r'])),
726            _ => Err(super::Malformed),
727        }) else {
728            report(Event::Malformed);
729            refuse(&mut writing, 400, "malformed");
730            return;
731        };
732
733        let Some(route) = resolve(&incoming.token) else {
734            report(Event::UnknownTab);
735            refuse(&mut writing, 401, "unknown tab");
736            return;
737        };
738        // From here the token is one a tab really has, so whatever this connection goes on to do
739        // counts as that tab asking: a request refused for its method or its path was still made,
740        // and a screen watching this is told the tab was busy rather than that it was idle.
741        let _asking = Counting::serving(activity, &incoming.token);
742        let Some(method) = (match incoming.method.as_str() {
743            "GET" => Some(Method::Get),
744            "POST" => Some(Method::Post),
745            _ => None,
746        }) else {
747            report(Event::NotAllowed { method: incoming.method.clone(), path: incoming.path.clone() });
748            refuse(&mut writing, 405, "method not served here");
749            return;
750        };
751        if !allowed(&incoming.method, &incoming.path) {
752            report(Event::NotAllowed { method: incoming.method.clone(), path: incoming.path.clone() });
753            refuse(&mut writing, 404, "not served here");
754            return;
755        }
756        // A route with no step to ask has nowhere to go. A profile that names a model brings one,
757        // and one that does not is QCode's own doing rather than the harness's.
758        if route.models.is_empty() {
759            refuse(&mut writing, 500, "the profile names no model");
760            return;
761        }
762
763        let headers: Vec<(String, String)> =
764            incoming.headers.into_iter().filter(|(name, _)| !STRIPPED.contains(&name.as_str())).collect();
765        let secret = secret_of(&route.entry);
766        // Under the provider's API rather than its bare address: a harness asks for
767        // `/v1/messages` the way it would ask Anthropic, OpenRouter answers that path only under
768        // `/api`, and Xiaomi only under `/anthropic` while it answers the other shape at its root.
769        let url = route.entry.api_address(Wire::of_path(&incoming.path), &incoming.path);
770
771        // Read whole rather than carried as it arrives: the model in the body is set to the one
772        // the person chose before anything goes out, and every step of a lineup is asked the very
773        // same bytes. Over the limit is refused here rather than asked about, since a provider has
774        // no business answering a request QCode itself will not carry.
775        let mut body = Vec::new();
776        let read = reading.take(bound(MOST_BODY).saturating_add(1)).read_to_end(&mut body);
777        if read.is_err() {
778            refuse(&mut writing, 400, "the request body could not be read");
779            return;
780        }
781        if body.len() > MOST_BODY {
782            refuse(&mut writing, 413, "the request body is too large");
783            return;
784        }
785
786        // The order the lineup is tried in, with the models that just refused at the back for a
787        // minute, so a busy one stops costing a request per turn. Fixed before the first ask, so
788        // every step of this request is asked of the same models in the same order.
789        let asked_at = Instant::now();
790        let key = |model: &str| (route.entry.tag.to_string(), route.lineup.clone(), model.to_owned());
791        let steps = order_now(cooling, &route, &key, asked_at);
792
793        let mut waiting = steps.iter();
794        // The first answer of this request that only failed for now, kept in case every step after
795        // it fails for good. Nothing is decided on it while there is still a step to ask: the step
796        // after a busy one may well answer, and an answer nobody was shown is not a claim.
797        let mut remembered: Option<Passing> = None;
798        while let Some(model) = waiting.next() {
799            let to = waiting.clone().next().map(String::as_str);
800            let ask = UpstreamAsk {
801                method,
802                url: url.clone(),
803                headers: headers.clone(),
804                secret: secret.clone(),
805                body: Box::new(io::Cursor::new(pinned(body.clone(), model))),
806            };
807            match upstream.call(ask) {
808                Ok(mut answer) => {
809                    // A `400` and a `404` only say which model they are about in what they say, and
810                    // an answer that may be kept has to be held rather than streamed, so that much of
811                    // the body is read before the head is decided on. What is read is kept: a step
812                    // that turns out not to be a reason to fall back still answers the harness whole,
813                    // and so does the last step of all.
814                    let mut said = Vec::new();
815                    if matches!(answer.status, 400 | 404) || (to.is_some() && passing(answer.status)) {
816                        let _ = answer.body.by_ref().take(bound(MOST_SAID)).read_to_end(&mut said);
817                    }
818                    let status = answer.status;
819                    // Only ever on the answer's head. Once a `2xx` has started, the body goes to the
820                    // harness as it arrives and this request's model is decided for good: a
821                    // stream that breaks halfway is a thing the harness must see, not a reason to
822                    // spend the same question twice.
823                    let giving_up = to.is_some() && (passing(status) || lasting(status, &said, model));
824                    if giving_up {
825                        let Some(to) = to else { return };
826                        cool(cooling, &key, model);
827                        report(Event::FellBack {
828                            token: incoming.token.clone(),
829                            tag: route.entry.tag.to_string(),
830                            lineup: route.lineup.clone(),
831                            from: (*model).to_owned(),
832                            to: to.to_owned(),
833                            status: Some(status),
834                        });
835                        // The first passing answer is kept: it is the model the person wrote first,
836                        // and a later one says less about what this request ran into.
837                        if remembered.is_none() && passing(status) {
838                            remembered = Some(Passing::Answer { status, headers: answer.headers, said });
839                        }
840                        continue;
841                    }
842                    // What is left is a step that is not given up on: the last of the request,
843                    // and an answer that is not a reason to fall back. That answer is the harness's as
844                    // it is, except where this request already produced a reason to wait and this one
845                    // will not go away by waiting: a passing reason tells the person to wait and a
846                    // lasting one tells them to change the lineup, and when one request ran into
847                    // both, the reason to wait is the true story of it. Nothing is cooled and no
848                    // fall back is said: the walking stops here either way, and what the last step
849                    // answered is still what happened.
850                    if let Some(kept) = remembered.filter(|_| lasting(status, &said, model)) {
851                        match kept.reframed() {
852                            Passing::Answer { status: shown, headers, said: body } => {
853                                report(Event::Forwarded {
854                                    tag: route.entry.tag.to_string(),
855                                    path: incoming.path,
856                                    status: shown,
857                                });
858                                let head = head_line(shown, &headers);
859                                if writing.write_all(head.as_bytes()).is_ok() {
860                                    let _ = writing.write_all(&body);
861                                }
862                            }
863                            Passing::Unreachable { reason } => {
864                                report(Event::Unreachable { tag: route.entry.tag.to_string(), reason });
865                                refuse(&mut writing, 502, "the provider could not be reached");
866                            }
867                        }
868                        return;
869                    }
870                    report(Event::Forwarded { tag: route.entry.tag.to_string(), path: incoming.path, status });
871                    let head = head_line(status, &answer.headers);
872                    if writing.write_all(head.as_bytes()).is_ok() {
873                        // The part of the body read to be decided on, then the rest of it as it
874                        // arrives, so an answer that starts is streamed all the same.
875                        let mut body = answer.body;
876                        if writing.write_all(&said).is_ok() {
877                            let _ = io::copy(&mut body, &mut writing);
878                        }
879                    }
880                    return;
881                }
882                Err(error) => {
883                    // The last step's own answer, whatever it is, is the harness's: a provider
884                    // that cannot be reached has said so in its own words, and a harness can show
885                    // those where a refusal QCode made up would only hide them. Being unreachable is
886                    // itself a reason to come back rather than a complaint about the model, so this
887                    // one is shown as it is and nothing kept from an earlier step stands in its place.
888                    let Some(to) = to else {
889                        report(Event::Unreachable { tag: route.entry.tag.to_string(), reason: error.reason });
890                        refuse(&mut writing, 502, "the provider could not be reached");
891                        return;
892                    };
893                    cool(cooling, &key, model);
894                    report(Event::FellBack {
895                        token: incoming.token.clone(),
896                        tag: route.entry.tag.to_string(),
897                        lineup: route.lineup.clone(),
898                        from: (*model).to_owned(),
899                        to: to.to_owned(),
900                        status: None,
901                    });
902                    // Kept like an answer, and replayed as the same refusal the last step's own
903                    // unreachable writes, since that is all there was to say about this provider.
904                    if remembered.is_none() {
905                        remembered = Some(Passing::Unreachable { reason: error.reason });
906                    }
907                }
908            }
909        }
910    }
911
912    /// Writes a short refusal and nothing else.
913    fn refuse(writing: &mut UnixStream, status: u16, why: &str) {
914        let head = head_line(status, &[("content-type".to_owned(), "application/json".to_owned())]);
915        if writing.write_all(head.as_bytes()).is_ok() {
916            let _ = writing.write_all(&refusal_body(why));
917        }
918    }
919
920    impl Drop for Listener {
921        /// Stops listening and takes the socket away, the same way [`crate::bridge::socket::
922        /// Listener`] does: one last connection of its own wakes the accepting thread so it sees
923        /// the listener is closed.
924        fn drop(&mut self) {
925            self.open.store(false, Ordering::SeqCst);
926            let _ = reach(&self.socket, |path| UnixStream::connect(path));
927            let _ = std::fs::remove_file(&self.socket);
928        }
929    }
930}
931
932/// One step of one route, named the way the cool-down names it: a model is only skipped for the
933/// provider tag and the lineup it refused in, since the same name on someone else's provider is a
934/// different machine with different models.
935type Step = (String, Option<String>, String);
936
937/// A step of a route and until when it is not to be asked.
938type Cooling = HashMap<Step, Instant>;
939
940/// The steps to ask at `now`: those not being skipped first, and the ones being skipped after
941/// them, each group in the order the profile wrote them. When every step is being skipped the
942/// order is the plain one — a model that keeps refusing is still the one the person chose, and
943/// asking it beats an answer from nobody.
944///
945/// A pure function of the clock rather than of the wall, so what a minute of skipping looks like
946/// is settled by looking at a minute later rather than by waiting for one.
947fn ordered(models: &[String], cooling: &Cooling, key: &dyn Fn(&str) -> Step, now: Instant) -> Vec<String> {
948    let (mut ready, mut waiting): (Vec<String>, Vec<String>) =
949        models.iter().cloned().partition(|model| cooling.get(&key(model)).is_none_or(|until| *until <= now));
950    ready.append(&mut waiting);
951    ready
952}
953
954/// The steps of `route` in the order they are asked at `now`, with the models that refused
955/// recently at the back until their cool-down is up and the ones that have run out of it forgotten.
956/// A lock that cannot be taken is not a reason to refuse a request: the order is then the plain
957/// one, which is what a relay with no memory of any refusal would do.
958fn order_now(cooling: &Mutex<Cooling>, route: &Route, key: &dyn Fn(&str) -> Step, now: Instant) -> Vec<String> {
959    let Ok(mut seen) = cooling.lock() else { return route.models.clone() };
960    seen.retain(|_, until| *until > now);
961    ordered(&route.models, &seen, key, now)
962}
963
964/// Whether `said` is an error body that names `model`, which is what makes a `400` or a `404` a
965/// complaint about that model rather than about the request. Read in bytes rather than in words:
966/// what a provider says about a model it does not have quotes the name, in a shape nobody can
967/// promise to write a rule for.
968fn names(said: &[u8], model: &str) -> bool {
969    let model = model.as_bytes();
970    !model.is_empty() && said.windows(model.len()).any(|window| window == model)
971}
972
973/// Whether `status` is a reason to come back rather than a complaint that will not go away: a
974/// provider that is busy, one that was too slow, and one that broke. Another model of the lineup is
975/// worth asking now, and so is this one in a minute, which is what makes such an answer the one to
976/// keep when every step that comes after it fails for good.
977fn passing(status: u16) -> bool {
978    matches!(status, 408 | 429) || status >= 500
979}
980
981/// Whether `status` is a complaint about the model that was asked, as `said` read up to [`MOST_SAID`]
982/// says: a provider that will not sell it, will not let this key near it, or has never heard of
983/// it. Such an answer does not change by waiting, so it is not what a person waiting is shown.
984///
985/// A `401` is none of these: there is one key and the next step would be refused the same way. Nor
986/// is a `400` that does not name the model — that is the request's own fault, and every model of the
987/// lineup would refuse it the same way.
988fn lasting(status: u16, said: &[u8], model: &str) -> bool {
989    matches!(status, 402 | 403) || (matches!(status, 400 | 404) && names(said, model))
990}
991
992/// An earlier step's answer, kept while the steps after it are asked in case none of them can answer:
993/// then this is what the harness is shown rather than a complaint that will not go away.
994enum Passing {
995    /// A provider that answered with a reason to wait, in its own status, its own headers and as much
996    /// of its own body as was read before the step was given up on.
997    Answer {
998        /// The status the provider answered with.
999        status: u16,
1000        /// The headers it answered with, minus the ones [`REPLAYED_WITHOUT`] names.
1001        headers: Vec<(String, String)>,
1002        /// Its body as it was read.
1003        said: Vec<u8>,
1004    },
1005    /// A provider that could not be reached at all, in the transport's own words. The harness is shown
1006    /// the relay's own refusal in its place, exactly as it is when the last step could not be reached.
1007    Unreachable {
1008        /// What the transport said.
1009        reason: String,
1010    },
1011}
1012
1013impl Passing {
1014    /// Takes away the headers of a replayed answer that do not go with the body it is shown behind.
1015    /// Called on the way out of the listener rather than on the way in, so what a step was given up
1016    /// on for is never touched.
1017    fn reframed(mut self) -> Self {
1018        if let Self::Answer { headers, .. } = &mut self {
1019            headers.retain(|(name, _)| !REPLAYED_WITHOUT.contains(&name.as_str()));
1020        }
1021        self
1022    }
1023}
1024
1025/// Remembers that `model` is not to be asked of this route again for [`COOLING`]. Nothing is done
1026/// with a lock that cannot be taken: the step is asked again on the next request rather than
1027/// never, which is what a plain order does.
1028fn cool(cooling: &Mutex<Cooling>, key: &dyn Fn(&str) -> Step, model: &str) {
1029    let Ok(mut cooling) = cooling.lock() else { return };
1030    cooling.insert(key(model), Instant::now() + COOLING);
1031}
1032
1033#[cfg(all(test, unix))]
1034mod tests {
1035    use std::io::{BufRead, BufReader, Read, Write};
1036    use std::os::unix::net::UnixStream;
1037    use std::path::PathBuf;
1038    use std::sync::atomic::{AtomicUsize, Ordering};
1039    use std::sync::mpsc::{self, Receiver, Sender};
1040    use std::sync::{Arc, Condvar, Mutex};
1041
1042    use serde_json::Value;
1043
1044    use super::super::{Key, ProviderKind, Tag};
1045    use super::*;
1046
1047    /// Not a key of anyone's: the characters spell what it is.
1048    const MADE_UP_KEY: &str = "not-a-real-key-0000-wxyz";
1049
1050    /// The model a profile is taken to have chosen, for the tests that are about the road and not
1051    /// about which model: no default, so a test that forgets to pin one cannot pass quietly.
1052    const CHOSEN: &str = "qwen/qwen3-coder:free";
1053
1054    struct Scratch(PathBuf);
1055
1056    impl Scratch {
1057        fn new(name: &str) -> Self {
1058            let stamp =
1059                std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_nanos();
1060            Self(std::env::temp_dir().join(format!("qcode-relay-{name}-{stamp}")))
1061        }
1062    }
1063
1064    impl Drop for Scratch {
1065        fn drop(&mut self) {
1066            let _ = std::fs::remove_dir_all(&self.0);
1067        }
1068    }
1069
1070    fn tag(name: &str) -> Tag {
1071        Tag::parse(name).expect("a tag")
1072    }
1073
1074    fn ollama_entry(key: Option<&str>) -> ProviderEntry {
1075        let mut entry = ProviderEntry::new(tag("ev"), ProviderKind::Ollama, "http://192.168.122.1:11434");
1076        entry.key = key.map(|key| Key::new(key).expect("a key"));
1077        entry
1078    }
1079
1080    /// One route answered for `token`, nothing for anything else: the entry and the single model
1081    /// it is asked for, which is what a profile that names a model rather than a lineup resolves
1082    /// to today.
1083    fn resolve_one(
1084        token: &'static str,
1085        entry: ProviderEntry,
1086        model: &'static str,
1087    ) -> impl Fn(&str) -> Option<Route> + Send + Sync {
1088        let models = [model.to_owned()];
1089        move |asked| (asked == token).then(|| Route { entry: entry.clone(), models: models.to_vec(), lineup: None })
1090    }
1091
1092    /// One route answered for `token` with `models` to try in order, under `lineup`'s name.
1093    fn resolve_steps(
1094        token: &'static str,
1095        entry: ProviderEntry,
1096        models: &'static [&'static str],
1097        lineup: Option<&'static str>,
1098    ) -> impl Fn(&str) -> Option<Route> + Send + Sync {
1099        let models: Vec<String> = models.iter().map(|model| (*model).to_owned()).collect();
1100        let lineup = lineup.map(str::to_owned);
1101        move |asked| {
1102            (asked == token).then(|| Route { entry: entry.clone(), models: models.clone(), lineup: lineup.clone() })
1103        }
1104    }
1105
1106    /// The events a listener was told about, in the order it was told them.
1107    fn taken(reported: &Arc<Mutex<Vec<Event>>>) -> Vec<Event> {
1108        reported.lock().expect("the events are not poisoned").clone()
1109    }
1110
1111    /// What reached the stand-in upstream, for a test to look at exactly what would have left
1112    /// the machine.
1113    struct SeenAsk {
1114        url: String,
1115        headers: Vec<(String, String)>,
1116        secret: Option<Secret>,
1117        body: Vec<u8>,
1118    }
1119
1120    /// An upstream that answers `body` for every request, and hands every [`UpstreamAsk`] it was
1121    /// given to `seen` as a [`SeenAsk`].
1122    fn canned(status: u16, body: &'static str, seen: Sender<SeenAsk>) -> Upstream {
1123        Upstream::new(move |mut ask| {
1124            let mut sent = Vec::new();
1125            let _ = ask.body.read_to_end(&mut sent);
1126            let _ = seen.send(SeenAsk {
1127                url: ask.url.clone(),
1128                headers: ask.headers.clone(),
1129                secret: ask.secret.clone(),
1130                body: sent,
1131            });
1132            Ok(UpstreamAnswer {
1133                status,
1134                headers: vec![("content-type".to_owned(), "application/json".to_owned())],
1135                body: Box::new(std::io::Cursor::new(body.as_bytes().to_vec())),
1136            })
1137        })
1138    }
1139
1140    /// What the stand-in says to one model: the status and the body of its answer.
1141    type Reply = (u16, &'static str);
1142
1143    /// An upstream that answers each model of `answers` with the [`Reply`] written beside it, and
1144    /// hands every [`UpstreamAsk`] to `seen` before it answers. A model asked that is not named
1145    /// here is answered `400` naming itself, so a test that expected another model to be asked
1146    /// fails on what the harness was told rather than on silence.
1147    fn by_model(answers: Vec<(&'static str, Reply)>, seen: Sender<SeenAsk>) -> Upstream {
1148        Upstream::new(move |mut ask| {
1149            let mut sent = Vec::new();
1150            let _ = ask.body.read_to_end(&mut sent);
1151            let asked: Option<String> =
1152                serde_json::from_slice::<Value>(&sent).ok().and_then(|body| body["model"].as_str().map(str::to_owned));
1153            let _ = seen.send(SeenAsk {
1154                url: ask.url.clone(),
1155                headers: ask.headers.clone(),
1156                secret: ask.secret.clone(),
1157                body: sent,
1158            });
1159            let reply = asked
1160                .as_deref()
1161                .and_then(|model| answers.iter().find(|(known, _)| *known == model).map(|(_, reply)| *reply))
1162                .unwrap_or((400, r#"{"error":"the stand-in has no answer for that model"}"#));
1163            let (status, body) = reply;
1164            Ok(UpstreamAnswer {
1165                status,
1166                headers: vec![("content-type".to_owned(), "application/json".to_owned())],
1167                body: Box::new(std::io::Cursor::new(body.as_bytes().to_vec())),
1168            })
1169        })
1170    }
1171
1172    /// Every model the stand-in was asked for, in the order they were asked, read out of the
1173    /// bodies: what the provider would have seen, whatever the harness's own body said.
1174    fn asked_models(seen: &Receiver<SeenAsk>, count: usize) -> Vec<String> {
1175        (0..count)
1176            .map(|_| {
1177                let ask = seen.try_recv().expect("the request reached the stand-in");
1178                serde_json::from_slice::<Value>(&ask.body).expect("a JSON body")["model"]
1179                    .as_str()
1180                    .expect("a model in the body")
1181                    .to_owned()
1182            })
1183            .collect()
1184    }
1185
1186    /// An upstream that never answers, for a test proving it was never called.
1187    fn unreachable_if_called(count: Arc<Mutex<usize>>) -> Upstream {
1188        Upstream::new(move |_| {
1189            *count.lock().expect("the lock") += 1;
1190            Ok(UpstreamAnswer { status: 200, headers: Vec::new(), body: Box::new(std::io::empty()) })
1191        })
1192    }
1193
1194    fn connect(socket: &Path) -> UnixStream {
1195        if socket.as_os_str().len() <= 107 {
1196            return UnixStream::connect(socket).expect("the socket answers");
1197        }
1198        let handle = std::fs::File::open(socket.parent().expect("a folder")).expect("the folder opens");
1199        let short = format!("/proc/self/fd/{}/{SOCKET_NAME}", std::os::fd::AsRawFd::as_raw_fd(&handle));
1200        UnixStream::connect(short).expect("the socket answers")
1201    }
1202
1203    /// Sends one request the way the relay script would: the head line, then the body, then the
1204    /// write side is closed. Answers the head line QCode sent back and the whole of its body.
1205    fn ask(socket: &Path, token: &str, method: &str, path: &str, body: &[u8]) -> (Value, Vec<u8>) {
1206        ask_with_headers(socket, token, method, path, &[], body)
1207    }
1208
1209    /// [`ask`], writing a body of `size` bytes that is never held whole in this test: the point
1210    /// of the bound is a body too large to keep, and a test that kept one would only move the
1211    /// cost. The byte is written over and over, so the body is the size asked for and nothing of
1212    /// it is a message.
1213    fn ask_huge(socket: &Path, token: &str, method: &str, path: &str, size: usize) -> (Value, Vec<u8>) {
1214        let mut stream = connect(socket);
1215        let head = json!({ "token": token, "method": method, "path": path, "headers": {} }).to_string();
1216        stream.write_all(format!("{head}\n").as_bytes()).expect("the head line is written");
1217        let written = io::copy(&mut std::io::repeat(b'x').take(bound(size)), &mut stream).expect("the body is written");
1218        assert_eq!(written, bound(size), "the whole body is on its way");
1219        stream.shutdown(std::net::Shutdown::Write).expect("the write side closes");
1220        let mut reading = BufReader::new(stream);
1221        let mut line = String::new();
1222        reading.read_line(&mut line).expect("an answer head comes");
1223        let head: Value = serde_json::from_str(line.trim_end()).expect("the head is JSON");
1224        let mut answer = Vec::new();
1225        reading.read_to_end(&mut answer).expect("the answer body is read");
1226        (head, answer)
1227    }
1228
1229    /// [`ask`], carrying `headers` the way a harness's own request would.
1230    fn ask_with_headers(
1231        socket: &Path,
1232        token: &str,
1233        method: &str,
1234        path: &str,
1235        headers: &[(&str, &str)],
1236        body: &[u8],
1237    ) -> (Value, Vec<u8>) {
1238        let mut stream = connect(socket);
1239        let headers: Map<String, Value> =
1240            headers.iter().map(|(name, value)| ((*name).to_owned(), json!(value))).collect();
1241        let head = json!({ "token": token, "method": method, "path": path, "headers": headers }).to_string();
1242        stream.write_all(format!("{head}\n").as_bytes()).expect("the head line is written");
1243        stream.write_all(body).expect("the body is written");
1244        stream.shutdown(std::net::Shutdown::Write).expect("the write side closes");
1245        let mut reading = BufReader::new(stream);
1246        let mut line = String::new();
1247        reading.read_line(&mut line).expect("an answer head comes");
1248        let head: Value = serde_json::from_str(line.trim_end()).expect("the head is JSON");
1249        let mut answer = Vec::new();
1250        reading.read_to_end(&mut answer).expect("the answer body is read");
1251        (head, answer)
1252    }
1253
1254    /// What a harness's body says, as the person reading this test would write it.
1255    const A_HARNESS_ASKING: &[u8] = br#"{"model":"claude-haiku-4-5","max_tokens":1024,"stream":true,"messages":[{"role":"user","content":"selam"}]}"#;
1256
1257    /// The same body with the model QCode pins to it, and nothing else changed.
1258    const PINNED_TO_THE_PROFILE: &str = r#"{"model":"qwen/qwen3-coder:free","max_tokens":1024,"stream":true,"messages":[{"role":"user","content":"selam"}]}"#;
1259
1260    #[test]
1261    fn the_provider_is_asked_for_the_model_the_profile_chose_and_not_the_one_the_harness_named() {
1262        let scratch = Scratch::new("pin");
1263        let (seen_tx, seen_rx) = mpsc::channel();
1264        let listener = Listener::open(
1265            &scratch.0,
1266            resolve_one("tok", ollama_entry(None), CHOSEN),
1267            canned(200, r#"{"ok":true}"#, seen_tx),
1268            |_| {},
1269        )
1270        .expect("the socket opens");
1271
1272        // Claude Code asks for its own small model by its own name for background work, which is
1273        // a model the person never chose and, on a provider that bills for it, a paid one.
1274        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1275        assert_eq!(head["status"], 200);
1276        let seen = seen_rx.recv().expect("the request reached the stand-in");
1277        let sent: Value = serde_json::from_slice(&seen.body).expect("the body that left is JSON");
1278        assert_eq!(sent, serde_json::from_str::<Value>(PINNED_TO_THE_PROFILE).expect("the expected body is JSON"));
1279    }
1280
1281    #[test]
1282    fn a_body_that_is_not_json_goes_on_byte_for_byte() {
1283        let scratch = Scratch::new("notjson");
1284        let (seen_tx, seen_rx) = mpsc::channel();
1285        let listener = Listener::open(
1286            &scratch.0,
1287            resolve_one("tok", ollama_entry(None), CHOSEN),
1288            canned(200, "{}", seen_tx),
1289            |_| {},
1290        )
1291        .expect("the socket opens");
1292
1293        let body = b"model=qwen3-coder:free&stream=1";
1294        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/messages", body);
1295        assert_eq!(head["status"], 200);
1296        let seen = seen_rx.recv().expect("the request reached the stand-in");
1297        assert_eq!(seen.body, body, "a body QCode cannot read a model out of is not one it touches");
1298    }
1299
1300    #[test]
1301    fn a_body_over_the_limit_is_refused_before_any_provider_is_asked() {
1302        let scratch = Scratch::new("huge");
1303        let calls = Arc::new(Mutex::new(0));
1304        let listener = Listener::open(
1305            &scratch.0,
1306            resolve_one("tok", ollama_entry(None), CHOSEN),
1307            unreachable_if_called(Arc::clone(&calls)),
1308            |_| {},
1309        )
1310        .expect("the socket opens");
1311
1312        let (head, body) = ask_huge(listener.socket(), "tok", "POST", "/v1/messages", MOST_BODY + 1);
1313        assert_eq!(head["status"], 413);
1314        assert!(!String::from_utf8_lossy(&body).contains('x'), "nothing of the refused body is sent back");
1315        assert_eq!(*calls.lock().expect("the lock"), 0, "a body too large never reaches a provider");
1316
1317        // The same relay still answers a body within the limit, so the bound is a bound and not a
1318        // refusal of the message shape.
1319        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1320        assert_eq!(head["status"], 200);
1321        assert_eq!(*calls.lock().expect("the lock"), 1);
1322    }
1323
1324    #[test]
1325    fn a_good_token_is_forwarded_with_the_key_added_and_the_answer_comes_back() {
1326        let scratch = Scratch::new("round");
1327        let entry = ollama_entry(None);
1328        let (seen_tx, seen_rx) = mpsc::channel();
1329        let upstream = canned(200, r#"{"ok":true}"#, seen_tx);
1330        let listener =
1331            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), upstream, |_| {}).expect("the socket opens");
1332
1333        let headers = [("content-type", "application/json"), ("authorization", "Bearer harness-own-key")];
1334        let (head, body) = ask_with_headers(listener.socket(), "tok", "POST", "/v1/messages", &headers, br#"{"hi":1}"#);
1335        assert_eq!(head["status"], 200);
1336        assert_eq!(body, br#"{"ok":true}"#);
1337        let seen = seen_rx.recv().expect("the request reached the stand-in");
1338        assert_eq!(seen.url, "http://192.168.122.1:11434/v1/messages");
1339        assert_eq!(seen.body, br#"{"hi":1}"#, "the harness's body reaches the provider whole");
1340        assert!(seen.headers.iter().any(|(name, value)| name == "content-type" && value == "application/json"));
1341        assert!(
1342            seen.headers.iter().all(|(name, _)| name != "authorization"),
1343            "the harness's own authorization header never reaches the provider"
1344        );
1345        assert!(seen.secret.is_none(), "an ollama entry with no key adds none");
1346    }
1347
1348    #[test]
1349    fn the_provider_gets_a_key_the_container_never_sees() {
1350        let scratch = Scratch::new("key");
1351        let entry = ollama_entry(Some(MADE_UP_KEY));
1352        let (seen_tx, seen_rx) = mpsc::channel();
1353        let upstream = canned(200, r#"{"ok":true}"#, seen_tx);
1354        let listener =
1355            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), upstream, |_| {}).expect("the socket opens");
1356
1357        let (head, body) = ask(listener.socket(), "tok", "GET", "/v1/models", b"");
1358        assert_eq!(head["status"], 200);
1359        assert!(!format!("{head}").contains(MADE_UP_KEY));
1360        assert!(!String::from_utf8_lossy(&body).contains(MADE_UP_KEY), "the container's own answer holds no key");
1361
1362        let seen = seen_rx.recv().expect("the request reached the stand-in");
1363        let secret = seen.secret.expect("the key was added for the provider");
1364        assert_eq!(secret.header, "Authorization");
1365        assert_eq!(secret.key.expose(), MADE_UP_KEY, "the real key reached the provider's own request");
1366    }
1367
1368    #[test]
1369    fn a_harness_asking_openrouter_the_way_it_asks_anthropic_reaches_its_api() {
1370        let scratch = Scratch::new("openrouter");
1371        let mut entry = ProviderEntry::new(tag("yol"), ProviderKind::OpenRouter, "https://openrouter.ai");
1372        entry.key = Some(Key::new(MADE_UP_KEY).expect("a key"));
1373        let (seen_tx, seen_rx) = mpsc::channel();
1374        let listener =
1375            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), canned(200, "{}", seen_tx), |_| {})
1376                .expect("the socket opens");
1377        // What Claude Code really sends: its own path, with the query string it adds.
1378        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/messages?beta=true", b"{}");
1379        assert_eq!(head["status"], 200);
1380        let seen = seen_rx.recv().expect("the request reached the stand-in");
1381        assert_eq!(seen.url, "https://openrouter.ai/api/v1/messages?beta=true", "the API, not the website");
1382    }
1383
1384    #[test]
1385    fn no_token_or_a_wrong_one_is_refused_and_nothing_is_forwarded() {
1386        let scratch = Scratch::new("badtoken");
1387        let entry = ollama_entry(None);
1388        let calls = Arc::new(Mutex::new(0));
1389        let listener = Listener::open(
1390            &scratch.0,
1391            resolve_one("tok", entry, CHOSEN),
1392            unreachable_if_called(Arc::clone(&calls)),
1393            |_| {},
1394        )
1395        .expect("the socket opens");
1396
1397        let (head, _) = ask(listener.socket(), "wrong", "GET", "/v1/models", b"");
1398        assert_eq!(head["status"], 401);
1399        assert_eq!(*calls.lock().expect("the lock"), 0, "a bad token never reaches the provider");
1400    }
1401
1402    #[test]
1403    fn a_path_that_is_not_allowed_is_refused() {
1404        let scratch = Scratch::new("badpath");
1405        let entry = ollama_entry(None);
1406        let calls = Arc::new(Mutex::new(0));
1407        let listener = Listener::open(
1408            &scratch.0,
1409            resolve_one("tok", entry, CHOSEN),
1410            unreachable_if_called(Arc::clone(&calls)),
1411            |_| {},
1412        )
1413        .expect("the socket opens");
1414
1415        let (head, _) = ask(listener.socket(), "tok", "GET", "/v1/admin", b"");
1416        assert_eq!(head["status"], 404);
1417        let (head, _) = ask(listener.socket(), "tok", "DELETE", "/v1/messages", b"");
1418        assert_eq!(head["status"], 405);
1419        assert_eq!(*calls.lock().expect("the lock"), 0, "a path outside the allowed set never reaches the provider");
1420    }
1421
1422    /// A body that hands out what it is given one piece at a time, blocking until the next piece
1423    /// or the end is sent. A relay that buffered a streamed answer before writing any of it would
1424    /// make the second half of this test indistinguishable from the first; only reading in
1425    /// pieces, ahead of the last piece being sent, proves it does not.
1426    struct Trickle(Receiver<Option<Vec<u8>>>, Vec<u8>);
1427
1428    impl Read for Trickle {
1429        fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
1430            if self.1.is_empty() {
1431                match self.0.recv() {
1432                    Ok(Some(piece)) => self.1 = piece,
1433                    _ => return Ok(0),
1434                }
1435            }
1436            let n = buf.len().min(self.1.len());
1437            buf[..n].copy_from_slice(&self.1[..n]);
1438            self.1.drain(..n);
1439            Ok(n)
1440        }
1441    }
1442
1443    /// A body that gives a piece of itself and then breaks, as a stream cut in the middle of an
1444    /// answer does. The error is not a refusal: nothing has refused, the answer simply stopped
1445    /// coming, which is what must not be mistaken for a reason to ask another model.
1446    struct Broken(Vec<u8>);
1447
1448    impl Read for Broken {
1449        fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
1450            if self.0.is_empty() {
1451                return Err(io::Error::other("the stream broke"));
1452            }
1453            let n = buf.len().min(self.0.len()).min(16);
1454            buf[..n].copy_from_slice(&self.0[..n]);
1455            self.0.drain(..n);
1456            Ok(n)
1457        }
1458    }
1459
1460    /// The models of a route in the order a test asks them, and the first that is expected to
1461    /// refuse with `status` while the second answers `200` with `body`. The two are the ones the
1462    /// fall-back tests use, so a test that reads them cannot pass with the order reversed.
1463    const A_LINEUP: &[&str] = &["a/one", "b/two"];
1464
1465    /// What the first model of [`A_LINEUP`] says when the provider will not have it, quoted the way
1466    /// a provider quotes a model it does not have.
1467    const NO_SUCH_MODEL: &str = r#"{"error":{"message":"model a/one not found"}}"#;
1468
1469    /// What the second model answers when it is asked.
1470    const THE_ANSWER: &str = r#"{"id":"chatcmpl-kyaz-1","choices":[{"text":"done"}]}"#;
1471
1472    /// A lineup of three, for the tests that need more than one step to have failed before the last
1473    /// one does: the first answer kept is only the first one when something answered before it.
1474    const THREE_MODELS: &[&str] = &["a/one", "b/two", "c/three"];
1475
1476    /// A stand-in that refuses `A_LINEUP`'s first model and answers its second, and records every
1477    /// request that reaches it.
1478    fn one_busy_one_willing(seen: Sender<SeenAsk>) -> Upstream {
1479        by_model(vec![(A_LINEUP[0], (429, r#"{"error":"busy"}"#)), (A_LINEUP[1], (200, THE_ANSWER))], seen)
1480    }
1481
1482    /// The order a route of two models is asked in at `now`, with the first one cooling down. The
1483    /// pure function behind the listener's own ordering, tested on a clock rather than on a wait.
1484    fn order_of(models: &[&str], cooling_until: Instant, now: Instant) -> Vec<String> {
1485        let models: Vec<String> = models.iter().map(|model| (*model).to_owned()).collect();
1486        let seen: Cooling =
1487            [(("ev".to_owned(), Some("bilim".to_owned()), "a/one".to_owned()), cooling_until)].into_iter().collect();
1488        ordered(&models, &seen, &|model| ("ev".to_owned(), Some("bilim".to_owned()), model.to_owned()), now)
1489    }
1490
1491    #[test]
1492    fn a_busy_first_model_hands_its_turn_to_the_next_and_the_harness_gets_that_answer() {
1493        let scratch = Scratch::new("fallback");
1494        let (seen_tx, seen_rx) = mpsc::channel();
1495        let reported: Arc<Mutex<Vec<Event>>> = Arc::new(Mutex::new(Vec::new()));
1496        let told = Arc::clone(&reported);
1497        let listener = Listener::open(
1498            &scratch.0,
1499            resolve_steps("tok", ollama_entry(None), A_LINEUP, Some("bilim")),
1500            one_busy_one_willing(seen_tx),
1501            move |event| told.lock().expect("the events are not poisoned").push(event),
1502        )
1503        .expect("the socket opens");
1504
1505        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1506        assert_eq!(head["status"], 200);
1507        assert_eq!(body, THE_ANSWER.as_bytes(), "the harness gets the answer of the model that had one");
1508        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP, "the lineup is asked in the order it was written");
1509        assert_eq!(
1510            taken(&reported),
1511            [
1512                Event::FellBack {
1513                    token: "tok".to_owned(),
1514                    tag: "ev".to_owned(),
1515                    lineup: Some("bilim".to_owned()),
1516                    from: "a/one".to_owned(),
1517                    to: "b/two".to_owned(),
1518                    status: Some(429),
1519                },
1520                Event::Forwarded { tag: "ev".to_owned(), path: "/v1/messages".to_owned(), status: 200 },
1521            ]
1522        );
1523    }
1524
1525    #[test]
1526    fn a_bad_request_about_the_model_it_names_hands_its_turn_to_the_next_and_one_about_something_else_does_not() {
1527        let scratch = Scratch::new("badmodel");
1528        for (named, expected) in [(true, 200), (false, 400)] {
1529            let (seen_tx, seen_rx) = mpsc::channel();
1530            let about_the_request = r#"{"error":{"message":"too many tokens"}}"#;
1531            let first = if named { NO_SUCH_MODEL } else { about_the_request };
1532            let answers = vec![(A_LINEUP[0], (400, first)), (A_LINEUP[1], (200, THE_ANSWER))];
1533            let listener = Listener::open(
1534                &scratch.0,
1535                resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1536                by_model(answers, seen_tx),
1537                |_| {},
1538            )
1539            .expect("the socket opens");
1540
1541            let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1542            assert_eq!(head["status"], expected, "whether the body names the model asked decides it");
1543            if named {
1544                assert_eq!(body, THE_ANSWER.as_bytes());
1545            } else {
1546                assert_eq!(body, about_the_request.as_bytes(), "the request's own complaint reaches the harness whole");
1547                assert_eq!(asked_models(&seen_rx, 1), ["a/one".to_owned()], "no other model was asked");
1548            }
1549        }
1550    }
1551
1552    #[test]
1553    fn a_refusal_of_the_key_is_never_a_reason_to_ask_another_model() {
1554        let scratch = Scratch::new("unauthorized");
1555        let (seen_tx, seen_rx) = mpsc::channel();
1556        let listener = Listener::open(
1557            &scratch.0,
1558            resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1559            by_model(vec![(A_LINEUP[0], (401, r#"{"error":"no key"}"#)), (A_LINEUP[1], (200, THE_ANSWER))], seen_tx),
1560            |_| {},
1561        )
1562        .expect("the socket opens");
1563
1564        // One key, one provider, one answer: the next model of the same provider would be refused
1565        // the same way, and the person is better served by what the provider actually said.
1566        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1567        assert_eq!(head["status"], 401);
1568        assert_eq!(body, br#"{"error":"no key"}"#);
1569        assert_eq!(asked_models(&seen_rx, 1), ["a/one".to_owned()], "no other model was asked");
1570    }
1571
1572    #[test]
1573    fn a_model_that_cannot_be_reached_hands_its_turn_to_the_next() {
1574        let scratch = Scratch::new("unreachable");
1575        let (seen_tx, seen_rx) = mpsc::channel();
1576        // An address that is not there for the first model and is for the second: what a provider
1577        // that has come back looks like to a request that is still on its way.
1578        let upstream = Upstream::new(move |mut ask| {
1579            let mut sent = Vec::new();
1580            let _ = ask.body.read_to_end(&mut sent);
1581            let asked: Option<String> =
1582                serde_json::from_slice::<Value>(&sent).ok().and_then(|body| body["model"].as_str().map(str::to_owned));
1583            let _ = seen_tx.send(SeenAsk {
1584                url: ask.url.clone(),
1585                headers: ask.headers.clone(),
1586                secret: ask.secret.clone(),
1587                body: sent,
1588            });
1589            if asked.as_deref() == Some(A_LINEUP[0]) {
1590                return Err(UpstreamError { url: ask.url.clone(), reason: "connection refused".to_owned() });
1591            }
1592            Ok(UpstreamAnswer {
1593                status: 200,
1594                headers: vec![("content-type".to_owned(), "application/json".to_owned())],
1595                body: Box::new(std::io::Cursor::new(THE_ANSWER.as_bytes().to_vec())),
1596            })
1597        });
1598        let reported: Arc<Mutex<Vec<Event>>> = Arc::new(Mutex::new(Vec::new()));
1599        let told = Arc::clone(&reported);
1600        let listener = Listener::open(
1601            &scratch.0,
1602            resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1603            upstream,
1604            move |event| told.lock().expect("the events are not poisoned").push(event),
1605        )
1606        .expect("the socket opens");
1607
1608        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1609        assert_eq!(head["status"], 200, "the second model answers");
1610        assert_eq!(body, THE_ANSWER.as_bytes());
1611        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP, "both were asked");
1612        assert_eq!(
1613            taken(&reported).first(),
1614            Some(&Event::FellBack {
1615                token: "tok".to_owned(),
1616                tag: "ev".to_owned(),
1617                lineup: None,
1618                from: "a/one".to_owned(),
1619                to: "b/two".to_owned(),
1620                status: None,
1621            }),
1622            "an answer that never came is a fall-back with no status to name"
1623        );
1624    }
1625
1626    #[test]
1627    fn the_last_step_s_own_refusal_reaches_the_harness_as_it_is() {
1628        let scratch = Scratch::new("laststep");
1629        let (seen_tx, seen_rx) = mpsc::channel();
1630        let said = r#"{"error":"the model server is down"}"#;
1631        let listener = Listener::open(
1632            &scratch.0,
1633            resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1634            by_model(vec![(A_LINEUP[0], (429, r#"{"error":"busy"}"#)), (A_LINEUP[1], (503, said))], seen_tx),
1635            |_| {},
1636        )
1637        .expect("the socket opens");
1638
1639        // The harness shows its own tab's error rather than dying, so the tab stays and the
1640        // person is still there when the provider comes back.
1641        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1642        assert_eq!(head["status"], 503);
1643        assert_eq!(body, said.as_bytes(), "the last step's body is not QCode's to change");
1644        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP, "both were asked and there was nowhere else to go");
1645    }
1646
1647    #[test]
1648    fn a_busy_first_model_is_the_answer_the_harness_gets_when_the_last_one_names_a_model_that_is_not_there() {
1649        let scratch = Scratch::new("passing");
1650        let (seen_tx, seen_rx) = mpsc::channel();
1651        let busy = r#"{"error":"busy, retry shortly"}"#;
1652        let missing = r#"{"error":{"message":"model b/two not found"}}"#;
1653        let reported: Arc<Mutex<Vec<Event>>> = Arc::new(Mutex::new(Vec::new()));
1654        let told = Arc::clone(&reported);
1655        let listener = Listener::open(
1656            &scratch.0,
1657            resolve_steps("tok", ollama_entry(None), A_LINEUP, Some("bilim")),
1658            by_model(vec![(A_LINEUP[0], (429, busy)), (A_LINEUP[1], (404, missing))], seen_tx),
1659            move |event| told.lock().expect("the events are not poisoned").push(event),
1660        )
1661        .expect("the socket opens");
1662
1663        // A person told to wait does wait; a person told a model does not exist goes and changes a
1664        // lineup that works. This is what the tab said before the lineup was ever written down.
1665        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1666        assert_eq!(head["status"], 429);
1667        assert_eq!(body, busy.as_bytes(), "the provider's own words, whole and unedited");
1668        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP, "both models were asked");
1669        // The last step is not cooled and is not said to be passed over: there was nowhere to fall
1670        // back to, and what the harness was shown is not what the last step answered.
1671        assert_eq!(
1672            taken(&reported).last(),
1673            Some(&Event::Forwarded { tag: "ev".to_owned(), path: "/v1/messages".to_owned(), status: 429 }),
1674            "what the harness got is the one thing said about this request: {reported:?}"
1675        );
1676    }
1677
1678    #[test]
1679    fn of_several_passing_answers_the_first_one_is_the_answer_the_harness_gets() {
1680        let scratch = Scratch::new("firstpassing");
1681        let (seen_tx, seen_rx) = mpsc::channel();
1682        let first = r#"{"error":"the model server is down"}"#;
1683        let second = r#"{"error":"busy, retry shortly"}"#;
1684        let listener = Listener::open(
1685            &scratch.0,
1686            resolve_steps("tok", ollama_entry(None), THREE_MODELS, Some("bilim")),
1687            by_model(
1688                vec![
1689                    (THREE_MODELS[0], (503, first)),
1690                    (THREE_MODELS[1], (429, second)),
1691                    (THREE_MODELS[2], (404, r#"{"error":{"message":"model c/three not found"}}"#)),
1692                ],
1693                seen_tx,
1694            ),
1695            |_| {},
1696        )
1697        .expect("the socket opens");
1698
1699        // The first model of the lineup is the one the person wrote first, so its answer is the one
1700        // worth showing even though a later model was passing as well.
1701        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1702        assert_eq!(head["status"], 503);
1703        assert_eq!(body, first.as_bytes());
1704        assert_eq!(asked_models(&seen_rx, 3), THREE_MODELS, "all three were asked");
1705    }
1706
1707    #[test]
1708    fn a_provider_that_could_not_be_reached_is_the_answer_the_harness_gets_when_the_last_one_is_not_there() {
1709        let scratch = Scratch::new("unreachable-passing");
1710        let (seen_tx, seen_rx) = mpsc::channel();
1711        // The provider is not there at all for the first model, and has never heard of the second:
1712        // what this request ran into is a provider that is down, not a lineup that is wrong.
1713        let upstream = Upstream::new(move |mut ask| {
1714            let mut sent = Vec::new();
1715            let _ = ask.body.read_to_end(&mut sent);
1716            let asked: Option<String> =
1717                serde_json::from_slice::<Value>(&sent).ok().and_then(|body| body["model"].as_str().map(str::to_owned));
1718            let _ = seen_tx.send(SeenAsk {
1719                url: ask.url.clone(),
1720                headers: ask.headers.clone(),
1721                secret: ask.secret.clone(),
1722                body: sent,
1723            });
1724            if asked.as_deref() == Some(A_LINEUP[0]) {
1725                return Err(UpstreamError { url: ask.url.clone(), reason: "connection refused".to_owned() });
1726            }
1727            Ok(UpstreamAnswer {
1728                status: 404,
1729                headers: vec![("content-type".to_owned(), "application/json".to_owned())],
1730                body: Box::new(std::io::Cursor::new(
1731                    r#"{"error":{"message":"model b/two not found"}}"#.as_bytes().to_vec(),
1732                )),
1733            })
1734        });
1735        let listener =
1736            Listener::open(&scratch.0, resolve_steps("tok", ollama_entry(None), A_LINEUP, None), upstream, |_| {})
1737                .expect("the socket opens");
1738
1739        // The relay's own refusal, exactly the one a last step that cannot be reached writes: what
1740        // the harness is shown is a provider to come back to rather than a model to go and find.
1741        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1742        assert_eq!(head["status"], 502);
1743        assert_eq!(body, refusal_body("the provider could not be reached"));
1744        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP, "both were asked");
1745    }
1746
1747    #[test]
1748    fn a_busy_first_model_does_not_take_the_place_of_a_complaint_about_the_request() {
1749        let scratch = Scratch::new("abouttherequest");
1750        let (seen_tx, seen_rx) = mpsc::channel();
1751        let about_the_request = r#"{"error":{"message":"too many tokens"}}"#;
1752        let listener = Listener::open(
1753            &scratch.0,
1754            resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1755            by_model(
1756                vec![
1757                    (A_LINEUP[0], (429, r#"{"error":"busy, retry shortly"}"#)),
1758                    (A_LINEUP[1], (400, about_the_request)),
1759                ],
1760                seen_tx,
1761            ),
1762            |_| {},
1763        )
1764        .expect("the socket opens");
1765
1766        // Nothing about the lineup is wrong: every model of it would refuse a request this long, so
1767        // being told to wait would send the person to wait for an answer that never comes.
1768        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1769        assert_eq!(head["status"], 400);
1770        assert_eq!(body, about_the_request.as_bytes(), "the request's own complaint reaches the harness whole");
1771        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP);
1772    }
1773
1774    #[test]
1775    fn a_route_of_one_model_that_does_not_exist_is_answered_as_it_is() {
1776        let scratch = Scratch::new("nostep");
1777        let (seen_tx, seen_rx) = mpsc::channel();
1778        let listener = Listener::open(
1779            &scratch.0,
1780            resolve_one("tok", ollama_entry(None), A_LINEUP[0]),
1781            by_model(vec![(A_LINEUP[0], (404, NO_SUCH_MODEL))], seen_tx),
1782            |_| {},
1783        )
1784        .expect("the socket opens");
1785
1786        // There was no earlier step to pass for now, so there is nothing to show in place of the
1787        // model's own complaint — which is the one thing a person can act on here.
1788        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1789        assert_eq!(head["status"], 404);
1790        assert_eq!(body, NO_SUCH_MODEL.as_bytes());
1791        assert_eq!(asked_models(&seen_rx, 1), [A_LINEUP[0].to_owned()]);
1792    }
1793
1794    #[test]
1795    fn a_replayed_answer_carries_neither_the_length_nor_the_framing_of_the_one_it_was_read_from() {
1796        let scratch = Scratch::new("replayedhead");
1797        let (seen_tx, seen_rx) = mpsc::channel();
1798        let busy = r#"{"error":"busy, retry shortly"}"#;
1799        // What a provider really sends: the length of the body and the framing it is written in.
1800        // The body was read before the step was given up on, so neither of them describes what the
1801        // harness would be sent, and a head the script writes as it stands would cut the answer.
1802        let upstream = Upstream::new(move |mut ask| {
1803            let mut sent = Vec::new();
1804            let _ = ask.body.read_to_end(&mut sent);
1805            let asked: Option<String> =
1806                serde_json::from_slice::<Value>(&sent).ok().and_then(|body| body["model"].as_str().map(str::to_owned));
1807            let _ = seen_tx.send(SeenAsk {
1808                url: ask.url.clone(),
1809                headers: ask.headers.clone(),
1810                secret: ask.secret.clone(),
1811                body: sent,
1812            });
1813            if asked.as_deref() == Some(A_LINEUP[0]) {
1814                return Ok(UpstreamAnswer {
1815                    status: 429,
1816                    headers: vec![
1817                        ("content-type".to_owned(), "application/json".to_owned()),
1818                        ("content-length".to_owned(), busy.len().to_string()),
1819                        ("transfer-encoding".to_owned(), "chunked".to_owned()),
1820                    ],
1821                    body: Box::new(std::io::Cursor::new(busy.as_bytes().to_vec())),
1822                });
1823            }
1824            Ok(UpstreamAnswer {
1825                status: 404,
1826                headers: vec![("content-type".to_owned(), "application/json".to_owned())],
1827                body: Box::new(std::io::Cursor::new(
1828                    r#"{"error":{"message":"model b/two not found"}}"#.as_bytes().to_vec(),
1829                )),
1830            })
1831        });
1832        let listener =
1833            Listener::open(&scratch.0, resolve_steps("tok", ollama_entry(None), A_LINEUP, None), upstream, |_| {})
1834                .expect("the socket opens");
1835
1836        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1837        assert_eq!(head["status"], 429);
1838        assert_eq!(body, busy.as_bytes(), "the whole of what was read is what the harness gets");
1839        for gone in ["content-length", "transfer-encoding"] {
1840            assert!(
1841                head["headers"].get(gone).is_none(),
1842                "{gone} describes a body that was read and may have been cut: {:?}",
1843                head["headers"]
1844            );
1845        }
1846        assert_eq!(head["headers"]["content-type"], "application/json", "the rest of the provider's head is its own");
1847        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP);
1848    }
1849
1850    #[test]
1851    fn an_answer_that_breaks_half_way_is_not_asked_a_second_time() {
1852        let scratch = Scratch::new("broken");
1853        let asked_count = Arc::new(Mutex::new(0));
1854        let counted = Arc::clone(&asked_count);
1855        let upstream = Upstream::new(move |_| {
1856            *counted.lock().expect("the lock") += 1;
1857            Ok(UpstreamAnswer {
1858                status: 200,
1859                headers: vec![("content-type".to_owned(), "text/event-stream".to_owned())],
1860                body: Box::new(Broken(b"data: one\n\ndata: tw".to_vec())),
1861            })
1862        });
1863        let listener =
1864            Listener::open(&scratch.0, resolve_steps("tok", ollama_entry(None), A_LINEUP, None), upstream, |_| {})
1865                .expect("the socket opens");
1866
1867        // What the harness got is what the model managed to say before the stream broke, and the
1868        // second model is never asked: a turn already begun cannot be begun again.
1869        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1870        assert_eq!(head["status"], 200);
1871        assert!(!body.is_empty(), "what came before the break is the harness's: {body:?}");
1872        assert_eq!(*asked_count.lock().expect("the lock"), 1, "one model was asked and no other");
1873    }
1874
1875    #[test]
1876    fn a_model_that_refused_is_not_asked_again_until_its_minute_is_up() {
1877        let scratch = Scratch::new("cooling");
1878        let (seen_tx, seen_rx) = mpsc::channel();
1879        let listener = Listener::open(
1880            &scratch.0,
1881            resolve_steps("tok", ollama_entry(None), A_LINEUP, None),
1882            one_busy_one_willing(seen_tx),
1883            |_| {},
1884        )
1885        .expect("the socket opens");
1886
1887        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1888        assert_eq!(head["status"], 200);
1889        assert_eq!(asked_models(&seen_rx, 2), A_LINEUP);
1890
1891        // The next request goes straight to the model that answered: a busy free model that is
1892        // asked first every time costs a round trip per turn for nothing.
1893        let (head, body) = ask(listener.socket(), "tok", "POST", "/v1/messages", A_HARNESS_ASKING);
1894        assert_eq!(head["status"], 200);
1895        assert_eq!(body, THE_ANSWER.as_bytes());
1896        assert_eq!(asked_models(&seen_rx, 1), ["b/two".to_owned()], "the model that refused is not asked again");
1897    }
1898
1899    #[test]
1900    fn a_cooled_model_is_asked_first_again_once_its_minute_is_up() {
1901        let now = Instant::now();
1902        assert_eq!(order_of(A_LINEUP, now + COOLING, now), ["b/two".to_owned(), "a/one".to_owned()]);
1903        assert_eq!(
1904            order_of(A_LINEUP, now + COOLING, now + COOLING + Duration::from_secs(1)),
1905            A_LINEUP.iter().map(|m| (*m).to_owned()).collect::<Vec<_>>(),
1906            "a model that is back is used in the order the person wrote it"
1907        );
1908    }
1909
1910    #[test]
1911    fn a_route_of_models_that_are_all_cooling_is_still_asked_in_the_order_it_was_written() {
1912        let now = Instant::now();
1913        let every: Vec<(String, Option<String>, String)> =
1914            A_LINEUP.iter().map(|m| ("ev".to_owned(), Some("bilim".to_owned()), (*m).to_owned())).collect();
1915        let seen: Cooling = every.iter().map(|step| (step.clone(), now + COOLING)).collect();
1916        let models: Vec<String> = A_LINEUP.iter().map(|m| (*m).to_owned()).collect();
1917        assert_eq!(
1918            ordered(&models, &seen, &|model| ("ev".to_owned(), Some("bilim".to_owned()), model.to_owned()), now),
1919            A_LINEUP.iter().map(|m| (*m).to_owned()).collect::<Vec<_>>(),
1920            "a model that keeps refusing is still the one the person chose"
1921        );
1922    }
1923
1924    #[test]
1925    fn a_cooling_step_is_only_skipped_for_the_provider_and_lineup_it_refused_in() {
1926        let now = Instant::now();
1927        let seen: Cooling = [("ev".to_owned(), Some("bilim".to_owned()), "a/one".to_owned())]
1928            .into_iter()
1929            .map(|step| (step, now + COOLING))
1930            .collect();
1931        let models = vec!["a/one".to_owned()];
1932        let order = |tag: &str, lineup: Option<&str>| {
1933            ordered(&models, &seen, &|model| (tag.to_owned(), lineup.map(str::to_owned), model.to_owned()), now)
1934        };
1935        assert_eq!(order("ev", Some("bilim")), ["a/one".to_owned()], "every step is cooling, so the plain order");
1936        assert_eq!(order("başka", Some("bilim")), ["a/one".to_owned()], "another provider's models are not");
1937        assert_eq!(order("ev", Some("başka")), ["a/one".to_owned()], "another lineup's steps are not");
1938        assert_eq!(order("ev", None), ["a/one".to_owned()], "nor a single model's own");
1939    }
1940
1941    #[test]
1942    fn a_streamed_answer_arrives_in_pieces_rather_than_only_at_the_end() {
1943        let scratch = Scratch::new("stream");
1944        let entry = ollama_entry(None);
1945        let (pieces_tx, pieces_rx) = mpsc::channel();
1946        // A `Receiver` is not `Sync`, and `Upstream` must be; held behind a lock, it is taken out
1947        // once, the one time this stand-in is ever called.
1948        let pieces_rx = Mutex::new(Some(pieces_rx));
1949        let upstream = Upstream::new(move |_| {
1950            let pieces_rx = pieces_rx.lock().expect("the lock").take().expect("the stand-in is called once");
1951            Ok(UpstreamAnswer {
1952                status: 200,
1953                headers: vec![("content-type".to_owned(), "text/event-stream".to_owned())],
1954                body: Box::new(Trickle(pieces_rx, Vec::new())),
1955            })
1956        });
1957        let listener =
1958            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), upstream, |_| {}).expect("the socket opens");
1959
1960        let mut stream = connect(listener.socket());
1961        let head = json!({ "token": "tok", "method": "GET", "path": "/v1/models", "headers": {} }).to_string();
1962        stream.write_all(format!("{head}\n").as_bytes()).expect("the head is written");
1963        stream.shutdown(std::net::Shutdown::Write).expect("the write side closes");
1964        stream.set_read_timeout(Some(std::time::Duration::from_secs(5))).expect("a read timeout");
1965        let mut reading = BufReader::new(stream);
1966        let mut line = String::new();
1967        reading.read_line(&mut line).expect("the answer head comes");
1968        assert_eq!(serde_json::from_str::<Value>(line.trim_end()).expect("json")["status"], 200);
1969
1970        pieces_tx.send(Some(b"first piece".to_vec())).expect("the first piece is sent");
1971        let mut buf = [0u8; 64];
1972        let n = reading.read(&mut buf).expect("the first piece arrives on its own");
1973        assert_eq!(&buf[..n], b"first piece", "only what was sent so far has arrived");
1974
1975        pieces_tx.send(Some(b", second piece".to_vec())).expect("the second piece is sent");
1976        pieces_tx.send(None).expect("the stream ends");
1977        let mut rest = Vec::new();
1978        reading.read_to_end(&mut rest).expect("the rest is read");
1979        assert_eq!(rest, b", second piece", "the second piece follows once it was sent");
1980    }
1981
1982    #[test]
1983    fn the_script_names_the_socket_the_port_and_the_headers_it_strips() {
1984        assert!(SCRIPT.contains(&format!("\"{SOCKET_NAME}\"")), "the script looks for this socket");
1985        assert!(SCRIPT.contains(&PORT.to_string()), "the script listens on this port");
1986        assert!(SCRIPT.contains(crate::bridge::TOKEN_VARIABLE), "the script reads the same token the bridge does");
1987        assert!(SCRIPT.to_lowercase().contains("authorization"), "the script strips the harness's own header");
1988        assert!(SCRIPT.to_lowercase().contains("x-api-key"), "the script strips the harness's other header");
1989        assert!(SCRIPT.contains("\"api-key\""), "and the one Xiaomi reads a key from");
1990    }
1991
1992    #[test]
1993    fn the_program_of_a_provider_tab_is_the_relay_with_the_harness_after_it() {
1994        let harness = ["claude".to_owned(), "--dangerously-skip-permissions".to_owned()];
1995        assert_eq!(
1996            wrapping(&harness),
1997            ["node", &script_in_container(), "claude", "--dangerously-skip-permissions"],
1998            "the harness is started by the relay, not beside it"
1999        );
2000        assert!(script_in_container().ends_with(SCRIPT_NAME), "the script is where the container sees the folder");
2001    }
2002
2003    /// The whole reason the relay starts the harness itself: a harness started beside it can ask
2004    /// before the server is up. The child here is the harness's stand-in and it does the one
2005    /// thing a harness does first — reach the relay's address — so a script that started it too
2006    /// early would fail this, and one that never started it would hang rather than answer.
2007    #[test]
2008    fn the_script_starts_what_it_is_given_only_once_it_is_listening_and_ends_with_it() {
2009        let Some(node) = node() else { return };
2010        let script = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("assets/provider/qcode-relay.mjs");
2011        // The stand-in reaches the address it was told of, as a harness does; the relay moves it
2012        // to another port when the usual one is taken (by another test, say).
2013        let reach = "const u=new URL(process.env.QCODE_TEST_BASE);\
2014             const s=require('node:net').connect(Number(u.port),'127.0.0.1');\
2015             s.on('connect',()=>{s.end();process.exit(23)});s.on('error',()=>process.exit(1));";
2016        let ran = std::process::Command::new(&node)
2017            .args([
2018                script.as_os_str(),
2019                std::ffi::OsStr::new(&node),
2020                std::ffi::OsStr::new("-e"),
2021                std::ffi::OsStr::new(reach),
2022            ])
2023            .env("QCODE_TEST_BASE", format!("http://127.0.0.1:{PORT}"))
2024            .output()
2025            .expect("the script runs");
2026        assert_eq!(
2027            ran.status.code(),
2028            Some(23),
2029            "the harness reached the relay and its own code came back: {}",
2030            String::from_utf8_lossy(&ran.stderr)
2031        );
2032    }
2033
2034    /// A second tab of the same profile shares its container, and so the relay's usual port. The
2035    /// second relay listens on a free port instead and tells its harness that one, in the
2036    /// environment and on the command line where QCode wrote the usual address; before, it died on
2037    /// the taken port and the tab never started (seen in the endurance trial, 2026-09-24).
2038    #[test]
2039    fn a_relay_whose_port_is_taken_listens_on_another_and_tells_its_harness_so() {
2040        let Some(node) = node() else { return };
2041        // Held for the whole test; if something else on this machine holds it, it is taken all the same.
2042        let _held = std::net::TcpListener::bind(("127.0.0.1", PORT));
2043        let script = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("assets/provider/qcode-relay.mjs");
2044        let address = format!("http://127.0.0.1:{PORT}/v1");
2045        let reach = "const u=new URL(process.env.QCODE_TEST_BASE),a=new URL(process.argv[1]);\
2046             console.log(u.port+' '+a.port);\
2047             const s=require('node:net').connect(Number(u.port),'127.0.0.1');\
2048             s.on('connect',()=>{s.end();process.exit(23)});s.on('error',()=>process.exit(1));";
2049        let ran = std::process::Command::new(&node)
2050            .args([
2051                script.as_os_str(),
2052                std::ffi::OsStr::new(&node),
2053                std::ffi::OsStr::new("-e"),
2054                std::ffi::OsStr::new(reach),
2055                std::ffi::OsStr::new(&address),
2056            ])
2057            .env("QCODE_TEST_BASE", &address)
2058            .output()
2059            .expect("the script runs");
2060        assert_eq!(
2061            ran.status.code(),
2062            Some(23),
2063            "the harness reached its relay: {}",
2064            String::from_utf8_lossy(&ran.stderr)
2065        );
2066        let printed = String::from_utf8_lossy(&ran.stdout);
2067        let ports: Vec<&str> = printed.split_whitespace().collect();
2068        assert_eq!(ports.len(), 2, "{printed}");
2069        assert_eq!(ports[0], ports[1], "the environment and the arguments name the same port: {printed}");
2070        assert_ne!(ports[0], PORT.to_string(), "not the taken one: {printed}");
2071    }
2072
2073    /// Node, when this machine has one. The script is run in a profile image, which always has
2074    /// it; a machine building QCode need not, and a missing node is not a failing QCode.
2075    fn node() -> Option<std::path::PathBuf> {
2076        let found = std::process::Command::new("sh").args(["-c", "command -v node"]).output().ok()?;
2077        found.status.success().then(|| std::path::PathBuf::from(String::from_utf8_lossy(&found.stdout).trim()))
2078    }
2079
2080    #[test]
2081    fn every_path_the_script_or_the_host_disagrees_on_is_refused_the_same_way() {
2082        assert!(allowed("POST", "/v1/messages"), "what Claude Code sends");
2083        assert!(allowed("POST", "/v1/chat/completions"), "what opencode sends");
2084        assert!(allowed("POST", "/v1/responses"), "what Codex sends");
2085        assert!(!allowed("GET", "/v1/responses"), "a stored answer is not asked for");
2086        assert!(!allowed("POST", "/v1/responses/compact"), "only the conversation itself");
2087        assert!(allowed("GET", "/v1/models"));
2088        assert!(allowed("GET", "/v1/models?x=1"), "a query string does not change what path was asked for");
2089        assert!(!allowed("POST", "/v1/models"));
2090        assert!(!allowed("GET", "/v1/messages"));
2091        assert!(!allowed("GET", "/v1/chat/completions"));
2092        assert!(!allowed("DELETE", "/v1/messages"));
2093        assert!(!allowed("GET", "/anything/else"));
2094        assert!(!allowed("POST", "/v1/completions"), "only a conversation, in either shape");
2095    }
2096
2097    #[test]
2098    fn an_openai_shaped_harness_is_carried_to_openrouter_whatever_shape_the_entry_was_measured_in() {
2099        let scratch = Scratch::new("openai-shape");
2100        let mut entry = ProviderEntry::new(tag("yol"), ProviderKind::OpenRouter, "https://openrouter.ai");
2101        entry.key = Some(Key::new(MADE_UP_KEY).expect("a key"));
2102        assert_eq!(entry.wire, super::super::Wire::Anthropic, "the entry was written down in the other shape");
2103        let (seen_tx, seen_rx) = mpsc::channel();
2104        let listener =
2105            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), canned(200, "{}", seen_tx), |_| {})
2106                .expect("the socket opens");
2107        let (head, _) = ask(listener.socket(), "tok", "POST", "/v1/chat/completions", b"{}");
2108        assert_eq!(head["status"], 200);
2109        let seen = seen_rx.recv().expect("the request reached the stand-in");
2110        assert_eq!(seen.url, "https://openrouter.ai/api/v1/chat/completions");
2111        assert_eq!(seen.secret.expect("the key is added").key.expose(), MADE_UP_KEY);
2112    }
2113
2114    /// A ready-made entry of `kind` at its first address, with the made-up key.
2115    fn ready_made(kind: ProviderKind) -> ProviderEntry {
2116        let mut entry = ProviderEntry::new(tag("hazir"), kind, kind.suggested_base());
2117        entry.key = Some(Key::new(MADE_UP_KEY).expect("a key"));
2118        entry
2119    }
2120
2121    /// Every header a harness could have been given a key of its own in, which must all stop at
2122    /// the relay whichever of them the provider reads.
2123    const HARNESS_OWN: [(&str, &str); 4] = [
2124        ("content-type", "application/json"),
2125        ("authorization", "Bearer harness-own-key"),
2126        ("x-api-key", "harness-own-key"),
2127        ("api-key", "harness-own-key"),
2128    ];
2129
2130    #[test]
2131    fn xiaomi_gets_its_key_in_api_key_at_the_root_each_shape_answers_under() {
2132        let scratch = Scratch::new("mimo");
2133        let (seen_tx, seen_rx) = mpsc::channel();
2134        let entry = ready_made(ProviderKind::MimoTokenPlan);
2135        let listener =
2136            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), canned(200, "{}", seen_tx), |_| {})
2137                .expect("the socket opens");
2138        for (method, path, url) in [
2139            ("POST", "/v1/messages?beta=true", "https://token-plan-ams.xiaomimimo.com/anthropic/v1/messages?beta=true"),
2140            ("POST", "/v1/chat/completions", "https://token-plan-ams.xiaomimimo.com/v1/chat/completions"),
2141            ("GET", "/v1/models", "https://token-plan-ams.xiaomimimo.com/v1/models"),
2142        ] {
2143            let (head, _) = ask_with_headers(listener.socket(), "tok", method, path, &HARNESS_OWN, b"{}");
2144            assert_eq!(head["status"], 200, "{path}");
2145            let seen = seen_rx.recv().expect("the request reached the stand-in");
2146            assert_eq!(seen.url, url, "{path}");
2147            let secret = seen.secret.expect("the key is added");
2148            assert_eq!((secret.header.as_str(), secret.prefix.as_str()), ("api-key", ""), "{path}");
2149            assert_eq!(secret.key.expose(), MADE_UP_KEY);
2150            for own in ["authorization", "x-api-key", "api-key"] {
2151                assert!(
2152                    seen.headers.iter().all(|(name, _)| name != own),
2153                    "{path}: the harness's own {own} never reaches the provider: {:?}",
2154                    seen.headers.iter().map(|(name, _)| name).collect::<Vec<_>>()
2155                );
2156            }
2157            assert!(seen.headers.iter().any(|(name, _)| name == "content-type"), "{path}: the rest is carried");
2158        }
2159    }
2160
2161    #[test]
2162    fn kimi_gets_its_key_in_x_api_key_under_coding_in_both_shapes() {
2163        let scratch = Scratch::new("kimi");
2164        let (seen_tx, seen_rx) = mpsc::channel();
2165        let entry = ready_made(ProviderKind::KimiCode);
2166        let listener =
2167            Listener::open(&scratch.0, resolve_one("tok", entry, CHOSEN), canned(200, "{}", seen_tx), |_| {})
2168                .expect("the socket opens");
2169        for (path, url) in [
2170            ("/v1/messages", "https://api.kimi.com/coding/v1/messages"),
2171            ("/v1/chat/completions", "https://api.kimi.com/coding/v1/chat/completions"),
2172        ] {
2173            let (head, body) = ask_with_headers(listener.socket(), "tok", "POST", path, &HARNESS_OWN, b"{}");
2174            assert_eq!(head["status"], 200);
2175            assert!(!format!("{head}").contains(MADE_UP_KEY) && !String::from_utf8_lossy(&body).contains(MADE_UP_KEY));
2176            let seen = seen_rx.recv().expect("the request reached the stand-in");
2177            assert_eq!(seen.url, url);
2178            let secret = seen.secret.expect("the key is added");
2179            assert_eq!((secret.header.as_str(), secret.prefix.as_str()), ("x-api-key", ""));
2180            assert!(
2181                seen.headers.iter().all(|(name, value)| !value.contains("harness-own-key") && name != "x-api-key"),
2182                "only QCode's key goes out: {:?}",
2183                seen.headers
2184            );
2185        }
2186    }
2187
2188    /// An upstream that holds every request until the test has let as many go, which is what makes a
2189    /// request in flight observable: without a stand-in that waits, every request here would be
2190    /// finished before a test could look at it.
2191    fn held(arrived: Sender<()>, released: Arc<(Mutex<usize>, Condvar)>) -> Upstream {
2192        let next = Arc::new(AtomicUsize::new(0));
2193        Upstream::new(move |_| {
2194            let at = next.fetch_add(1, Ordering::SeqCst);
2195            arrived.send(()).expect("the test is waiting for it");
2196            let (free, wake) = &*released;
2197            let mut free = free.lock().expect("the lock");
2198            while *free <= at {
2199                free = wake.wait(free).expect("the lock");
2200            }
2201            Ok(UpstreamAnswer { status: 200, headers: Vec::new(), body: Box::new(io::empty()) })
2202        })
2203    }
2204
2205    /// Lets `count` of the requests [`held`] is holding go.
2206    fn release(released: &Arc<(Mutex<usize>, Condvar)>, count: usize) {
2207        let (free, wake) = &**released;
2208        *free.lock().expect("the lock") = count;
2209        wake.notify_all();
2210    }
2211
2212    /// The handle a screen would keep, and what it says while a request is on its way to a model
2213    /// and once it is not.
2214    #[test]
2215    fn a_request_on_its_way_to_a_model_is_busy_and_then_is_not() {
2216        let scratch = Scratch::new("busy");
2217        let (arrived_tx, arrived_rx) = mpsc::channel();
2218        let released = Arc::new((Mutex::new(0), Condvar::new()));
2219        let listener = Listener::open(
2220            &scratch.0,
2221            resolve_one("tok", ollama_entry(None), CHOSEN),
2222            held(arrived_tx, Arc::clone(&released)),
2223            |_| {},
2224        )
2225        .expect("the socket opens");
2226        let activity = listener.activity();
2227
2228        // Nothing has asked yet: a screen reading this before the first turn of a conversation is
2229        // told the tab is idle, which is what it wants to show.
2230        assert!(!activity.busy("tok"), "a tab that has not asked is not busy");
2231        assert_eq!(activity.last_finished("tok"), None, "and has no finished time");
2232
2233        let asking = std::thread::spawn({
2234            let socket = listener.socket().to_owned();
2235            move || ask(&socket, "tok", "POST", "/v1/messages", A_HARNESS_ASKING)
2236        });
2237        arrived_rx.recv().expect("the request reached the provider");
2238        assert!(activity.busy("tok"), "the tab is waiting on the model: this is what a screen shows");
2239        assert_eq!(activity.last_finished("tok"), None, "nothing of this token has finished yet");
2240
2241        release(&released, 1);
2242        let (head, _) = asking.join().expect("the answer came back");
2243        assert_eq!(head["status"], 200);
2244        assert!(!activity.busy("tok"), "the answer is here, so the tab is idle again");
2245        let finished = activity.last_finished("tok").expect("a request that finished has a time");
2246        assert!(finished <= Instant::now(), "and it is a time that has already been");
2247    }
2248
2249    /// A token no tab has is not busy and has no time: the two answers a screen needs before a tab
2250    /// of this workspace has ever made a request.
2251    #[test]
2252    fn a_token_no_tab_has_is_not_busy_and_has_no_time() {
2253        let scratch = Scratch::new("unknown-busy");
2254        let (seen_tx, _seen_rx) = mpsc::channel();
2255        let listener = Listener::open(
2256            &scratch.0,
2257            resolve_one("tok", ollama_entry(None), CHOSEN),
2258            canned(200, "{}", seen_tx),
2259            |_| {},
2260        )
2261        .expect("the socket opens");
2262        let activity = listener.activity();
2263
2264        assert!(!activity.busy("never-asked"), "no tab of that token has asked");
2265        assert_eq!(activity.last_finished("never-asked"), None, "and nothing of it has finished");
2266
2267        // A request refused before its token was read is not a tab asking either: there was no tab
2268        // behind it, so nothing is counted and nothing has a time.
2269        let (head, _) = ask(listener.socket(), "wrong", "GET", "/v1/models", b"");
2270        assert_eq!(head["status"], 401);
2271        assert!(!activity.busy("wrong"), "a refused request leaves no tab busy");
2272        assert_eq!(activity.last_finished("wrong"), None, "and no time behind it");
2273    }
2274
2275    /// Two tabs of one workspace can be on different providers, so what is counted is counted for
2276    /// the token a request came with and for no other: the screen writes what it is told under the
2277    /// tab it belongs to.
2278    #[test]
2279    fn a_request_is_counted_for_its_own_token_and_never_for_another_tabs() {
2280        let scratch = Scratch::new("two-busy");
2281        let (arrived_tx, arrived_rx) = mpsc::channel();
2282        let released = Arc::new((Mutex::new(0), Condvar::new()));
2283        let listener = Listener::open(
2284            &scratch.0,
2285            resolve_one("tok", ollama_entry(None), CHOSEN),
2286            held(arrived_tx, Arc::clone(&released)),
2287            |_| {},
2288        )
2289        .expect("the socket opens");
2290        let activity = listener.activity();
2291
2292        let asking = std::thread::spawn({
2293            let socket = listener.socket().to_owned();
2294            move || ask(&socket, "tok", "POST", "/v1/messages", A_HARNESS_ASKING)
2295        });
2296        arrived_rx.recv().expect("the request reached the provider");
2297        assert!(activity.busy("tok"), "the tab that asked is waiting on the model");
2298        assert!(!activity.busy("other"), "and no other tab of the workspace is");
2299
2300        release(&released, 1);
2301        let _ = asking.join().expect("the answer came back");
2302        assert!(!activity.busy("tok"));
2303        assert!(activity.last_finished("other").is_none(), "the other tab never asked, so it has no time");
2304    }
2305}