use crate::template::GitRunner;
use std::fs;
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::path::Path;
use std::process::Command;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tempfile::TempDir;
use super::prompt_end_to_end::{
HAPPY_SSE, scaffold_repo, write_brazen_config, write_global_models,
};
fn lernie_bin() -> std::path::PathBuf {
crate::test_support::lernie_binary()
}
type Received = Arc<Mutex<Vec<serde_json::Value>>>;
fn spawn_recording_server() -> (String, Received) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let received: Received = Arc::new(Mutex::new(Vec::new()));
let sink = Arc::clone(&received);
std::thread::spawn(move || {
for stream in listener.incoming() {
let mut stream = stream.expect("accept");
if let Some(body) = read_http_body(&mut stream)
&& let Ok(value) = serde_json::from_slice::<serde_json::Value>(&body)
{
sink.lock().unwrap().push(value);
}
let resp = format!(
"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\n\
content-length: {}\r\nconnection: close\r\n\r\n{HAPPY_SSE}",
HAPPY_SSE.len()
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
}
});
(format!("http://127.0.0.1:{port}"), received)
}
fn read_http_body(stream: &mut TcpStream) -> Option<Vec<u8>> {
let mut head = Vec::new();
let mut byte = [0u8; 1];
while !head.ends_with(b"\r\n\r\n") {
if stream.read(&mut byte).ok()? == 0 {
return None;
}
head.push(byte[0]);
}
let len: usize = String::from_utf8_lossy(&head)
.to_lowercase()
.lines()
.find_map(|l| {
l.strip_prefix("content-length:")
.map(str::trim)
.map(str::to_owned)
})?
.parse()
.ok()?;
let mut body = vec![0u8; len];
stream.read_exact(&mut body).ok()?;
Some(body)
}
const PARENT: &str = "20260101-p1";
fn git_run(dest: &Path, args: &[&str]) {
crate::template::RealGit::new()
.run(dest, args)
.unwrap_or_else(|e| panic!("git {args:?}: {e}"));
}
fn fabricate_parent_that_used_bash(repo: &Path) {
let bare = repo.join("repo.git");
let worktree = repo.join("agents").join(PARENT);
git_run(
&bare,
&[
"worktree",
"add",
"-b",
&format!("agents/{PARENT}"),
worktree.to_str().unwrap(),
"config/default",
],
);
let messages = worktree.join("messages");
fs::create_dir_all(&messages).unwrap();
fs::write(worktree.join("goal.md"), "run a command\n").unwrap();
fs::write(messages.join("001-user.md"), "run a command\n").unwrap();
fs::write(
messages.join("002-claude-sonnet-5.json"),
serde_json::json!([{
"type": "tool_use",
"id": "toolu_bash",
"name": "bash",
"input": {"command": "echo hello-from-tool"}
}])
.to_string(),
)
.unwrap();
fs::write(
messages.join("003-tool.json"),
serde_json::json!([{
"type": "tool_result",
"tool_use_id": "toolu_bash",
"content": [{"type": "text", "text": "hello-from-tool\n"}],
"is_error": false
}])
.to_string(),
)
.unwrap();
git_run(&worktree, &["add", "-A"]);
git_run(&worktree, &["commit", "-m", "checkpoint"]);
}
fn wait_for_terminal_end(path: &Path, ctx: &Ctx, deadline: Duration) -> Vec<serde_json::Value> {
let start = Instant::now();
loop {
let lines: Vec<serde_json::Value> = fs::read(path)
.unwrap_or_default()
.split(|b| *b == b'\n')
.filter(|l| !l.is_empty())
.filter_map(|l| serde_json::from_slice(l).ok())
.collect();
if lines.last().is_some_and(|e| e["type"] == "end") {
return lines;
}
assert!(
start.elapsed() < deadline,
"the compactor's step never completed: {path:?}"
);
ctx.advance();
}
}
struct Ctx {
repo: std::path::PathBuf,
harness: std::path::PathBuf,
brazen_config: std::path::PathBuf,
agent: String,
}
impl Ctx {
fn advance(&self) {
let _ = Command::new(lernie_bin())
.arg("advance")
.arg(&self.repo)
.arg(&self.agent)
.env("LERNIE_HOME", &self.harness)
.env("BRAZEN_CONFIG", &self.brazen_config)
.output();
std::thread::sleep(Duration::from_millis(100));
}
}
fn wait_for_request(received: &Received, ctx: &Ctx, deadline: Duration) -> serde_json::Value {
let start = Instant::now();
loop {
if let Some(first) = received.lock().unwrap().first() {
return first.clone();
}
assert!(
start.elapsed() < deadline,
"the compactor never reached the provider row"
);
ctx.advance();
}
}
#[test]
fn a_compactor_over_a_tool_using_transcript_reaches_the_wire() {
let (endpoint, received) = spawn_recording_server();
let holder = TempDir::new().unwrap();
let harness = holder.path().join("harness");
fs::create_dir_all(&harness).unwrap();
write_global_models(&harness);
let brazen_config = write_brazen_config(holder.path(), &endpoint);
let repo = holder.path().join("ws");
scaffold_repo(&repo, &harness);
fabricate_parent_that_used_bash(&repo);
let out = Command::new(lernie_bin())
.args(["dispatch", "compactor"])
.arg(&repo)
.arg(PARENT)
.env("LERNIE_HOME", &harness)
.env("BRAZEN_CONFIG", &brazen_config)
.output()
.expect("spawn lernie dispatch compactor");
assert!(
out.status.success(),
"lernie dispatch compactor: {}",
String::from_utf8_lossy(&out.stderr)
);
let child = String::from_utf8(out.stdout).unwrap().trim().to_string();
let step_dir = repo.join(format!("steps/{child}/001"));
let deadline = Duration::from_secs(120);
let ctx = Ctx {
repo: repo.clone(),
harness: harness.clone(),
brazen_config: brazen_config.clone(),
agent: child,
};
let request = wait_for_request(&received, &ctx, deadline);
let names: Vec<&str> = request["tools"]
.as_array()
.expect("a tools array reached the wire")
.iter()
.map(|t| t["name"].as_str().unwrap())
.collect();
assert!(names.contains(&"write_summary"), "{names:?}");
assert!(names.contains(&"mark_for_deletion"), "{names:?}");
assert!(
names.contains(&"bash"),
"the inherited transcript's tool must be declared: {names:?}"
);
assert!(
request["messages"].to_string().contains("toolu_bash"),
"{}",
request["messages"]
);
let lines = wait_for_terminal_end(&step_dir.join("response.json"), &ctx, deadline);
assert!(
lines.iter().all(|e| e["type"] != "error"),
"the compactor's first call must not error: {lines:?}"
);
}