use super::*;
use std::sync::Arc;
use leviath_agent_client::PROTOCOL_VERSION;
use leviath_core::interaction::{InteractionKind, InteractionRequest};
use leviath_runtime::control_socket::{ControlClient, bind_control_listener, control_id};
use leviath_runtime::host::WorldEvent;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, DuplexStream};
use tokio::task::JoinHandle;
const RUN_ID: &str = "coder-test-run";
const HARNESS_DEFAULT_CWD: &str = "/harness-launch-dir";
struct AbortOnDrop(JoinHandle<()>);
impl Drop for AbortOnDrop {
fn drop(&mut self) {
self.0.abort();
}
}
struct ScriptedDaemon {
client: ControlClient,
_dir: tempfile::TempDir,
_accept: AbortOnDrop,
}
impl ScriptedDaemon {
fn new(
events: Vec<WorldEvent>,
responder: impl Fn(ControlRequest) -> ControlResponse + Send + Sync + 'static,
) -> Self {
let dir = tempfile::tempdir().unwrap();
let id = control_id(dir.path());
let mut listener = bind_control_listener(&id).unwrap();
let events = Arc::new(events);
let responder = Arc::new(responder);
let accept = tokio::spawn(async move {
loop {
let Ok(Some(stream)) = listener.accept().await else {
break;
};
let events = events.clone();
let responder = responder.clone();
tokio::spawn(async move {
let (read_half, mut write_half) = tokio::io::split(stream);
let mut lines = BufReader::new(read_half).lines();
let Ok(Some(line)) = lines.next_line().await else {
return;
};
let Ok(req) = serde_json::from_str::<ControlRequest>(&line) else {
return;
};
match req {
ControlRequest::Subscribe => {
for ev in events.iter() {
let mut out = serde_json::to_string(ev).unwrap();
out.push('\n');
if write_half.write_all(out.as_bytes()).await.is_err() {
return;
}
}
std::future::pending::<()>().await;
}
other => {
let mut out = serde_json::to_string(&responder(other)).unwrap();
out.push('\n');
let _ = write_half.write_all(out.as_bytes()).await;
}
}
});
}
});
Self {
client: ScriptedDaemon::client_at(&dir),
_dir: dir,
_accept: AbortOnDrop(accept),
}
}
fn client_at(dir: &tempfile::TempDir) -> ControlClient {
ControlClient::new(control_id(dir.path()))
}
fn client(&self) -> ControlClient {
self.client.clone()
}
}
struct Harness {
to_server: DuplexStream,
from_server: BufReader<DuplexStream>,
runs_dir: tempfile::TempDir,
_daemon: ScriptedDaemon,
_server: AbortOnDrop,
}
impl Harness {
fn start(daemon: ScriptedDaemon, args: AgentClientArgs) -> Self {
let (to_server, server_in) = tokio::io::duplex(64 * 1024);
let (server_out, from_server) = tokio::io::duplex(1024 * 1024);
let runs_dir = tempfile::tempdir().unwrap();
let control = daemon.client();
let runs_path = runs_dir.path().to_path_buf();
let server = tokio::spawn(async move {
let _ = serve_over(
BufReader::new(server_in),
server_out,
control,
args,
runs_path,
HARNESS_DEFAULT_CWD.to_string(),
)
.await;
});
Self {
to_server,
from_server: BufReader::new(from_server),
runs_dir,
_daemon: daemon,
_server: AbortOnDrop(server),
}
}
async fn send(&mut self, line: &str) {
self.to_server.write_all(line.as_bytes()).await.unwrap();
self.to_server.write_all(b"\n").await.unwrap();
self.to_server.flush().await.unwrap();
}
async fn close_input(&mut self) {
self.to_server.shutdown().await.unwrap();
}
async fn recv(&mut self) -> JsonRpcMessage {
let mut line = String::new();
let n = self.from_server.read_line(&mut line).await.unwrap();
assert_ne!(n, 0, "server closed its output before sending a message");
serde_json::from_str(line.trim()).unwrap()
}
async fn recv_until(&mut self, pred: impl Fn(&JsonRpcMessage) -> bool) -> JsonRpcMessage {
loop {
let msg = self.recv().await;
if pred(&msg) {
return msg;
}
}
}
fn write_meta_status(&self, status: &str) {
let dir = self.runs_dir.path().join(RUN_ID);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("meta.json"), format!(r#"{{"status":"{status}"}}"#)).unwrap();
}
fn write_output(&self, idx: usize, text: &str) {
let path = self
.runs_dir
.path()
.join(RUN_ID)
.join("stages")
.join(idx.to_string())
.join("output.log");
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(path, text).unwrap();
}
}
fn completed(status: &str) -> WorldEvent {
WorldEvent::Completed {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
status: status.to_string(),
final_output: None,
}
}
fn completed_with_answer(content: &str, format: Option<&str>) -> WorldEvent {
WorldEvent::Completed {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
status: "complete".to_string(),
final_output: Some(leviath_core::output::FinalOutput::new(
content,
format.map(str::to_string),
"summary".to_string(),
0,
)),
}
}
fn status_event() -> WorldEvent {
WorldEvent::Status {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
status: "active".to_string(),
stage: "implement".to_string(),
iteration: 1,
tool_calls: 0,
accepts_messages: false,
}
}
fn spawned_event() -> WorldEvent {
WorldEvent::Spawned {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
blueprint: "coder".to_string(),
}
}
fn tokens_event() -> WorldEvent {
WorldEvent::Tokens {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
prompt_tokens: 10,
completion_tokens: 5,
cached_tokens: 0,
cache_write_tokens: 0,
}
}
fn context_event(total: usize, max: usize) -> WorldEvent {
WorldEvent::Context {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
total_tokens: total,
max_tokens: max,
}
}
fn approval_event() -> WorldEvent {
WorldEvent::Interaction {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
request: InteractionRequest {
id: "appr-1".to_string(),
kind: InteractionKind::ToolApproval,
prompt: "Run bash `ls`?".to_string(),
options: vec![],
tool_name: Some("bash".to_string()),
tool_arguments: None,
required: true,
stage_name: "implement".to_string(),
body: None,
body_format: Default::default(),
},
}
}
fn free_text_event() -> WorldEvent {
WorldEvent::Interaction {
run_id: RUN_ID.to_string(),
agent_id: RUN_ID.to_string(),
request: InteractionRequest::free_text("q1", "What color?", "implement", true),
}
}
fn spawn_ok(req: ControlRequest) -> ControlResponse {
match req {
ControlRequest::Spawn { .. } => ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
},
_ => ControlResponse::Ok { ok: true },
}
}
fn blueprint_args() -> (tempfile::TempDir, AgentClientArgs) {
let root = tempfile::tempdir().unwrap();
let dir = root.path().join("coder");
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(
dir.join("agent.leviath"),
r#"
[agent]
name = "coder"
version = "1.0.0"
description = "test"
[stages.implement]
system_prompt = "Do it"
"#,
)
.unwrap();
let args = AgentClientArgs {
agent: Some(dir.to_string_lossy().to_string()),
yolo: false,
no_seed_commands: false,
allow: vec![],
max_depth: None,
output_format: None,
output_instructions: None,
};
(root, args)
}
fn is_result(msg: &JsonRpcMessage) -> bool {
msg.result.is_some() || msg.error.is_some()
}
fn update_kind(msg: &JsonRpcMessage) -> Option<String> {
if msg.method.as_deref() != Some("session/update") {
return None;
}
msg.params
.as_ref()?
.get("update")?
.get("sessionUpdate")?
.as_str()
.map(str::to_string)
}
async fn opened_session(daemon: ScriptedDaemon, with_caps: bool) -> (Harness, tempfile::TempDir) {
let (bp, args) = blueprint_args();
let mut h = Harness::start(daemon, args);
let init = if with_caps {
r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":1,"clientCapabilities":{"fs":{"readTextFile":true}}}}"#
} else {
r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":1}}"#
};
h.send(init).await;
let _ = h.recv().await; h.send(r#"{"jsonrpc":"2.0","id":2,"method":"session/new","params":{"cwd":"/tmp"}}"#)
.await;
let _ = h.recv().await; (h, bp)
}
struct FailingWriter;
impl tokio::io::AsyncWrite for FailingWriter {
fn poll_write(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
_buf: &[u8],
) -> std::task::Poll<std::io::Result<usize>> {
std::task::Poll::Ready(Err(std::io::Error::other("write failed")))
}
fn poll_flush(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::task::Poll::Ready(Ok(()))
}
fn poll_shutdown(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::task::Poll::Ready(Ok(()))
}
}
struct FailAfter {
remaining: std::sync::atomic::AtomicUsize,
}
impl tokio::io::AsyncWrite for FailAfter {
fn poll_write(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
buf: &[u8],
) -> std::task::Poll<std::io::Result<usize>> {
use std::sync::atomic::Ordering;
if self.remaining.load(Ordering::SeqCst) == 0 {
return std::task::Poll::Ready(Err(std::io::Error::other("write failed")));
}
self.remaining.fetch_sub(1, Ordering::SeqCst);
std::task::Poll::Ready(Ok(buf.len()))
}
fn poll_flush(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::task::Poll::Ready(Ok(()))
}
fn poll_shutdown(
self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
) -> std::task::Poll<std::io::Result<()>> {
std::task::Poll::Ready(Ok(()))
}
}
#[tokio::test]
async fn output_failing_mid_turn_ends_the_turn() {
let daemon = ScriptedDaemon::new(vec![status_event()], spawn_ok);
let control = daemon.client();
let runs = tempfile::tempdir().unwrap();
let out = runs
.path()
.join(RUN_ID)
.join("stages")
.join("0")
.join("output.log");
std::fs::create_dir_all(out.parent().unwrap()).unwrap();
std::fs::write(out, "mid-turn output\n").unwrap();
let (_bp, args) = blueprint_args();
let cwd = args.agent.clone().unwrap();
let session_new = serde_json::json!({
"jsonrpc": "2.0",
"id": 2,
"method": "session/new",
"params": {"cwd": cwd},
});
let script = format!(
"{}\n{}\n{}\n",
r#"{"jsonrpc":"2.0","id":1,"method":"initialize"}"#,
session_new,
r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#,
);
let writer = FailAfter {
remaining: std::sync::atomic::AtomicUsize::new(2),
};
let result = tokio::time::timeout(
std::time::Duration::from_secs(5),
serve_over(
BufReader::new(std::io::Cursor::new(script.into_bytes())),
writer,
control,
args,
runs.path().to_path_buf(),
HARNESS_DEFAULT_CWD.to_string(),
),
)
.await;
assert!(result.unwrap().is_ok());
drop(daemon);
}
#[tokio::test]
async fn a_broken_output_stream_winds_the_server_down() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let control = daemon.client();
let runs = tempfile::tempdir().unwrap().path().to_path_buf();
let input = std::io::Cursor::new(
b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"initialize\"}\n".to_vec(),
);
let result = tokio::time::timeout(
std::time::Duration::from_secs(5),
serve_over(
BufReader::new(input),
FailingWriter,
control,
AgentClientArgs::default(),
runs,
HARNESS_DEFAULT_CWD.to_string(),
),
)
.await;
assert!(result.unwrap().is_ok());
drop(daemon);
}
#[tokio::test]
async fn initialize_advertises_agent_identity_and_capabilities() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":1}}"#)
.await;
let resp = h.recv().await;
let result = resp.result.unwrap();
assert_eq!(result["protocolVersion"], PROTOCOL_VERSION);
assert_eq!(result["agentInfo"]["name"], "leviath");
assert_eq!(result["agentCapabilities"]["loadSession"], false);
assert_eq!(
result["agentCapabilities"]["promptCapabilities"]["embeddedContext"],
true
);
h.close_input().await;
}
#[tokio::test]
async fn initialize_tolerates_absent_params() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"initialize"}"#)
.await;
assert_eq!(
h.recv().await.result.unwrap()["protocolVersion"],
PROTOCOL_VERSION
);
h.close_input().await;
}
#[tokio::test]
async fn session_new_opens_a_session_and_returns_its_id() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let (bp, args) = blueprint_args();
let mut h = Harness::start(daemon, args);
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"session/new","params":{"cwd":"/tmp"}}"#)
.await;
let resp = h.recv().await;
assert!(
resp.result.unwrap()["sessionId"]
.as_str()
.unwrap()
.starts_with("coder")
);
drop(bp);
h.close_input().await;
}
#[tokio::test]
async fn empty_cwd_defaults_to_the_launch_directory() {
let captured = Arc::new(std::sync::Mutex::new(None));
let cap = captured.clone();
let daemon = ScriptedDaemon::new(vec![completed("complete")], move |req| match req {
ControlRequest::Spawn { args } => {
*cap.lock().unwrap() = Some(args.workdir.clone());
ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
}
}
_ => ControlResponse::Ok { ok: true },
});
let (_bp, args) = blueprint_args();
let mut h = Harness::start(daemon, args);
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":1}}"#)
.await;
let _ = h.recv().await;
h.send(r#"{"jsonrpc":"2.0","id":2,"method":"session/new","params":{"cwd":""}}"#)
.await;
let _ = h.recv().await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let _ = h.recv_until(is_result).await;
assert_eq!(
captured.lock().unwrap().as_deref(),
Some(HARNESS_DEFAULT_CWD)
);
h.close_input().await;
}
#[tokio::test]
async fn session_new_logs_and_ignores_supplied_mcp_servers() {
crate::test_support::with_tracing(|| {});
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let (_bp, args) = blueprint_args();
let mut h = Harness::start(daemon, args);
h.send(
r#"{"jsonrpc":"2.0","id":1,"method":"session/new","params":{"cwd":"/tmp","mcpServers":[{"name":"x","command":"y"}]}}"#,
)
.await;
assert!(h.recv().await.result.is_some());
h.close_input().await;
}
#[tokio::test]
async fn session_new_errors_when_no_blueprint_resolves() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
let empty = tempfile::tempdir().unwrap();
let msg = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": "session/new",
"params": {"cwd": empty.path().to_string_lossy()},
});
h.send(&msg.to_string()).await;
let resp = h.recv().await;
assert_eq!(resp.error.unwrap().code, error_codes::INVALID_PARAMS);
h.close_input().await;
}
#[tokio::test]
async fn unknown_method_is_method_not_found() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send(r#"{"jsonrpc":"2.0","id":9,"method":"does/not-exist"}"#)
.await;
let resp = h.recv().await;
assert_eq!(resp.id.unwrap(), serde_json::json!(9));
assert_eq!(resp.error.unwrap().code, error_codes::METHOD_NOT_FOUND);
h.close_input().await;
}
#[tokio::test]
async fn invalid_json_is_a_parse_error() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send("this is not json").await;
let resp = h.recv().await;
assert_eq!(resp.error.unwrap().code, error_codes::PARSE_ERROR);
h.close_input().await;
}
#[tokio::test]
async fn blank_lines_and_bare_notifications_are_ignored() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send("").await; h.send(r#"{"jsonrpc":"2.0","method":"initialized"}"#).await; h.send(r#"{"jsonrpc":"2.0","id":7}"#).await; h.send(r#"{"jsonrpc":"2.0","id":8,"method":"initialize"}"#)
.await;
assert_eq!(h.recv().await.id.unwrap(), serde_json::json!(8));
h.close_input().await;
}
#[tokio::test]
async fn prompt_without_a_session_is_an_error() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"hi"}]}}"#)
.await;
assert_eq!(
h.recv().await.error.unwrap().code,
error_codes::INVALID_REQUEST
);
h.close_input().await;
}
#[tokio::test]
async fn prompt_with_no_usable_text_is_an_error() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"image"}]}}"#)
.await;
assert_eq!(
h.recv().await.error.unwrap().code,
error_codes::INVALID_PARAMS
);
h.close_input().await;
}
#[tokio::test]
async fn prompt_streams_output_then_ends_the_turn() {
let daemon = ScriptedDaemon::new(
vec![
spawned_event(),
status_event(),
tokens_event(),
completed("complete"),
],
spawn_ok,
);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_output(0, "working on it\n");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let chunk = h
.recv_until(|m| update_kind(m).as_deref() == Some("agent_message_chunk"))
.await;
assert_eq!(
chunk.params.unwrap()["update"]["content"]["text"],
"working on it\n"
);
let result = h.recv_until(is_result).await;
assert_eq!(result.result.unwrap()["stopReason"], "end_turn");
h.close_input().await;
}
#[tokio::test]
async fn prompt_maps_error_completion_to_refusal() {
let daemon = ScriptedDaemon::new(vec![completed("error")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let result = h.recv_until(is_result).await;
assert_eq!(result.result.unwrap()["stopReason"], "refusal");
h.close_input().await;
}
#[tokio::test]
async fn prompt_maps_cancelled_completion_to_cancelled() {
let daemon = ScriptedDaemon::new(vec![completed("cancelled")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"cancelled"
);
h.close_input().await;
}
#[tokio::test]
async fn context_events_become_usage_updates() {
let daemon = ScriptedDaemon::new(
vec![context_event(120, 8000), completed("complete")],
spawn_ok,
);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let usage = h
.recv_until(|m| update_kind(m).as_deref() == Some("usage_update"))
.await;
let update = &usage.params.unwrap()["update"];
assert_eq!(update["used"], 120);
assert_eq!(update["size"], 8000);
h.close_input().await;
}
#[tokio::test]
async fn events_for_other_runs_are_ignored() {
let mut other = completed("complete");
if let WorldEvent::Completed { run_id, .. } = &mut other {
*run_id = "someone-else".to_string();
}
let daemon = ScriptedDaemon::new(vec![other, completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[test]
fn world_event_run_id_reads_the_run_id_from_a_log_event() {
let ev = WorldEvent::Log {
run_id: "run-log".to_string(),
agent_id: "a".to_string(),
line: "some output".to_string(),
};
assert_eq!(ev.run_id(), "run-log");
}
#[tokio::test]
async fn a_refused_spawn_ends_the_turn_as_refusal() {
let daemon = ScriptedDaemon::new(vec![], |req| match req {
ControlRequest::Spawn { .. } => ControlResponse::Error {
message: "bad blueprint".to_string(),
},
_ => ControlResponse::Ok { ok: true },
});
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"refusal"
);
h.close_input().await;
}
#[tokio::test]
async fn a_second_prompt_is_delivered_as_a_message() {
let daemon = ScriptedDaemon::new(vec![completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"one"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.send(r#"{"jsonrpc":"2.0","id":4,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"two"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_second_prompt_that_cannot_be_delivered_ends_the_turn() {
let daemon = ScriptedDaemon::new(vec![completed("complete")], |req| match req {
ControlRequest::Spawn { .. } => ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
},
ControlRequest::Message { .. } => ControlResponse::Ok { ok: false },
_ => ControlResponse::Ok { ok: true },
});
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"one"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.send(r#"{"jsonrpc":"2.0","id":4,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"two"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn an_interaction_without_capabilities_is_surfaced_but_does_not_end_the_turn() {
let daemon = ScriptedDaemon::new(vec![free_text_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let chunk = h
.recv_until(|m| update_kind(m).as_deref() == Some("agent_message_chunk"))
.await;
assert!(
chunk.params.unwrap()["update"]["content"]["text"]
.as_str()
.unwrap()
.contains("What color?")
);
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_tool_approval_without_capabilities_keeps_the_turn_alive_until_done() {
let daemon = ScriptedDaemon::new(vec![approval_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_run_that_settles_into_complete_interactive_ends_the_turn_via_meta() {
let daemon = ScriptedDaemon::new(vec![status_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_meta_status("complete_interactive");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_done_run_status_is_detected_on_the_poll_tick_with_no_events() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_meta_status("complete");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_waiting_input_run_status_does_not_end_the_turn() {
let daemon = ScriptedDaemon::new(vec![status_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_meta_status("waiting_input");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn an_error_run_status_ends_the_turn_as_refusal_via_meta() {
let daemon = ScriptedDaemon::new(vec![status_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_meta_status("error");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"refusal"
);
h.close_input().await;
}
#[tokio::test]
async fn a_cancelled_run_status_ends_the_turn_as_cancelled_via_meta() {
let daemon = ScriptedDaemon::new(vec![status_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_meta_status("cancelled");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"cancelled"
);
h.close_input().await;
}
#[tokio::test]
async fn a_capable_client_parks_a_non_approval_interaction() {
let daemon = ScriptedDaemon::new(vec![free_text_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let chunk = h
.recv_until(|m| update_kind(m).as_deref() == Some("agent_message_chunk"))
.await;
assert!(
chunk.params.unwrap()["update"]["content"]["text"]
.as_str()
.unwrap()
.contains("What color?")
);
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_tool_approval_with_capabilities_becomes_a_permission_request() {
let daemon = ScriptedDaemon::new(vec![approval_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let perm = h
.recv_until(|m| m.method.as_deref() == Some("session/request_permission"))
.await;
let perm_id = perm.id.clone().unwrap();
assert_eq!(perm.params.unwrap()["toolCall"]["title"], "bash");
h.send(&format!(
r#"{{"jsonrpc":"2.0","id":{},"result":{{"outcome":{{"outcome":"selected","optionId":"allow-once"}}}}}}"#,
perm_id
))
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_permission_response_with_no_result_denies_and_continues() {
let daemon = ScriptedDaemon::new(vec![approval_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let perm = h
.recv_until(|m| m.method.as_deref() == Some("session/request_permission"))
.await;
h.send(&format!(
r#"{{"jsonrpc":"2.0","id":{},"error":{{"code":-32603,"message":"nope"}}}}"#,
perm.id.unwrap()
))
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn unrelated_messages_during_a_permission_wait_are_skipped() {
let daemon = ScriptedDaemon::new(vec![approval_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let perm = h
.recv_until(|m| m.method.as_deref() == Some("session/request_permission"))
.await;
h.send("not json").await;
h.send(r#"{"jsonrpc":"2.0","id":999,"result":{"unrelated":true}}"#)
.await;
h.send(&format!(
r#"{{"jsonrpc":"2.0","id":{},"result":{{"outcome":{{"outcome":"selected","optionId":"allow-always"}}}}}}"#,
perm.id.unwrap()
))
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_cancel_during_a_permission_wait_cancels_the_run() {
let daemon = ScriptedDaemon::new(vec![approval_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let _ = h
.recv_until(|m| m.method.as_deref() == Some("session/request_permission"))
.await;
h.send(r#"{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"s"}}"#)
.await;
h.close_input().await;
let result = h.recv_until(is_result).await;
assert_eq!(result.result.unwrap()["stopReason"], "end_turn");
}
#[tokio::test]
async fn eof_during_a_permission_wait_parks_the_turn() {
let daemon = ScriptedDaemon::new(vec![approval_event()], spawn_ok);
let (mut h, _bp) = opened_session(daemon, true).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let _ = h
.recv_until(|m| m.method.as_deref() == Some("session/request_permission"))
.await;
h.close_input().await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
}
#[tokio::test]
async fn a_cancel_notification_between_turns_cancels_the_run() {
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let flag = cancelled.clone();
let daemon = ScriptedDaemon::new(vec![completed("complete")], move |req| match req {
ControlRequest::Spawn { .. } => ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
},
ControlRequest::Cancel { .. } => {
flag.store(true, std::sync::atomic::Ordering::SeqCst);
ControlResponse::Ok { ok: true }
}
_ => ControlResponse::Ok { ok: true },
});
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let _ = h.recv_until(is_result).await;
h.send(r#"{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"s"}}"#)
.await;
h.send(r#"{"jsonrpc":"2.0","id":5,"method":"initialize"}"#)
.await;
let _ = h.recv_until(is_result).await;
h.close_input().await;
assert!(cancelled.load(std::sync::atomic::Ordering::SeqCst));
}
#[tokio::test]
async fn a_cancel_notification_with_no_active_run_is_a_no_op() {
let daemon = ScriptedDaemon::new(vec![], spawn_ok);
let mut h = Harness::start(daemon, AgentClientArgs::default());
h.send(r#"{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"s"}}"#)
.await;
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"initialize"}"#)
.await;
assert!(h.recv().await.result.is_some());
h.close_input().await;
}
#[tokio::test]
async fn a_cancel_notification_mid_turn_cancels_the_run() {
let cancelled = Arc::new(std::sync::atomic::AtomicBool::new(false));
let flag = cancelled.clone();
let daemon = ScriptedDaemon::new(vec![status_event()], move |req| match req {
ControlRequest::Spawn { .. } => ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
},
ControlRequest::Cancel { .. } => {
flag.store(true, std::sync::atomic::Ordering::SeqCst);
ControlResponse::Ok { ok: true }
}
_ => ControlResponse::Ok { ok: true },
});
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
h.send(r#"{"jsonrpc":"2.0","method":"initialized"}"#).await;
h.send("not json at all").await;
h.send(r#"{"jsonrpc":"2.0","method":"session/cancel","params":{"sessionId":"s"}}"#)
.await;
h.close_input().await;
let _ = h.recv_until(is_result).await;
assert!(cancelled.load(std::sync::atomic::Ordering::SeqCst));
}
#[tokio::test]
async fn an_unreachable_daemon_makes_a_prompt_refuse() {
let bad = ControlClient::new(control_id(std::path::Path::new("/no/such/daemon")));
let (bp, args) = blueprint_args();
let (to_server, server_in) = tokio::io::duplex(64 * 1024);
let (server_out, from_server) = tokio::io::duplex(1024 * 1024);
let runs_dir = tempfile::tempdir().unwrap();
let runs_path = runs_dir.path().to_path_buf();
let server = tokio::spawn(async move {
let _ = serve_over(
BufReader::new(server_in),
server_out,
bad,
args,
runs_path,
HARNESS_DEFAULT_CWD.to_string(),
)
.await;
});
let mut h = Harness {
to_server,
from_server: BufReader::new(from_server),
runs_dir,
_daemon: ScriptedDaemon::new(vec![], spawn_ok),
_server: AbortOnDrop(server),
};
let _bp = bp;
h.send(r#"{"jsonrpc":"2.0","id":1,"method":"session/new","params":{"cwd":"/tmp"}}"#)
.await;
let _ = h.recv().await;
h.send(r#"{"jsonrpc":"2.0","id":2,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"refusal"
);
h.close_input().await;
}
#[tokio::test]
async fn oversized_output_is_split_across_chunks() {
let daemon = ScriptedDaemon::new(vec![status_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
let big = "x".repeat(leviath_agent_client::MAX_FRAME_BYTES + 100);
h.write_output(0, &big);
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let mut assembled = String::new();
loop {
let msg = h.recv().await;
if let Some("agent_message_chunk") = update_kind(&msg).as_deref() {
assembled.push_str(
msg.params.unwrap()["update"]["content"]["text"]
.as_str()
.unwrap(),
);
} else if is_result(&msg) {
break;
}
}
assert_eq!(assembled, big);
h.close_input().await;
}
#[tokio::test]
async fn output_is_flushed_on_the_poll_tick_between_events() {
let dir = tempfile::tempdir().unwrap();
let id = control_id(dir.path());
let mut listener = bind_control_listener(&id).unwrap();
let accept = tokio::spawn(async move {
loop {
let Ok(Some(stream)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let (read_half, mut write_half) = tokio::io::split(stream);
let mut lines = BufReader::new(read_half).lines();
let Ok(Some(line)) = lines.next_line().await else {
return;
};
let req = serde_json::from_str::<ControlRequest>(&line).unwrap();
match req {
ControlRequest::Spawn { .. } => {
let mut out = serde_json::to_string(&ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
})
.unwrap();
out.push('\n');
let _ = write_half.write_all(out.as_bytes()).await;
}
ControlRequest::Subscribe => {
tokio::time::sleep(std::time::Duration::from_millis(400)).await;
let mut out = serde_json::to_string(&completed("complete")).unwrap();
out.push('\n');
let _ = write_half.write_all(out.as_bytes()).await;
std::future::pending::<()>().await;
}
_ => {}
}
});
}
});
let daemon = ScriptedDaemon {
client: ControlClient::new(control_id(dir.path())),
_dir: dir,
_accept: AbortOnDrop(accept),
};
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_output(0, "tick-flushed output\n");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let chunk = h
.recv_until(|m| update_kind(m).as_deref() == Some("agent_message_chunk"))
.await;
assert_eq!(
chunk.params.unwrap()["update"]["content"]["text"],
"tick-flushed output\n"
);
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
#[tokio::test]
async fn a_closed_event_stream_ends_the_turn() {
let dir = tempfile::tempdir().unwrap();
let id = control_id(dir.path());
let mut listener = bind_control_listener(&id).unwrap();
let accept = tokio::spawn(async move {
loop {
let Ok(Some(stream)) = listener.accept().await else {
break;
};
tokio::spawn(async move {
let (read_half, mut write_half) = tokio::io::split(stream);
let mut lines = BufReader::new(read_half).lines();
let Ok(Some(line)) = lines.next_line().await else {
return;
};
let req = serde_json::from_str::<ControlRequest>(&line).unwrap();
if let ControlRequest::Spawn { .. } = req {
let mut out = serde_json::to_string(&ControlResponse::Spawned {
run_id: RUN_ID.to_string(),
})
.unwrap();
out.push('\n');
let _ = write_half.write_all(out.as_bytes()).await;
}
});
}
});
let daemon = ScriptedDaemon {
client: ControlClient::new(control_id(dir.path())),
_dir: dir,
_accept: AbortOnDrop(accept),
};
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
assert_eq!(
h.recv_until(is_result).await.result.unwrap()["stopReason"],
"end_turn"
);
h.close_input().await;
}
mod run_status_helpers {
use super::*;
use leviath_core::run_meta::RunStatus;
#[test]
fn read_run_status_reads_the_persisted_status_ignoring_extra_fields() {
let dir = tempfile::tempdir().unwrap();
let run = dir.path().join("r1");
std::fs::create_dir_all(&run).unwrap();
std::fs::write(
run.join("meta.json"),
r#"{"status":"complete_interactive","run_id":"r1","extra":1}"#,
)
.unwrap();
assert_eq!(
read_run_status(dir.path(), "r1"),
Some(RunStatus::CompleteInteractive)
);
}
#[test]
fn read_run_status_is_none_when_missing_or_malformed() {
let dir = tempfile::tempdir().unwrap();
assert_eq!(read_run_status(dir.path(), "nope"), None);
let run = dir.path().join("bad");
std::fs::create_dir_all(&run).unwrap();
std::fs::write(run.join("meta.json"), "not json").unwrap();
assert_eq!(read_run_status(dir.path(), "bad"), None);
}
}
#[tokio::test]
async fn a_submitted_answer_arrives_as_the_turns_closing_message() {
let daemon = ScriptedDaemon::new(
vec![
status_event(),
completed_with_answer("<report><finding/></report>", Some("xml")),
],
spawn_ok,
);
let (mut h, _bp) = opened_session(daemon, false).await;
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let mut assembled = String::new();
loop {
let msg = h.recv().await;
if let Some("agent_message_chunk") = update_kind(&msg).as_deref() {
assembled.push_str(
msg.params.unwrap()["update"]["content"]["text"]
.as_str()
.unwrap(),
);
} else if is_result(&msg) {
break;
}
}
assert!(
assembled.contains("<report><finding/></report>"),
"got: {assembled}"
);
assert!(assembled.contains("final output (xml)"), "got: {assembled}");
h.close_input().await;
}
#[tokio::test]
async fn a_run_with_no_answer_adds_no_closing_message() {
let daemon = ScriptedDaemon::new(vec![status_event(), completed("complete")], spawn_ok);
let (mut h, _bp) = opened_session(daemon, false).await;
h.write_output(0, "streamed work");
h.send(r#"{"jsonrpc":"2.0","id":3,"method":"session/prompt","params":{"prompt":[{"type":"text","text":"go"}]}}"#)
.await;
let mut assembled = String::new();
loop {
let msg = h.recv().await;
if let Some("agent_message_chunk") = update_kind(&msg).as_deref() {
assembled.push_str(
msg.params.unwrap()["update"]["content"]["text"]
.as_str()
.unwrap(),
);
} else if is_result(&msg) {
break;
}
}
assert_eq!(assembled, "streamed work");
h.close_input().await;
}