use harness_context::default_world;
use harness_core::Task;
use harness_loop::{AgentLoop, TelemetryHook};
use harness_models::{MockModel, MockResponse};
use harness_tools_fs::ReadFile;
use serde_json::json;
use std::io::Write;
use std::sync::{Arc, Mutex};
use tracing_subscriber::fmt::MakeWriter;
#[derive(Clone)]
struct BufWriter(Arc<Mutex<Vec<u8>>>);
impl Write for BufWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> MakeWriter<'a> for BufWriter {
type Writer = BufWriter;
fn make_writer(&'a self) -> BufWriter {
self.clone()
}
}
fn task(desc: &str) -> Task {
Task {
description: desc.into(),
source: None,
deadline: None,
}
}
#[tokio::test]
async fn emits_structured_run_telemetry() {
let buf = Arc::new(Mutex::new(Vec::new()));
let subscriber = tracing_subscriber::fmt()
.with_writer(BufWriter(buf.clone()))
.with_max_level(tracing::Level::DEBUG)
.without_time()
.with_ansi(false)
.finish();
let output = {
let _guard = tracing::subscriber::set_default(subscriber);
let ws = std::env::temp_dir().join(format!("telem-test-{}", std::process::id()));
std::fs::create_dir_all(&ws).unwrap();
let mut world = default_world(&ws);
let model = MockModel::new()
.script(
MockResponse::tool_call("read_file", json!({"path": "x.txt"})).with_usage(100, 20),
)
.script(MockResponse::text("done").with_usage(50, 10));
let outcome = AgentLoop::new(model)
.with_tool(Arc::new(ReadFile))
.with_hook(Arc::new(TelemetryHook::new()))
.run(task("read a file"), &mut world)
.await
.unwrap();
assert!(matches!(outcome, harness_loop::Outcome::Done { .. }));
String::from_utf8(buf.lock().unwrap().clone()).unwrap()
};
assert!(output.contains("agent_run"), "missing run span:\n{output}");
assert!(output.contains("run.start"), "missing run.start:\n{output}");
assert!(
output.contains("model.complete"),
"missing model.complete:\n{output}"
);
assert!(
output.contains("input_tokens=100"),
"missing token field:\n{output}"
);
assert!(
output.contains("gen_ai.usage.input_tokens=100"),
"missing gen_ai usage convention:\n{output}"
);
assert!(
output.contains("gen_ai.tool.name=read_file"),
"missing gen_ai.tool.name convention:\n{output}"
);
assert!(
output.contains("gen_ai.operation.name=\"invoke_agent\""),
"missing invoke_agent on run span:\n{output}"
);
assert!(output.contains("tool.call"), "missing tool.call:\n{output}");
assert!(output.contains("read_file"), "missing tool name:\n{output}");
assert!(output.contains("run.end"), "missing run.end:\n{output}");
assert!(
output.contains("total_tokens=180"),
"run.end must total the turns (100+20+50+10):\n{output}"
);
assert!(
output.contains("model_calls=2"),
"run.end must count model calls:\n{output}"
);
assert!(
output.contains("tool_calls=1") && output.contains("tool_failures=1"),
"run.end must count tool calls and failures:\n{output}"
);
}
#[tokio::test]
async fn compaction_reports_what_it_saved() {
use harness_core::{Budget, CompactError, CompactionStage, Compactor, Context};
use std::sync::atomic::{AtomicU32, Ordering};
struct StubCompactor {
used: AtomicU32,
}
#[async_trait::async_trait]
impl Compactor for StubCompactor {
fn budget(&self, _ctx: &Context) -> Budget {
Budget {
used: self.used.load(Ordering::SeqCst),
window: 1000,
}
}
async fn compact(
&self,
_stage: CompactionStage,
_ctx: &mut Context,
) -> Result<(), CompactError> {
self.used.fetch_sub(200, Ordering::SeqCst);
Ok(())
}
}
let buf = Arc::new(Mutex::new(Vec::new()));
let subscriber = tracing_subscriber::fmt()
.with_writer(BufWriter(buf.clone()))
.with_max_level(tracing::Level::DEBUG)
.without_time()
.with_ansi(false)
.finish();
let output = {
let _guard = tracing::subscriber::set_default(subscriber);
let ws = std::env::temp_dir().join(format!("telem-compact-{}", std::process::id()));
std::fs::create_dir_all(&ws).unwrap();
let mut world = default_world(&ws);
AgentLoop::new(MockModel::new().script(MockResponse::text("done").with_usage(500, 10)))
.with_compactor(Arc::new(StubCompactor {
used: AtomicU32::new(900),
}))
.with_hook(Arc::new(TelemetryHook::new()))
.run(task("go"), &mut world)
.await
.unwrap();
let _ = std::fs::remove_dir_all(&ws);
String::from_utf8(buf.lock().unwrap().clone()).unwrap()
};
assert!(
output.contains("tokens_before=900") && output.contains("tokens_after=700"),
"first stage must report its before/after:\n{output}"
);
assert!(
output.contains("tokens_saved=200"),
"each stage must report what it bought:\n{output}"
);
assert!(
output.contains("compactions=2") && output.contains("tokens_saved=400"),
"run.end must total the compactions and their saving:\n{output}"
);
}
#[tokio::test]
async fn streaming_reports_time_to_first_token() {
use futures::StreamExt;
use harness_core::{
Context, Model, ModelDelta, ModelError, ModelInfo, ModelOutput, StopReason,
};
use std::time::Duration;
struct SlowStart;
#[async_trait::async_trait]
impl Model for SlowStart {
async fn complete(&self, _ctx: &Context) -> Result<ModelOutput, ModelError> {
Ok(ModelOutput {
text: Some("hello world".into()),
stop_reason: StopReason::EndTurn,
..Default::default()
})
}
async fn stream(
&self,
_ctx: &Context,
) -> Result<futures::stream::BoxStream<'static, Result<ModelDelta, ModelError>>, ModelError>
{
let s = futures::stream::iter(vec![
Ok(ModelDelta::Text("hel".into())),
Ok(ModelDelta::Text("lo ".into())),
Ok(ModelDelta::Text("world".into())),
Ok(ModelDelta::Stop(StopReason::EndTurn)),
])
.then(|d| async move {
tokio::time::sleep(Duration::from_millis(40)).await;
d
});
Ok(s.boxed())
}
fn info(&self) -> ModelInfo {
ModelInfo {
handle: "slow".into(),
provider: "test".into(),
model: "slow".into(),
context_window: 8192,
input_cost_usd_per_million_tokens: None,
output_cost_usd_per_million_tokens: None,
supports_tool_use: false,
supports_streaming: true,
supports_web_grounding: false,
}
}
}
let buf = Arc::new(Mutex::new(Vec::new()));
let subscriber = tracing_subscriber::fmt()
.with_writer(BufWriter(buf.clone()))
.with_max_level(tracing::Level::DEBUG)
.without_time()
.with_ansi(false)
.finish();
let output = {
let _guard = tracing::subscriber::set_default(subscriber);
let ws = std::env::temp_dir().join(format!("telem-ttft-{}", std::process::id()));
std::fs::create_dir_all(&ws).unwrap();
let mut world = default_world(&ws);
AgentLoop::new(SlowStart)
.with_streaming(true)
.with_hook(Arc::new(TelemetryHook::new()))
.run(task("say hello"), &mut world)
.await
.unwrap();
let _ = std::fs::remove_dir_all(&ws);
String::from_utf8(buf.lock().unwrap().clone()).unwrap()
};
assert!(
output.contains("model.first_token"),
"a streamed call must report its first fragment:\n{output}"
);
assert!(
output.contains("first_token_ms="),
"the run summary must carry it:\n{output}"
);
let ttft: u64 = output
.split("first_token_ms=")
.nth(1)
.and_then(|s| s.split_whitespace().next())
.and_then(|s| s.parse().ok())
.expect("first_token_ms in run.end");
let total: u64 = output
.rsplit("duration_ms=")
.next()
.and_then(|s| s.split_whitespace().next())
.and_then(|s| s.parse().ok())
.expect("duration_ms in run.end");
assert!(
ttft < total,
"TTFT ({ttft}ms) must be shorter than the whole run ({total}ms)"
);
}