1use 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#[derive(Debug, Clone)]
110pub struct Route {
111 pub entry: ProviderEntry,
114 pub models: Vec<String>,
117 pub lineup: Option<String>,
121}
122
123pub const SCRIPT_NAME: &str = "qcode-relay.mjs";
125
126pub const SOCKET_NAME: &str = "relay.sock";
128
129pub const PORT: u16 = 41417;
133
134pub const SCRIPT: &str = include_str!("../../assets/provider/qcode-relay.mjs");
137
138#[must_use]
141pub fn script_in_container() -> String {
142 format!("{}/{SCRIPT_NAME}", crate::base::paths::MCP_DIR)
143}
144
145#[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
157const PATIENCE: Duration = Duration::from_secs(120);
161
162const MOST_HEAD: usize = 64 * 1024;
165
166const MOST_CONNECTIONS: usize = 16;
169
170const MOST_BODY: usize = 64 * 1024 * 1024;
176
177const MOST_SAID: usize = 64 * 1024;
183
184const COOLING: Duration = Duration::from_secs(60);
188
189const STRIPPED: [&str; 7] =
196 ["authorization", "x-api-key", "api-key", "host", "content-length", "transfer-encoding", "connection"];
197
198const REPLAYED_WITHOUT: [&str; 2] = ["content-length", "transfer-encoding"];
203
204#[derive(Debug, Clone, PartialEq, Eq)]
206struct Incoming {
207 token: String,
209 method: String,
211 path: String,
213 headers: Vec<(String, String)>,
215}
216
217#[derive(Debug, Clone, Copy, PartialEq, Eq)]
219struct Malformed;
220
221fn 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
242pub const RESPONSES_PATH: &str = "/v1/responses";
247
248fn path_only(path: &str) -> &str {
251 path.split('?').next().unwrap_or(path)
252}
253
254fn 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
265fn 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
277fn refusal_body(why: &str) -> Vec<u8> {
281 json!({ "error": why }).to_string().into_bytes()
282}
283
284fn bound(count: usize) -> u64 {
286 u64::try_from(count).unwrap_or(u64::MAX)
287}
288
289fn 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
304pub struct UpstreamAsk {
306 pub method: Method,
308 pub url: String,
310 pub headers: Vec<(String, String)>,
312 pub secret: Option<Secret>,
314 pub body: Box<dyn Read + Send>,
316}
317
318pub struct UpstreamAnswer {
321 pub status: u16,
323 pub headers: Vec<(String, String)>,
325 pub body: Box<dyn Read + Send>,
327}
328
329#[derive(Debug, Clone, PartialEq, Eq)]
331pub struct UpstreamError {
332 pub url: String,
335 pub reason: String,
337}
338
339#[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 #[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 #[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 pub(super) fn call(&self, ask: UpstreamAsk) -> Result<UpstreamAnswer, UpstreamError> {
399 (self.0)(ask)
400 }
401}
402
403#[derive(Debug, Clone, PartialEq, Eq)]
407pub enum Event {
408 UnknownTab,
410 Malformed,
412 NotAllowed {
414 method: String,
416 path: String,
418 },
419 Unreachable {
421 tag: String,
423 reason: String,
425 },
426 Forwarded {
428 tag: String,
430 path: String,
432 status: u16,
434 },
435 FellBack {
439 token: String,
443 tag: String,
445 lineup: Option<String>,
447 from: String,
449 to: String,
451 status: Option<u16>,
453 },
454}
455
456#[derive(Debug, Clone, Copy, Default)]
459struct Asking {
460 serving: usize,
462 last_finished: Option<Instant>,
464}
465
466#[derive(Debug, Clone, Default)]
474pub struct Activity(Arc<Mutex<HashMap<String, Asking>>>);
475
476impl Activity {
477 #[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 #[must_use]
490 pub fn last_finished(&self, token: &str) -> Option<Instant> {
491 self.0.lock().ok()?.get(token)?.last_finished
492 }
493}
494
495struct Counting {
501 asking: Arc<Mutex<HashMap<String, Asking>>>,
502 token: String,
503}
504
505impl Counting {
506 fn serving(activity: &Activity, token: &str) -> Self {
508 let counting = Self { asking: Arc::clone(&activity.0), token: token.to_owned() };
509 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#[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 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 #[must_use]
569 pub fn socket(&self) -> &Path {
570 &self.socket
571 }
572
573 #[must_use]
578 pub fn activity(&self) -> Activity {
579 self.activity.clone()
580 }
581}
582
583fn 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 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 let cooling = Arc::new(Mutex::new(Cooling::new()));
640 let cooling_in_threads = Arc::clone(&cooling);
641 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 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 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 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 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 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 let url = route.entry.api_address(Wire::of_path(&incoming.path), &incoming.path);
770
771 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 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 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 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 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 if remembered.is_none() && passing(status) {
838 remembered = Some(Passing::Answer { status, headers: answer.headers, said });
839 }
840 continue;
841 }
842 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 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 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 if remembered.is_none() {
905 remembered = Some(Passing::Unreachable { reason: error.reason });
906 }
907 }
908 }
909 }
910 }
911
912 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 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
932type Step = (String, Option<String>, String);
936
937type Cooling = HashMap<Step, Instant>;
939
940fn 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
954fn 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
964fn 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
973fn passing(status: u16) -> bool {
978 matches!(status, 408 | 429) || status >= 500
979}
980
981fn lasting(status: u16, said: &[u8], model: &str) -> bool {
989 matches!(status, 402 | 403) || (matches!(status, 400 | 404) && names(said, model))
990}
991
992enum Passing {
995 Answer {
998 status: u16,
1000 headers: Vec<(String, String)>,
1002 said: Vec<u8>,
1004 },
1005 Unreachable {
1008 reason: String,
1010 },
1011}
1012
1013impl Passing {
1014 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
1025fn 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 const MADE_UP_KEY: &str = "not-a-real-key-0000-wxyz";
1049
1050 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 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 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 fn taken(reported: &Arc<Mutex<Vec<Event>>>) -> Vec<Event> {
1108 reported.lock().expect("the events are not poisoned").clone()
1109 }
1110
1111 struct SeenAsk {
1114 url: String,
1115 headers: Vec<(String, String)>,
1116 secret: Option<Secret>,
1117 body: Vec<u8>,
1118 }
1119
1120 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 type Reply = (u16, &'static str);
1142
1143 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 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 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 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 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 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 const A_HARNESS_ASKING: &[u8] = br#"{"model":"claude-haiku-4-5","max_tokens":1024,"stream":true,"messages":[{"role":"user","content":"selam"}]}"#;
1256
1257 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 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 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 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 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 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 const A_LINEUP: &[&str] = &["a/one", "b/two"];
1464
1465 const NO_SUCH_MODEL: &str = r#"{"error":{"message":"model a/one not found"}}"#;
1468
1469 const THE_ANSWER: &str = r#"{"id":"chatcmpl-kyaz-1","choices":[{"text":"done"}]}"#;
1471
1472 const THREE_MODELS: &[&str] = &["a/one", "b/two", "c/three"];
1475
1476 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 #[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 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 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 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 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 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 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 #[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 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 #[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 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 #[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}