Files
boxmaker/docs/plans/M4a/files/crates/gatewayd/tests/deliver.rs
T
kyleandClaude Opus 5.5 0339dc13b2 Plan M4a: gatewayd in 15 tasks, with skeletons and given tests
Each task's tests were run against a reference at its end state; the end states were replayed
from master in order with the gate at each step (650 to 762 tests); each skeleton compiles
against its tests and fails them. The reference is kept off this machine. Lessons T27 (every
wait in a test has a limit) and T28 (mutate the reference before hand-over) come from this work.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-09-23 19:05:44 -07:00

305 lines
9.3 KiB
Rust

//! 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<Vec<(String, String, String)>>,
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<String> {
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<String> {
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::<Vec<_>>(),
[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::<Vec<_>>()
);
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)]
);
}