//! A turn on `loop.sock` and its answer in the thread, against a fake `loopd`: the answer, long //! answers, approvals, errors, an unknown session, and a loop that is not there (M4a spec, section //! 8). Do not edit. #[path = "support/fake_loop.rs"] mod fake_loop; #[path = "support/tmp.rs"] mod tmp; use std::sync::Mutex; use std::time::Duration; use fake_loop::{Reply, done, error, event, serve_loop}; use gatewayd::deliver::{EMPTY_ANSWER, LOOP_DOWN, MAX_POST, Poster, deliver, split_answer}; use gatewayd::mm::MmError; use gatewayd::sessions::{Batch, Thread}; use proto::{DataClass, ErrorCode, SessionId, Timestamp, TurnEvent}; use tmp::TempDir; const DM: &str = "d0000000000000000000000000"; const ROOT: &str = "r0000000000000000000000000"; #[derive(Default)] struct Record { posts: Mutex>, fail: bool, } impl Poster for Record { fn post(&self, channel: &str, root: &str, text: &str) -> Result<(), MmError> { if self.fail { return Err(MmError::Status(500, "down".to_string())); } self.posts .lock() .unwrap() .push((channel.to_string(), root.to_string(), text.to_string())); Ok(()) } } impl Record { fn texts(&self) -> Vec { let posts = self.posts.lock().unwrap(); assert!( posts.iter().all(|(c, r, _)| c == DM && r == ROOT), "every post in the thread" ); posts.iter().map(|(_, _, t)| t.clone()).collect() } } fn batch(resume: bool, text: &str) -> Batch { Batch { session: SessionId::new(&format!("mm-{ROOT}")).unwrap(), thread: Thread { channel: DM.to_string(), root: ROOT.to_string(), }, resume, text: text.to_string(), } } fn run(dir: &TempDir, poster: &Record, batch: &Batch) -> Vec { let log = Mutex::new(Vec::new()); deliver(poster, &dir.path().join("loop.sock"), batch, &|line| { log.lock().unwrap().push(line.to_string()) }); log.into_inner().unwrap() } #[test] fn the_answer_is_posted_and_nothing_else() { let dir = TempDir::new("deliver-answer"); let turns = serve_loop(&dir.path().join("loop.sock"), |_, _| { vec![ event(TurnEvent::Progress { total: 10, cache: 0, processed: 10, }), event(TurnEvent::Reasoning { text: "private thoughts".to_string(), }), event(TurnEvent::ToolCallStarted { name: "read_file".to_string(), }), event(TurnEvent::ToolResult { name: "read_file".to_string(), class: DataClass::Private, truncated: false, }), event(TurnEvent::Content { text: "The ans".to_string(), }), done("The answer."), ] }); let poster = Record::default(); let log = run(&dir, &poster, &batch(true, "one\n\ntwo")); assert_eq!(poster.texts(), ["The answer."]); assert!(log.is_empty(), "{log:?}"); let turn = turns.recv_timeout(Duration::from_secs(5)).unwrap(); assert_eq!( (turn.session.as_str(), turn.content.as_str(), turn.resume), (format!("mm-{ROOT}").as_str(), "one\n\ntwo", true) ); } #[test] fn an_approval_is_announced_once_before_the_answer() { let dir = TempDir::new("deliver-approval"); let expires = Timestamp::from_unix_millis(1_758_650_000_000).unwrap(); let _turns = serve_loop(&dir.path().join("loop.sock"), move |_, _| { vec![ event(TurnEvent::ApprovalPending { approval: 42, tool: "shell".to_string(), expires, }), done("done"), ] }); let poster = Record::default(); run(&dir, &poster, &batch(false, "go")); assert_eq!( poster.texts(), [ "waiting for approval 42: approve or deny it with `bxctl` (Mattermost approvals arrive in M4b)", "done" ] ); } #[test] fn errors_are_posted_with_their_code() { let dir = TempDir::new("deliver-error"); let _turns = serve_loop(&dir.path().join("loop.sock"), |_, _| { vec![error( ErrorCode::Inference, "the model server failed\nsee docs/runbook.md#loopd-selftest-failed", )] }); let poster = Record::default(); run(&dir, &poster, &batch(false, "go")); assert_eq!( poster.texts(), ["Error: inference: the model server failed\nsee docs/runbook.md#loopd-selftest-failed"] ); } #[test] fn a_reply_in_a_thread_loopd_does_not_know_creates_the_session() { let dir = TempDir::new("deliver-unknown"); let turns = serve_loop(&dir.path().join("loop.sock"), |n, _| { if n == 0 { vec![error(ErrorCode::NoSuchSession, "no such session")] } else { vec![done("hello")] } }); let poster = Record::default(); run(&dir, &poster, &batch(true, "hi")); assert_eq!(poster.texts(), ["hello"]); let first = turns.recv_timeout(Duration::from_secs(5)).unwrap(); let second = turns.recv_timeout(Duration::from_secs(5)).unwrap(); assert_eq!( (first.resume, second.resume, second.content.as_str()), (true, false, "hi") ); } #[test] fn a_new_session_is_not_retried() { let dir = TempDir::new("deliver-noretry"); let turns = serve_loop(&dir.path().join("loop.sock"), |_, _| { vec![error(ErrorCode::NoSuchSession, "odd")] }); let poster = Record::default(); run(&dir, &poster, &batch(false, "hi")); assert_eq!(poster.texts(), ["Error: no_such_session: odd"]); assert!(turns.recv_timeout(Duration::from_secs(5)).is_ok()); assert!( turns.recv_timeout(Duration::from_millis(200)).is_err(), "one turn only" ); } #[test] fn a_loop_that_is_not_there_or_goes_away() { let dir = TempDir::new("deliver-down"); let poster = Record::default(); let log = run(&dir, &poster, &batch(false, "hi")); assert_eq!(poster.texts(), [LOOP_DOWN]); assert_eq!(log.len(), 1, "{log:?}"); assert!( log[0].starts_with(&format!("gatewayd: mm-{ROOT}: cannot connect to ")), "{log:?}" ); for (n, replies) in [ vec![ event(TurnEvent::Content { text: "x".to_string(), }), Reply::Close, ], vec![Reply::Bytes(b"\x00\x00\x00\x05{bad}".to_vec())], vec![ Reply::Frame(proto::Envelope { v: 1, id: 2, r#final: true, msg: proto::Message::Ok(proto::Empty {}), }), done("late"), ], ] .into_iter() .enumerate() { let dir = TempDir::new(&format!("deliver-early-{n}")); let replies = Mutex::new(Some(replies)); let _turns = serve_loop(&dir.path().join("loop.sock"), move |_, _| { replies.lock().unwrap().take().unwrap_or_default() }); let poster = Record::default(); let log = run(&dir, &poster, &batch(false, "hi")); assert_eq!(poster.texts(), [LOOP_DOWN], "case {n}"); assert_eq!(log.len(), 1, "case {n}: {log:?}"); } } #[test] fn a_post_that_fails_is_logged() { let dir = TempDir::new("deliver-postfail"); let _turns = serve_loop(&dir.path().join("loop.sock"), |_, _| vec![done("lost")]); let poster = Record { fail: true, ..Record::default() }; let log = run(&dir, &poster, &batch(false, "hi")); assert_eq!( log, [format!( "gatewayd: cannot post in {DM} (thread {ROOT}): status 500: \"down\"" )] ); } #[test] fn long_answers_are_split_at_newlines() { let short = "a".repeat(MAX_POST); assert_eq!(split_answer(&short), [short.as_str()]); let over = "a".repeat(MAX_POST + 1); assert_eq!(split_answer(&over), ["a".repeat(MAX_POST), "a".to_string()]); let lines = format!( "{}\n{}\n{}", "a".repeat(10_000), "b".repeat(5_000), "c".repeat(2_000) ); assert_eq!( split_answer(&lines), [ format!("{}\n{}", "a".repeat(10_000), "b".repeat(5_000)), "c".repeat(2_000) ] ); let wide = "é".repeat(MAX_POST + 5); let parts = split_answer(&wide); assert_eq!( parts.iter().map(|p| p.chars().count()).collect::>(), [MAX_POST, 5], "characters, not bytes" ); let leading = format!("\n{}", "x".repeat(MAX_POST + 1)); let parts = split_answer(&leading); assert!( parts .iter() .all(|p| !p.is_empty() && p.chars().count() <= MAX_POST), "{:?}", parts.iter().map(|p| p.len()).collect::>() ); assert_eq!(parts.concat(), leading, "a hard cut drops nothing"); assert_eq!(split_answer(""), [EMPTY_ANSWER]); assert_eq!(split_answer(" \n "), [EMPTY_ANSWER]); } #[test] fn a_long_answer_is_posted_in_order() { let dir = TempDir::new("deliver-long"); let answer = format!("{}\n{}", "a".repeat(MAX_POST - 1), "b".repeat(MAX_POST)); let sent = answer.clone(); let _turns = serve_loop(&dir.path().join("loop.sock"), move |_, _| vec![done(&sent)]); let poster = Record::default(); run(&dir, &poster, &batch(false, "hi")); assert_eq!( poster.texts(), ["a".repeat(MAX_POST - 1), "b".repeat(MAX_POST)] ); }