use std::io::{BufRead, BufReader, Read, Write};
use std::net::{TcpListener, TcpStream};
pub fn inprocess(script: &str) -> Result<String, String> {
use std::collections::HashMap;
use std::sync::{Mutex, OnceLock};
static SERVERS: OnceLock<Mutex<HashMap<String, String>>> = OnceLock::new();
let map = SERVERS.get_or_init(|| Mutex::new(HashMap::new()));
if let Some(addr) = map.lock().expect("mock servers").get(script) {
return Ok(addr.clone());
}
let addr_file = std::env::temp_dir().join(format!(
"agentd-mockintel-{}-{}.addr",
std::process::id(),
crate::sha::sha256_hex(script.as_bytes())
.chars()
.take(12)
.collect::<String>()
));
let _ = std::fs::remove_file(&addr_file);
let (af, sc) = (addr_file.to_string_lossy().to_string(), script.to_string());
std::thread::Builder::new()
.name("mock-intel".into())
.spawn(move || run(&af, &sc))
.map_err(|e| format!("mock intelligence: spawn: {e}"))?;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
if let Ok(addr) = std::fs::read_to_string(&addr_file) {
let addr = addr.trim().to_string();
if !addr.is_empty() {
map.lock()
.expect("mock servers")
.insert(script.to_string(), addr.clone());
return Ok(addr);
}
}
if std::time::Instant::now() > deadline {
return Err("mock intelligence: server never announced its address".into());
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
}
pub fn run(addr_file: &str, script: &str) -> i32 {
let listener = match TcpListener::bind("127.0.0.1:0") {
Ok(l) => l,
Err(e) => {
eprintln!("mock-llm: bind 127.0.0.1:0: {e}");
return crate::exit::GENERIC;
}
};
if let Err(e) = crate::announce_addr(addr_file, &listener) {
eprintln!("mock-llm: write {addr_file}: {e}");
return crate::exit::GENERIC;
}
for stream in listener.incoming().flatten() {
let script = script.to_string();
std::thread::spawn(move || handle(stream, &script));
}
0
}
fn handle(mut stream: TcpStream, script: &str) {
let Some(body) = read_request_body(&mut stream) else {
return;
};
if let Some(path) = script.strip_prefix("file:") {
let payload = match std::fs::read_to_string(path)
.map_err(|e| e.to_string())
.and_then(|t| serde_json::from_str::<serde_json::Value>(&t).map_err(|e| e.to_string()))
{
Ok(playbook) => playbook_response(&playbook, &body),
Err(e) => final_answer(&format!("mock-llm: cannot load playbook {path}: {e}")),
};
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
payload.len(),
payload
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
return;
}
if script == "echo-system" {
let sys = serde_json::from_str::<serde_json::Value>(&body)
.ok()
.and_then(|v| {
v["messages"].as_array().and_then(|ms| {
ms.iter()
.find(|m| m["role"] == "system")
.and_then(|m| m["content"].as_str().map(str::to_string))
})
})
.unwrap_or_else(|| "(no system message)".to_string());
let payload = final_answer(&sys);
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
payload.len(),
payload
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
return;
}
let tool_results =
body.matches("\"role\":\"tool\"").count() + body.matches("\"role\": \"tool\"").count();
let saw_tool_result = tool_results > 0;
let script = match script {
"slow" => {
std::thread::sleep(std::time::Duration::from_secs(5));
"final"
}
"hang" => {
std::thread::sleep(std::time::Duration::from_secs(12));
"final"
}
other => other,
};
let payload = if script == "gate" {
match tool_results {
0 => tool_call(
"workflow.define",
r#"{"workflow":{"start":"gate","nodes":{
"gate":{"kind":"human","payload":{"question":"approve the deploy?"},
"timeout_ms":30000,"writes":"verdict",
"edges":{"replied":"done","timeout":"esc"}},
"done":{"kind":"halt","status":"completed","result_from":"verdict"},
"esc":{"kind":"halt","status":"refused"}}}}"#,
),
1 => tool_call("workflow.run", r#"{"workflow_id":"w1"}"#),
_ => final_answer("gate flow complete"),
}
} else {
response_json(script, saw_tool_result)
};
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
payload.len(),
payload
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
}
fn read_request_body(stream: &mut TcpStream) -> Option<String> {
let mut reader = BufReader::new(stream);
let mut content_length = 0usize;
loop {
let mut line = String::new();
if reader.read_line(&mut line).ok()? == 0 {
return None; }
let t = line.trim_end();
if t.is_empty() {
break; }
if let Some(v) = t
.strip_prefix("Content-Length:")
.or_else(|| t.strip_prefix("content-length:"))
{
content_length = v.trim().parse().unwrap_or(0);
}
}
let mut body = vec![0u8; content_length];
reader.read_exact(&mut body).ok()?;
Some(String::from_utf8_lossy(&body).into_owned())
}
fn response_json(script: &str, saw_tool_result: bool) -> String {
match (script, saw_tool_result) {
("read", false) => tool_call("resource.read", r#"{"uri":"file:///in.json"}"#),
("read", true) => final_answer("read complete"),
("mcp-call", false) => tool_call("bench_echo", r#"{"query":"ping"}"#),
("mcp-call", true) => final_answer("bench tool called"),
("shell-call", false) => tool_call("bash", r#"{"command":"echo pong > out.txt"}"#),
("shell-call", true) => final_answer("ran the command"),
("schedule", false) => tool_call(
"schedule",
r#"{"after_seconds":1,"instruction":"follow up"}"#,
),
("schedule", true) => final_answer("scheduled a follow-up"),
("subscribe", false) => tool_call("subscribe", r#"{"uri":"file:///watch.json"}"#),
("subscribe", true) => final_answer("now watching the resource"),
("a2a-delegate", false) => tool_call(
"a2a.delegate",
r#"{"peer":"peer","objective":"summarize the mesh","output_contract":"one line"}"#,
),
("a2a-delegate", true) => final_answer("delegated over a2a"),
("spawn-churn", _) => tool_call(
"subagent.spawn",
r#"{"instruction":"do a trivial subtask","detach":true}"#,
),
("wf-tool", false) => tool_call("billing.refund", r#"{"order_id":"A1"}"#),
("wf-tool", true) => final_answer("refund started"),
("wf-once", false) => tool_call("workflow.run", r#"{"name":"loop"}"#),
("wf-once", true) => final_answer("started the workflow"),
("json", _) => final_answer(r#"{"verdict":"approve","score":9}"#),
_ => final_answer("mock-llm done"),
}
}
fn playbook_response(playbook: &serde_json::Value, body: &str) -> String {
let tool_results =
body.matches("\"role\":\"tool\"").count() + body.matches("\"role\": \"tool\"").count();
let turn = playbook
.get("match")
.and_then(serde_json::Value::as_array)
.and_then(|rules| {
rules.iter().find(|r| {
r.get("when_contains")
.and_then(serde_json::Value::as_str)
.is_some_and(|needle| body.contains(needle))
})
})
.or_else(|| {
let turns = playbook.get("turns")?.as_array()?;
turns.get(tool_results.min(turns.len().saturating_sub(1)))
});
let Some(turn) = turn else {
return final_answer("mock-llm: empty playbook");
};
if let Some(ms) = turn.get("delay_ms").and_then(serde_json::Value::as_u64) {
std::thread::sleep(std::time::Duration::from_millis(ms));
}
let usage = turn.get("usage").cloned();
let mut resp: serde_json::Value = if let Some(calls) =
turn.get("tool_calls").and_then(serde_json::Value::as_array)
{
let tool_calls: Vec<serde_json::Value> = calls
.iter()
.enumerate()
.map(|(i, c)| {
let args = c.get("arguments").cloned().unwrap_or(serde_json::json!({}));
let args = match args {
serde_json::Value::String(s) => s,
other => other.to_string(),
};
serde_json::json!({"id": format!("call_{}", i + 1), "type": "function",
"function": {"name": c.get("name").and_then(serde_json::Value::as_str).unwrap_or(""), "arguments": args}})
})
.collect();
serde_json::json!({
"choices": [{"message": {"role": "assistant", "content": serde_json::Value::Null, "tool_calls": tool_calls},
"finish_reason": "tool_calls"}],
"usage": {"prompt_tokens": 11, "completion_tokens": 7}
})
} else {
let content = match turn.get("content") {
Some(serde_json::Value::String(s)) => s.clone(),
Some(other) => other.to_string(),
None => "mock-llm done".to_string(),
};
serde_json::json!({
"choices": [{"message": {"role": "assistant", "content": content}, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 11, "completion_tokens": 5}
})
};
if let Some(u) = usage {
resp["usage"] = u;
}
resp.to_string()
}
fn final_answer(text: &str) -> String {
serde_json::json!({
"choices": [{"message": {"role": "assistant", "content": text}, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 11, "completion_tokens": 5}
})
.to_string()
}
fn tool_call(name: &str, args: &str) -> String {
serde_json::json!({
"choices": [{
"message": {
"role": "assistant",
"content": serde_json::Value::Null,
"tool_calls": [{"id": "call_1", "type": "function", "function": {"name": name, "arguments": args}}]
},
"finish_reason": "tool_calls"
}],
"usage": {"prompt_tokens": 11, "completion_tokens": 7}
})
.to_string()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::intel::openai;
#[test]
fn final_script_parses_to_a_completed_answer() {
let resp = openai::parse_response(response_json("final", false).as_bytes()).unwrap();
assert_eq!(resp.text.as_deref(), Some("mock-llm done"));
assert!(resp.tool_calls.is_empty());
}
#[test]
fn read_script_calls_then_answers() {
let turn1 = openai::parse_response(response_json("read", false).as_bytes()).unwrap();
assert!(turn1.wants_tools());
assert_eq!(turn1.tool_calls[0].name, "resource.read");
assert_eq!(turn1.tool_calls[0].arguments["uri"], "file:///in.json");
let turn2 = openai::parse_response(response_json("read", true).as_bytes()).unwrap();
assert!(!turn2.wants_tools());
assert_eq!(turn2.text.as_deref(), Some("read complete"));
}
#[test]
fn schedule_script_calls_the_schedule_tool() {
let turn1 = openai::parse_response(response_json("schedule", false).as_bytes()).unwrap();
assert_eq!(turn1.tool_calls[0].name, "schedule");
assert_eq!(turn1.tool_calls[0].arguments["after_seconds"], 1);
}
#[test]
fn spawn_churn_never_converges() {
for saw_tool in [false, true] {
let turn =
openai::parse_response(response_json("spawn-churn", saw_tool).as_bytes()).unwrap();
assert!(turn.wants_tools(), "spawn-churn must keep calling tools");
assert_eq!(turn.tool_calls[0].name, "subagent.spawn");
assert_eq!(
turn.tool_calls[0].arguments["instruction"],
"do a trivial subtask"
);
}
}
#[test]
fn a_playbook_answers_by_match_rule_then_by_turn_index() {
let pb: serde_json::Value = serde_json::json!({
"turns": [
{"tool_calls": [{"name": "memory.set", "arguments": {"key": "k", "value": 1}}]},
{"content": "all done", "usage": {"prompt_tokens": 500, "completion_tokens": 50}}
],
"match": [{"when_contains": "PREFLIGHT", "content": {"intent": "status"}}]
});
let t0 = openai::parse_response(
playbook_response(&pb, r#"{"messages":[{"role":"user","content":"hi"}]}"#).as_bytes(),
)
.unwrap();
assert!(t0.wants_tools());
assert_eq!(t0.tool_calls[0].name, "memory.set");
assert_eq!(t0.tool_calls[0].arguments["value"], 1);
let t1 = openai::parse_response(
playbook_response(&pb, r#"{"messages":[{"role":"tool","content":"ok"}]}"#).as_bytes(),
)
.unwrap();
assert!(!t1.wants_tools());
assert_eq!(t1.text.as_deref(), Some("all done"));
assert_eq!(t1.usage.input_tokens, 500);
assert_eq!(t1.usage.output_tokens, 50);
let t9 = openai::parse_response(
playbook_response(&pb, r#"[{"role":"tool"},{"role":"tool"},{"role":"tool"}]"#)
.as_bytes(),
)
.unwrap();
assert_eq!(t9.text.as_deref(), Some("all done"));
let m = openai::parse_response(
playbook_response(
&pb,
r#"{"messages":[{"role":"system","content":"PREFLIGHT"}]}"#,
)
.as_bytes(),
)
.unwrap();
assert_eq!(m.text.as_deref(), Some(r#"{"intent":"status"}"#));
}
#[test]
fn subscribe_script_calls_the_subscribe_tool() {
let turn1 = openai::parse_response(response_json("subscribe", false).as_bytes()).unwrap();
assert_eq!(turn1.tool_calls[0].name, "subscribe");
assert_eq!(turn1.tool_calls[0].arguments["uri"], "file:///watch.json");
let turn2 = openai::parse_response(response_json("subscribe", true).as_bytes()).unwrap();
assert!(!turn2.wants_tools());
}
}