use crate::agentloop::action::{SelfHandler, ToolClass};
use crate::agentloop::stop::{Outcome, TerminalStatus};
use crate::intel::client::IntelClient;
use crate::mcp::client::McpClient;
use crate::obs::log::Logger;
use crate::subagent::protocol::ALLOWED_TOOLS_ROLE;
use crate::supervisor::budget::Budget;
use crate::wire::intel::{Message, Request, ToolDef, Usage};
use serde_json::{Value, json};
use std::collections::HashMap;
use std::fmt;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Instant;
const PER_CALL_MAX_TOKENS: u32 = 4096;
const SYSTEM_PROMPT: &str = "You are agentd, an autonomous agent. Accomplish the user's \
instruction by calling the available tools and reasoning over their results. Call a tool when you \
need information or need to act. When the task is complete, reply with your final answer and do \
NOT call a tool. If the task cannot be done, say so plainly. Be concise and factual.";
#[derive(Debug)]
pub enum LoopAbort {
Intel(String),
Mcp(String),
}
impl fmt::Display for LoopAbort {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
LoopAbort::Intel(m) => write!(f, "intelligence: {m}"),
LoopAbort::Mcp(m) => write!(f, "mcp: {m}"),
}
}
}
pub struct LoopInput {
pub instruction: String,
pub output_contract: Option<String>,
pub seed: Vec<(String, String)>,
pub model: String,
pub max_steps: u32,
pub max_tokens: u64,
pub deadline: Instant,
pub cancel: Option<Arc<AtomicBool>>,
}
pub struct Session<'a> {
servers: &'a [McpClient],
tools: Vec<ToolDef>,
tool_to_server: HashMap<String, usize>,
resources: ResourceCatalogue,
model: String,
messages: Vec<Message>,
allowed: Vec<Vec<String>>,
}
impl<'a> Session<'a> {
pub fn prepare(
servers: &'a [McpClient],
input: &LoopInput,
self_handler: &mut dyn SelfHandler,
) -> Result<Session<'a>, LoopAbort> {
let allowed = seed_grants(&input.seed);
let (mut tools, mut tool_to_server) = build_catalogue(servers)?;
let code = crate::tools::defs();
if !code.is_empty() {
tools.retain(|t| !crate::tools::is_registered(&t.name));
tools.extend(code);
}
tools.extend(self_handler.tools());
let resources = collect_resources(servers);
if !resources.owner.is_empty() || self_handler.serves_self_resources() {
tools.push(resource_read_tool_def());
}
narrow_catalogue(&allowed, &mut tools, &mut tool_to_server);
let mut messages = vec![Message::system(system_prompt(
input.output_contract.as_deref(),
))];
if let Some(note) = resources.catalogue_note() {
messages.push(Message::system(note));
}
for (role, content) in &input.seed {
if role == ALLOWED_TOOLS_ROLE {
continue;
}
messages.push(seed_message(role, content));
}
messages.push(Message::user(&input.instruction));
Ok(Session {
servers,
tools,
tool_to_server,
resources,
model: input.model.clone(),
messages,
allowed,
})
}
pub fn refresh_tools(&mut self, self_handler: &mut dyn SelfHandler) -> Result<(), LoopAbort> {
let (mut tools, mut tool_to_server) = build_catalogue(self.servers)?;
let code = crate::tools::defs();
if !code.is_empty() {
tools.retain(|t| !crate::tools::is_registered(&t.name));
tools.extend(code);
}
tools.extend(self_handler.tools());
if !self.resources.owner.is_empty() || self_handler.serves_self_resources() {
tools.push(resource_read_tool_def());
}
narrow_catalogue(&self.allowed, &mut tools, &mut tool_to_server);
self.tools = tools;
self.tool_to_server = tool_to_server;
Ok(())
}
pub fn tools_len(&self) -> usize {
self.tools.len()
}
pub fn tool_class(&self, name: &str) -> ToolClass {
if crate::tools::is_registered(name) {
ToolClass::Code
} else if self.tool_to_server.contains_key(name) {
ToolClass::Mcp
} else {
ToolClass::SelfControl
}
}
pub fn tool_permitted(&self, name: &str) -> bool {
grant_permits(&self.allowed, name)
}
pub fn deliver(&mut self, content: &str) {
self.messages.push(Message::user(content));
}
pub fn set_model(&mut self, model: &str) {
self.model = model.to_string();
}
pub fn model(&self) -> &str {
&self.model
}
pub fn transcript_len(&self) -> usize {
self.messages.len()
}
pub fn truncate_transcript(&mut self, len: usize) {
if len < self.messages.len() {
self.messages.truncate(len);
}
}
pub fn run_turn(
&mut self,
intel: &IntelClient,
self_handler: &mut dyn SelfHandler,
log: &Logger,
budget: &mut Budget,
cancel: Option<&Arc<AtomicBool>>,
) -> Result<(Outcome, Usage), LoopAbort> {
let mut last_text: Option<String> = None;
let run_start = crate::obs::otel::now_unix_nanos();
let mut run_span = crate::obs::otel::run_begin(log.ctx().trace_id.as_deref(), run_start);
let (mut tok_in, mut tok_out) = (0u64, 0u64);
log.info(
"loop.start",
json!({"tools": self.tools.len(), "servers": self.servers.len(), "resources": self.resources.owner.len(), "max_steps": budget.max_steps()}),
);
loop {
if cancel.is_some_and(|c| c.load(Ordering::Relaxed)) {
log.warn(
"loop.final",
json!({"status": "cancelled", "steps": budget.steps()}),
);
run_span.finish(&self.model, tok_in, tok_out, false);
return Ok((
Outcome {
status: TerminalStatus::Cancelled,
partial: last_text.is_some(),
result: json!(last_text.unwrap_or_default()),
scheduled: self_handler.take_scheduled(),
subscriptions: self_handler.take_subscriptions(),
},
Usage {
input_tokens: tok_in,
output_tokens: tok_out,
},
));
}
if let Some(status) = budget.exceeded() {
log.warn("loop.final", json!({"status": status.as_str(), "steps": budget.steps(), "tokens": budget.tokens()}));
run_span.finish(&self.model, tok_in, tok_out, false);
return Ok((
Outcome {
status,
partial: last_text.is_some(),
result: json!(last_text.unwrap_or_default()),
scheduled: self_handler.take_scheduled(),
subscriptions: self_handler.take_subscriptions(),
},
Usage {
input_tokens: tok_in,
output_tokens: tok_out,
},
));
}
log.debug(
"loop.step",
json!({"step": budget.steps(), "tokens": budget.tokens(), "messages": self.messages.len()}),
);
let req = Request {
model: self.model.clone(),
messages: self.messages.clone(),
tools: self.tools.clone(),
max_tokens: PER_CALL_MAX_TOKENS,
temperature: Some(0.0),
};
log.debug(
"intel.call",
json!({"step": budget.steps(), "messages": self.messages.len()}),
);
let chat_start = crate::obs::otel::now_unix_nanos();
let resp = intel
.complete(&req)
.map_err(|e| LoopAbort::Intel(e.to_string()))?;
budget.record_usage(resp.usage);
budget.record_step();
tok_in += resp.usage.input_tokens;
tok_out += resp.usage.output_tokens;
run_span.record_chat(
&self.model,
resp.usage.input_tokens,
resp.usage.output_tokens,
true,
chat_start,
);
log.debug(
"intel.result",
json!({"tool_calls": resp.tool_calls.len(), "tokens_in": resp.usage.input_tokens, "tokens_out": resp.usage.output_tokens}),
);
if resp.wants_tools() {
if let Some(t) = resp.text.as_deref().filter(|t| !t.is_empty()) {
last_text = Some(t.to_string());
}
let tool_calls = resp.tool_calls.clone();
self.messages.push(Message::Assistant {
text: resp.text,
tool_calls: tool_calls.clone(),
});
for tc in &tool_calls {
let mut call = json!({"tool": tc.name, "id": tc.id});
if log.content_capture() {
call["args"] = json!(truncate_for_log(&tc.arguments.to_string()));
}
log.info("tool.call", call);
let tool_start = crate::obs::otel::now_unix_nanos();
let (content, is_error) = if !self.tool_permitted(&tc.name) {
(
format!(
"error: tool '{}' is not in this subagent's allowed tools",
tc.name
),
true,
)
} else if tc.name == "resource.read" {
let uri = tc
.arguments
.get("uri")
.and_then(Value::as_str)
.unwrap_or("")
.trim();
if uri.starts_with("agentd://") || uri.starts_with("agent://") {
self_handler.read_resource(uri).unwrap_or_else(|| {
(format!("unknown agentd resource: {uri}"), true)
})
} else {
read_resource_tool(self.servers, &self.resources.owner, &tc.arguments)
}
} else {
match self_handler.handle(&tc.name, &tc.arguments) {
Some(r) => r, None => match crate::tools::dispatch(&tc.name, &tc.arguments) {
Some(r) => r,
None => dispatch_tool(
self.servers,
&self.tool_to_server,
&tc.name,
&tc.arguments,
),
},
}
};
run_span.record_tool(&tc.name, !is_error, tool_start);
let mut result =
json!({"tool": tc.name, "is_error": is_error, "bytes": content.len()});
if log.content_capture() {
result["content"] = json!(truncate_for_log(&content));
}
log.info("tool.result", result);
self.messages
.push(Message::tool_result(&tc.id, content, is_error));
}
continue;
}
let text = resp.text.clone().or(last_text).unwrap_or_default();
self.messages.push(Message::Assistant {
text: Some(text.clone()),
tool_calls: Vec::new(),
});
log.info(
"loop.final",
json!({"status": "completed", "steps": budget.steps(), "tokens": budget.tokens()}),
);
run_span.finish(&self.model, tok_in, tok_out, true);
return Ok((
Outcome {
status: TerminalStatus::Completed,
partial: false,
result: json!(text),
scheduled: self_handler.take_scheduled(),
subscriptions: self_handler.take_subscriptions(),
},
Usage {
input_tokens: tok_in,
output_tokens: tok_out,
},
));
}
}
}
pub fn run_loop(
intel: &IntelClient,
servers: &[McpClient],
input: &LoopInput,
self_handler: &mut dyn SelfHandler,
log: &Logger,
) -> Result<(Outcome, Usage), LoopAbort> {
let mut session = Session::prepare(servers, input, self_handler)?;
let mut budget = Budget::new(input.max_steps, input.max_tokens, input.deadline);
session.run_turn(intel, self_handler, log, &mut budget, input.cancel.as_ref())
}
const CONTENT_LOG_CAP: usize = 4096;
fn truncate_for_log(s: &str) -> String {
if s.chars().count() <= CONTENT_LOG_CAP {
return s.to_string();
}
let mut t: String = s.chars().take(CONTENT_LOG_CAP).collect();
t.push_str(&format!(
"…(+{} more bytes)",
s.len().saturating_sub(t.len())
));
t
}
fn build_catalogue(
servers: &[McpClient],
) -> Result<(Vec<ToolDef>, HashMap<String, usize>), LoopAbort> {
let mut tools = Vec::new();
let mut routing = HashMap::new();
for (i, server) in servers.iter().enumerate() {
let listed = server
.list_tools()
.map_err(|e| LoopAbort::Mcp(e.to_string()))?;
for t in listed {
routing.entry(t.name.clone()).or_insert(i);
tools.push(ToolDef {
name: t.name,
description: t.description.unwrap_or_default(),
input_schema: t.input_schema,
});
}
}
Ok((tools, routing))
}
fn seed_grants(seed: &[(String, String)]) -> Vec<Vec<String>> {
seed.iter()
.filter(|(role, _)| role == ALLOWED_TOOLS_ROLE)
.map(|(_, content)| crate::subagent::protocol::parse_allowed_tools(content))
.collect()
}
fn grant_permits(grants: &[Vec<String>], name: &str) -> bool {
grants
.iter()
.all(|g| g.iter().any(|p| crate::registry::pattern_matches(p, name)))
}
fn narrow_catalogue(
grants: &[Vec<String>],
tools: &mut Vec<ToolDef>,
routing: &mut HashMap<String, usize>,
) {
if grants.is_empty() {
return;
}
tools.retain(|t| grant_permits(grants, &t.name));
routing.retain(|name, _| grant_permits(grants, name));
}
fn dispatch_tool(
servers: &[McpClient],
routing: &HashMap<String, usize>,
name: &str,
arguments: &Value,
) -> (String, bool) {
match routing.get(name) {
Some(&i) => {
if let Err(e) = crate::mcp::pace::take(servers[i].name()) {
return (e, true);
}
match servers[i].call_tool(name, Some(arguments.clone())) {
Ok(res) => (res.text(), res.is_error()),
Err(e) => (format!("tool transport error: {e}"), true),
}
}
None => (format!("error: no such tool '{name}'"), true),
}
}
fn system_prompt(contract: Option<&str>) -> String {
match contract {
Some(c) if !c.is_empty() => format!("{SYSTEM_PROMPT}\n\nOutput contract:\n{c}"),
_ => SYSTEM_PROMPT.to_string(),
}
}
fn seed_message(role: &str, content: &str) -> Message {
match role {
"system" => Message::system(content),
"assistant" => Message::Assistant {
text: Some(content.to_string()),
tool_calls: Vec::new(),
},
_ => Message::user(content),
}
}
const RESOURCE_CAP: usize = 50;
struct ResourceCatalogue {
owner: HashMap<String, usize>,
entries: Vec<(String, String)>, truncated: bool,
}
impl ResourceCatalogue {
fn catalogue_note(&self) -> Option<String> {
if self.entries.is_empty() {
return None;
}
let mut s = String::from(
"Available MCP resources — read the current content of any with the resource.read tool:\n",
);
for (uri, label) in &self.entries {
if label.is_empty() {
s.push_str(&format!("- {uri}\n"));
} else {
s.push_str(&format!("- {uri} — {label}\n"));
}
}
if self.truncated {
s.push_str(&format!(
"(… more than {RESOURCE_CAP} resources; list truncated)\n"
));
}
Some(s)
}
}
fn collect_resources(servers: &[McpClient]) -> ResourceCatalogue {
let mut owner = HashMap::new();
let mut entries = Vec::new();
let mut truncated = false;
'outer: for (i, s) in servers.iter().enumerate() {
let Ok(list) = s.list_resources() else {
continue;
};
for r in list {
if entries.len() >= RESOURCE_CAP {
truncated = true;
break 'outer;
}
if !owner.contains_key(&r.uri) {
let label = r.title.or(r.name).or(r.description).unwrap_or_default();
owner.insert(r.uri.clone(), i);
entries.push((r.uri, label));
}
}
}
ResourceCatalogue {
owner,
entries,
truncated,
}
}
fn resource_read_tool_def() -> ToolDef {
ToolDef {
name: "resource.read".into(),
description: "Read the current content of an available MCP resource by its uri (see the \
resource catalogue). Use this to pull a resource's body when you need it."
.into(),
input_schema: json!({
"type": "object",
"properties": {"uri": {"type": "string", "description": "the resource uri to read"}},
"required": ["uri"]
}),
}
}
fn read_resource_tool(
servers: &[McpClient],
owner: &HashMap<String, usize>,
args: &Value,
) -> (String, bool) {
let uri = args.get("uri").and_then(Value::as_str).unwrap_or("").trim();
if uri.is_empty() {
return ("error: resource.read requires a 'uri'".into(), true);
}
let candidates: Vec<usize> = match owner.get(uri) {
Some(i) => vec![*i],
None => (0..servers.len()).collect(), };
for i in candidates {
if let Ok(r) = servers[i].read_resource(uri) {
return (r.text(), false);
}
}
(format!("resource.read: no server could read '{uri}'"), true)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn resource_catalogue_note_lists_uris() {
let c = ResourceCatalogue {
owner: HashMap::new(),
entries: vec![
("file:///a.json".into(), "inbox".into()),
("db://orders".into(), String::new()),
],
truncated: false,
};
let note = c.catalogue_note().unwrap();
assert!(note.contains("resource.read"));
assert!(note.contains("file:///a.json — inbox"));
assert!(note.contains("- db://orders\n"));
}
#[test]
fn empty_catalogue_is_no_note() {
let c = ResourceCatalogue {
owner: HashMap::new(),
entries: vec![],
truncated: false,
};
assert!(c.catalogue_note().is_none());
}
#[test]
fn resource_read_rejects_missing_uri() {
let (msg, err) = read_resource_tool(&[], &HashMap::new(), &json!({}));
assert!(err);
assert!(msg.contains("uri"));
}
#[test]
fn resource_read_no_server_is_an_error_observation() {
let (msg, err) = read_resource_tool(&[], &HashMap::new(), &json!({"uri": "file:///x"}));
assert!(err);
assert!(msg.contains("file:///x"));
}
#[test]
fn system_prompt_appends_contract() {
let p = system_prompt(Some("Return JSON."));
assert!(p.contains("Output contract:"));
assert!(p.contains("Return JSON."));
assert_eq!(system_prompt(None), SYSTEM_PROMPT);
}
#[test]
fn a_code_registered_tool_classifies_code_and_wins_a_name_collision() {
let _guard = crate::tools::test_registry_guard();
crate::tools::register(crate::tools::CodeTool::new(
"runner.code_tool",
"a native tool",
json!({"type": "object"}),
|_| Ok(json!("native")),
))
.expect("register");
let mut tool_to_server = HashMap::new();
tool_to_server.insert("runner.code_tool".to_string(), 0usize);
let sess = Session {
servers: &[],
tools: vec![],
tool_to_server,
resources: ResourceCatalogue {
owner: HashMap::new(),
entries: vec![],
truncated: false,
},
model: "m".into(),
messages: vec![],
allowed: Vec::new(),
};
assert_eq!(
sess.tool_class("runner.code_tool"),
ToolClass::Code,
"code wins the collision — a server cannot steal a registered tool's calls"
);
let (content, is_err) =
crate::tools::dispatch("runner.code_tool", &json!({})).expect("code tool dispatches");
assert!(!is_err);
assert_eq!(content, "native");
assert!(crate::tools::unregister("runner.code_tool"));
}
#[test]
fn catalogue_partitions_into_mcp_and_self_control_classes() {
use crate::agentloop::action::SELF_CONTROL_TOOLS;
let mcp = ["db.query", "http.get"];
let mut tool_to_server = HashMap::new();
let mut tools: Vec<ToolDef> = Vec::new();
for n in mcp {
tool_to_server.insert(n.to_string(), 0usize);
tools.push(ToolDef {
name: n.into(),
description: String::new(),
input_schema: json!({}),
});
}
for n in SELF_CONTROL_TOOLS {
tools.push(ToolDef {
name: (*n).into(),
description: String::new(),
input_schema: json!({}),
});
}
let sess = Session {
servers: &[],
tools,
tool_to_server,
resources: ResourceCatalogue {
owner: HashMap::new(),
entries: vec![],
truncated: false,
},
model: "m".into(),
messages: vec![],
allowed: Vec::new(),
};
for n in mcp {
assert_eq!(sess.tool_class(n), ToolClass::Mcp, "{n} is an MCP tool");
}
for n in SELF_CONTROL_TOOLS {
assert_eq!(
sess.tool_class(n),
ToolClass::SelfControl,
"{n} is self/control"
);
}
let (mut n_mcp, mut n_self, mut n_code) = (0usize, 0usize, 0usize);
for t in &sess.tools {
match sess.tool_class(&t.name) {
ToolClass::Mcp => n_mcp += 1,
ToolClass::SelfControl => n_self += 1,
ToolClass::Code => n_code += 1,
}
}
assert_eq!(n_code, 0, "no code tools registered here");
assert_eq!(n_mcp, mcp.len(), "every MCP tool classified");
assert_eq!(
n_self,
SELF_CONTROL_TOOLS.len(),
"every self tool classified"
);
for bad in [
"exec", "shell", "bash", "sh", "command", "system", "eval", "run",
] {
assert!(
!SELF_CONTROL_TOOLS.contains(&bad),
"no local-exec self-tool: {bad}"
);
}
}
#[test]
fn dispatch_unknown_tool_is_error_observation() {
let routing = HashMap::new();
let (content, is_error) = dispatch_tool(&[], &routing, "ghost", &Value::Null);
assert!(is_error);
assert!(content.contains("ghost"));
}
#[test]
fn loop_abort_display() {
assert!(LoopAbort::Intel("down".into()).to_string().contains("down"));
}
#[test]
fn truncate_for_log_caps_and_marks() {
let short = "{\"a\":1}";
assert_eq!(truncate_for_log(short), short); let big = "x".repeat(CONTENT_LOG_CAP + 500);
let out = truncate_for_log(&big);
assert!(out.len() < big.len());
assert!(
out.contains("more bytes"),
"truncation is marked: {}",
&out[out.len() - 32..]
);
let multi = "é".repeat(CONTENT_LOG_CAP + 10);
let _ = truncate_for_log(&multi);
}
#[test]
fn refresh_tools_picks_up_a_changed_handler_catalogue() {
let _guard = crate::tools::test_registry_guard();
struct GrowingHandler {
grown: bool,
}
impl SelfHandler for GrowingHandler {
fn tools(&self) -> Vec<ToolDef> {
let mut t = vec![ToolDef {
name: "alpha".into(),
description: String::new(),
input_schema: Value::Null,
}];
if self.grown {
t.push(ToolDef {
name: "beta".into(),
description: String::new(),
input_schema: Value::Null,
});
}
t
}
fn handle(&mut self, _name: &str, _args: &Value) -> Option<(String, bool)> {
None
}
}
let input = LoopInput {
instruction: "x".into(),
output_contract: None,
seed: Vec::new(),
model: "m".into(),
max_steps: 5,
max_tokens: 1000,
deadline: std::time::Instant::now() + std::time::Duration::from_secs(5),
cancel: None,
};
let mut handler = GrowingHandler { grown: false };
let mut session = Session::prepare(&[], &input, &mut handler).unwrap();
let before = session.tools_len();
let transcript = session.transcript_len();
handler.grown = true;
session.refresh_tools(&mut handler).unwrap();
assert_eq!(session.tools_len(), before + 1, "the new tool is live");
assert_eq!(session.transcript_len(), transcript, "transcript untouched");
assert_eq!(session.tool_class("beta"), ToolClass::SelfControl);
}
#[test]
fn a_seed_grant_narrows_the_catalogue_the_dispatch_and_nothing_else() {
let _guard = crate::tools::test_registry_guard();
struct TwoTools;
impl SelfHandler for TwoTools {
fn tools(&self) -> Vec<ToolDef> {
["alpha", "beta"]
.into_iter()
.map(|n| ToolDef {
name: n.into(),
description: String::new(),
input_schema: Value::Null,
})
.collect()
}
fn handle(&mut self, _name: &str, _args: &Value) -> Option<(String, bool)> {
Some(("served".into(), false))
}
}
let grant = LoopInput {
instruction: "x".into(),
output_contract: None,
seed: vec![
(
crate::subagent::protocol::ALLOWED_TOOLS_ROLE.to_string(),
"[\"alpha\"]".to_string(),
),
("user".to_string(), "a real seed message".to_string()),
],
model: "m".into(),
max_steps: 5,
max_tokens: 1000,
deadline: std::time::Instant::now() + std::time::Duration::from_secs(5),
cancel: None,
};
let mut handler = TwoTools;
let narrowed = Session::prepare(&[], &grant, &mut handler).unwrap();
assert_eq!(narrowed.tools_len(), 1, "only the granted tool is offered");
assert!(narrowed.tool_permitted("alpha"));
assert!(
!narrowed.tool_permitted("beta"),
"a filtered-out tool is refused at dispatch, not served"
);
assert_eq!(narrowed.transcript_len(), 3);
let mut plain = grant;
plain.seed.remove(0);
let wide = Session::prepare(&[], &plain, &mut handler).unwrap();
assert_eq!(wide.tools_len(), 2);
assert!(wide.tool_permitted("beta"));
}
#[cfg(unix)]
mod usage_producer {
use super::*;
use crate::intel::client::IntelClient;
use crate::obs::log::{Comp, Level, LogCtx, Logger};
use std::time::{Duration, Instant};
struct NoopHandler;
impl SelfHandler for NoopHandler {
fn tools(&self) -> Vec<ToolDef> {
Vec::new()
}
fn handle(&mut self, _name: &str, _args: &Value) -> Option<(String, bool)> {
None
}
}
fn test_log() -> Logger {
Logger::new(
LogCtx {
run_id: "r".into(),
agent_id: "0".into(),
agent_path: "0".into(),
comp: Comp::Agent,
pid: 0,
trace_id: None,
},
Level::Error, )
}
fn start_mock_llm(addr_file: &std::path::Path, script: &'static str) -> String {
let s = addr_file.to_str().unwrap().to_string();
std::thread::spawn(move || {
crate::intel::mock::run(&s, script);
});
let deadline = Instant::now() + Duration::from_secs(3);
while !addr_file.exists() {
assert!(Instant::now() < deadline, "mock-llm never announced");
std::thread::sleep(Duration::from_millis(10));
}
let addr = std::fs::read_to_string(addr_file).expect("read mock-llm addr-file");
format!("http://{}", addr.trim())
}
fn input(instruction: &str) -> LoopInput {
LoopInput {
instruction: instruction.into(),
output_contract: None,
seed: Vec::new(),
model: "mock".into(),
max_steps: 8,
max_tokens: 100_000,
deadline: Instant::now() + Duration::from_secs(10),
cancel: None,
}
}
#[test]
fn run_turn_returns_the_turns_token_usage() {
let dir = tempfile::tempdir().unwrap();
let sock = dir.path().join("llm.addr");
let url = start_mock_llm(&sock, "final");
let intel = IntelClient::from_parts(&url, None).unwrap();
let inp = input("do the thing");
let mut handler = NoopHandler;
let mut session = Session::prepare(&[], &inp, &mut handler).unwrap();
let mut budget = Budget::new(inp.max_steps, inp.max_tokens, inp.deadline);
let (outcome, usage) = session
.run_turn(&intel, &mut handler, &test_log(), &mut budget, None)
.expect("turn runs against the mock LLM");
assert_eq!(outcome.status, TerminalStatus::Completed);
assert_eq!(
usage.input_tokens, 11,
"input tokens surfaced from the model"
);
assert_eq!(
usage.output_tokens, 5,
"output tokens surfaced from the model"
);
assert!(usage.total() > 0, "the rolled-up Usage is non-zero");
}
#[test]
fn run_loop_returns_the_runs_total_token_usage() {
let dir = tempfile::tempdir().unwrap();
let sock = dir.path().join("llm.addr");
let url = start_mock_llm(&sock, "read");
let intel = IntelClient::from_parts(&url, None).unwrap();
let inp = input("read the resource");
let mut handler = NoopHandler;
let (outcome, usage) =
run_loop(&intel, &[], &inp, &mut handler, &test_log()).expect("one-shot run");
assert_eq!(outcome.status, TerminalStatus::Completed);
assert_eq!(usage.input_tokens, 22, "summed input over both model calls");
assert_eq!(
usage.output_tokens, 12,
"summed output over both model calls"
);
}
}
}