use harn_session_store::{
AppendEvent, CreateSession, MemorySessionStore, SessionEventKind, SessionStore,
};
use serde_json::json;
use super::*;
fn custom(kind: &str) -> SessionEventKind {
SessionEventKind::Custom {
custom_type: kind.to_string(),
}
}
fn transcript_event(kind: &str, metadata: serde_json::Value) -> serde_json::Value {
json!({
"transcript_event": {
"id": format!("event-{kind}"),
"kind": kind,
"role": "assistant",
"text": "",
"metadata": metadata,
}
})
}
fn llm_call(input: i64, output: i64, cost: f64) -> AppendEvent {
AppendEvent::new(
custom("llm_call"),
transcript_event(
"llm_call",
json!({
"cache_read_tokens": 0,
"cache_write_tokens": 10948,
"cost_usd": cost,
"input_tokens": input,
"model": "gpt-5.6-luna",
"output_tokens": output,
"provider": "openai",
}),
),
)
}
fn tool_call(id: &str, name: &str) -> AppendEvent {
AppendEvent::new(
SessionEventKind::ToolCall,
transcript_event(
"tool_call",
json!({
"raw_input": {"file": "test/unit/cart_test.rb", "intent": "read"},
"status": "pending",
"tool_call_id": id,
"tool_name": name,
}),
),
)
}
fn tool_update(id: &str, name: &str, status: &str, duration_ms: i64) -> AppendEvent {
AppendEvent::new(
custom("tool_call_update"),
transcript_event(
"tool_call_update",
json!({
"duration_ms": duration_ms,
"status": status,
"tool_call_id": id,
"tool_name": name,
}),
),
)
}
fn tool_result(id: &str, text: &str) -> AppendEvent {
AppendEvent::new(
SessionEventKind::ToolResult,
json!({
"transcript_event": {
"id": format!("result-{id}"),
"kind": "tool_result",
"role": "tool",
"text": text,
"metadata": {"tool_call_id": id},
}
}),
)
}
fn iteration_start(iteration: i64) -> AppendEvent {
AppendEvent::new(
custom("loop_checkpoint"),
transcript_event(
"loop_checkpoint",
json!({"iteration": iteration, "kind": "iteration_start"}),
),
)
}
fn terminal(final_status: &str, stop_reason: &str) -> AppendEvent {
AppendEvent::new(
custom("agent_run_terminal"),
transcript_event(
"agent_run_terminal",
json!({
"error": null,
"final_status": final_status,
"stop_reason": stop_reason,
"terminal_class": null,
}),
),
)
}
fn user_message(text: &str) -> AppendEvent {
AppendEvent::new(
SessionEventKind::Message,
json!({
"transcript_event": {"kind": "message", "role": "user", "text": text},
"raw_message": {"content": text, "role": "user"},
}),
)
.with_actor("user")
}
async fn capstone_like_store() -> (MemorySessionStore, String) {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession {
id: Some("019fc7e6-3103-7610-81ed-91599858fa1a".to_string()),
..CreateSession::default()
})
.await
.expect("create session");
let id = meta.id.clone();
for event in [
user_message("Migrate the three unit test files."),
iteration_start(1),
llm_call(10951, 112, 0.002872),
tool_call("call_A", "look"),
tool_update("call_A", "look", "in_progress", 0),
tool_result("call_A", "[result of look]\n1\tclass CartTest"),
tool_update("call_A", "look", "completed", 5),
iteration_start(2),
llm_call(20000, 300, 0.004),
tool_call("call_B", "edit"),
tool_update("call_B", "edit", "failed", 12),
terminal("done", "pace_cutoff"),
] {
store.append(&id, event).await.expect("append");
}
(store, id)
}
#[tokio::test]
async fn a_headless_session_projects_the_run_record_no_host_ever_wrote() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
assert_eq!(run.id, id);
assert_eq!(run.workflow_id, AGENT_SESSION_WORKFLOW_ID);
assert_eq!(run.task, "Migrate the three unit test files.");
assert_eq!(run.status, "completed");
assert!(
run.finished_at.is_some(),
"a run with a terminal event has an end time even when its session was never closed"
);
assert_eq!(run.root_run_id.as_deref(), Some(id.as_str()));
let usage = run.usage.as_ref().expect("usage");
assert_eq!(usage.call_count, 2);
assert_eq!(usage.input_tokens, 30951);
assert_eq!(usage.output_tokens, 412);
assert!((usage.total_cost - 0.006872).abs() < 1e-9);
assert_eq!(usage.models, vec!["gpt-5.6-luna".to_string()]);
assert_eq!(
run.metadata.get("stop_reason").and_then(|v| v.as_str()),
Some("pace_cutoff"),
"the reason a run was cut short is the first thing a reader asks for"
);
assert_eq!(
run.metadata.get("iterations").and_then(|v| v.as_u64()),
Some(2)
);
}
#[tokio::test]
async fn tool_calls_join_their_updates_and_results_by_provider_call_id() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
assert_eq!(run.tool_recordings.len(), 2);
let look = &run.tool_recordings[0];
assert_eq!(look.tool_name, "look");
assert_eq!(look.tool_use_id, "call_A");
assert_eq!(
look.duration_ms, 5,
"duration comes from the terminal update, not the in_progress one"
);
assert!(look.result.contains("class CartTest"));
assert_eq!(look.iteration, 1);
assert!(!look.args_hash.is_empty());
let edit = &run.tool_recordings[1];
assert_eq!(edit.tool_name, "edit");
assert_eq!(edit.duration_ms, 12);
assert_eq!(
edit.iteration, 2,
"a call is attributed to the iteration that was open when it was made"
);
}
#[tokio::test]
async fn interleaved_tool_calls_do_not_cross_attribute_results() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
let id = meta.id.clone();
for event in [
iteration_start(1),
tool_call("call_first", "look"),
tool_call("call_second", "run"),
tool_result("call_second", "second result"),
tool_update("call_second", "run", "completed", 900),
tool_result("call_first", "first result"),
tool_update("call_first", "look", "completed", 3),
] {
store.append(&id, event).await.expect("append");
}
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
let by_id: std::collections::BTreeMap<_, _> = run
.tool_recordings
.iter()
.map(|record| (record.tool_use_id.as_str(), record))
.collect();
assert_eq!(by_id["call_first"].result, "first result");
assert_eq!(by_id["call_first"].duration_ms, 3);
assert_eq!(by_id["call_second"].result, "second result");
assert_eq!(by_id["call_second"].duration_ms, 900);
}
#[tokio::test]
async fn a_projection_says_it_is_one_and_names_what_it_could_not_recover() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
let projected = run
.metadata
.get("projected_from")
.expect("a projected record must be identifiable as one");
assert_eq!(
projected.get("source").and_then(|v| v.as_str()),
Some(PROJECTION_SOURCE)
);
assert_eq!(
projected.get("session_id").and_then(|v| v.as_str()),
Some(id.as_str())
);
assert_eq!(
projected.get("session_status").and_then(|v| v.as_str()),
Some("open")
);
let unrecovered: Vec<&str> = projected
.get("not_recoverable_from_session")
.and_then(|v| v.as_array())
.expect("not_recoverable_from_session")
.iter()
.filter_map(|v| v.as_str())
.collect();
assert_eq!(unrecovered, UNRECOVERABLE_FIELDS.to_vec());
}
#[tokio::test]
async fn the_unrecoverable_field_list_matches_what_the_projector_actually_leaves_empty() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
let is_empty: std::collections::BTreeMap<&str, bool> = [
(
"usage.total_duration_ms",
run.usage
.as_ref()
.is_none_or(|usage| usage.total_duration_ms == 0),
),
(
"trace_spans[].duration_ms",
run.trace_spans.iter().all(|span| span.duration_ms == 0),
),
("policy", run.policy == RunRecord::default().policy),
("replay_fixture", run.replay_fixture.is_none()),
("id", run.id.is_empty()),
("task", run.task.is_empty()),
("status", run.status.is_empty()),
("started_at", run.started_at.is_empty()),
("finished_at", run.finished_at.is_none()),
("usage", run.usage.is_none()),
("tool_recordings", run.tool_recordings.is_empty()),
("metadata", run.metadata.is_empty()),
]
.into_iter()
.collect();
let actually_empty: Vec<&str> = is_empty
.iter()
.filter(|(_, empty)| **empty)
.map(|(field, _)| *field)
.collect();
let mut declared = UNRECOVERABLE_FIELDS.to_vec();
declared.sort_unstable();
assert_eq!(
actually_empty, declared,
"UNRECOVERABLE_FIELDS must name exactly the fields this projector leaves at their default"
);
}
#[tokio::test]
async fn child_sessions_project_as_child_runs_from_the_stores_own_lineage() {
let store = MemorySessionStore::default();
let parent = store
.create(CreateSession::default())
.await
.expect("create parent");
for name in ["worker-a", "worker-b"] {
store
.create(CreateSession {
parent_session_id: Some(parent.id.clone()),
persona: Some(name.to_string()),
title: Some(format!("{name} task")),
..CreateSession::default()
})
.await
.expect("create child");
}
let run = project_run_record_from_session(&store, &parent.id)
.await
.expect("project");
assert_eq!(run.child_runs.len(), 2);
let names: Vec<&str> = run
.child_runs
.iter()
.map(|child| child.worker_name.as_str())
.collect();
assert_eq!(names, vec!["worker-a", "worker-b"]);
assert!(run
.child_runs
.iter()
.all(|child| child.parent_session_id.as_deref() == Some(parent.id.as_str())));
}
#[tokio::test]
async fn a_loop_that_ended_in_error_projects_as_a_failed_run() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
store
.append(&meta.id, terminal("error", "tool_failure"))
.await
.expect("append");
let run = project_run_record_from_session(&store, &meta.id)
.await
.expect("project");
assert_eq!(run.status, "failed");
}
#[tokio::test]
async fn a_still_running_session_projects_as_running_rather_than_complete() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
store
.append(&meta.id, iteration_start(1))
.await
.expect("append");
let run = project_run_record_from_session(&store, &meta.id)
.await
.expect("project");
assert_eq!(run.status, "running");
assert!(run.finished_at.is_none());
assert!(
run.usage.is_none(),
"a run with no LLM calls reports no usage rather than a zeroed envelope"
);
}
#[tokio::test]
async fn an_unknown_session_names_the_session_and_how_to_find_a_real_one() {
let store = MemorySessionStore::default();
let error = project_run_record_from_session(&store, "019f-not-a-session")
.await
.expect_err("unknown session must fail");
let message = error.to_string();
assert!(
message.contains("019f-not-a-session"),
"error must name the session that was not found: {message}"
);
assert!(
message.contains("harn session list"),
"error must point at the surface that lists real sessions: {message}"
);
}
#[tokio::test]
async fn listing_a_workspace_without_a_session_store_is_empty_rather_than_an_error() {
let dir = tempfile::tempdir().expect("tempdir");
let sessions = list_session_runs(dir.path(), None).await.expect("list");
assert!(sessions.is_empty());
}
#[tokio::test]
async fn a_result_arriving_after_a_rejection_does_not_clear_it() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
let id = meta.id.clone();
for event in [
tool_call("call_R", "run"),
tool_update("call_R", "run", "rejected", 0),
tool_result("call_R", "[error] the user declined this command"),
] {
store.append(&id, event).await.expect("append");
}
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
let record = &run.tool_recordings[0];
assert!(
record.is_rejected,
"a rejected call stays rejected once its result lands"
);
assert!(record.result.contains("declined"));
}
#[tokio::test]
async fn a_grandchild_reports_the_top_of_its_chain_as_the_root_run() {
let store = MemorySessionStore::default();
let root = store
.create(CreateSession::default())
.await
.expect("create root");
let middle = store
.create(CreateSession {
parent_session_id: Some(root.id.clone()),
..CreateSession::default()
})
.await
.expect("create middle");
let leaf = store
.create(CreateSession {
parent_session_id: Some(middle.id.clone()),
..CreateSession::default()
})
.await
.expect("create leaf");
let projected = project_run_record_from_session(&store, &leaf.id)
.await
.expect("project");
assert_eq!(projected.parent_run_id.as_deref(), Some(middle.id.as_str()));
assert_eq!(
projected.root_run_id.as_deref(),
Some(root.id.as_str()),
"root must be the chain's top, not one hop up"
);
let root_projected = project_run_record_from_session(&store, &root.id)
.await
.expect("project root");
assert_eq!(root_projected.parent_run_id, None);
assert_eq!(
root_projected.root_run_id.as_deref(),
Some(root.id.as_str())
);
}
#[tokio::test]
async fn every_recorded_provider_call_becomes_an_llm_call_span() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
let spans: Vec<_> = run
.trace_spans
.iter()
.filter(|span| span.kind == "llm_call")
.collect();
assert_eq!(
spans.len(),
run.usage.as_ref().expect("usage").call_count as usize,
"the per-call view and the aggregate must agree on how many calls there were"
);
let first = spans[0];
assert_eq!(first.trace_id, id);
assert_eq!(first.name, "gpt-5.6-luna");
assert_eq!(first.cost_usd, Some(0.002872));
assert_eq!(
first.metadata.get("input_tokens").and_then(|v| v.as_i64()),
Some(10951)
);
assert_eq!(
first.metadata.get("model").and_then(|v| v.as_str()),
Some("gpt-5.6-luna")
);
let ids: std::collections::BTreeSet<_> = spans.iter().map(|span| span.span_id).collect();
assert_eq!(ids.len(), spans.len());
}
#[tokio::test]
async fn a_projected_span_declares_that_its_duration_is_not_a_measurement() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
for span in run.trace_spans.iter().filter(|s| s.kind == "llm_call") {
assert_eq!(span.duration_ms, 0);
assert_eq!(
span.metadata
.get("duration_available")
.and_then(|v| v.as_bool()),
Some(false),
"a zero duration must be labelled as absent evidence"
);
assert_eq!(span.ttft_ms, None);
}
}
#[tokio::test]
async fn the_cost_aggregate_is_exact_rather_than_float_accumulated() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
for cost in [0.1, 0.2, 0.3] {
store
.append(&meta.id, llm_call(10, 1, cost))
.await
.expect("append");
}
let run = project_run_record_from_session(&store, &meta.id)
.await
.expect("project");
let usage = run.usage.as_ref().expect("usage");
assert_eq!(
usage.total_cost, 0.6,
"0.1 + 0.2 + 0.3 accumulated as f64 gives 0.6000000000000001"
);
let span_total: f64 = run
.trace_spans
.iter()
.filter_map(|span| span.cost_usd)
.sum();
assert!((span_total - usage.total_cost).abs() < 1e-9);
}
fn llm_call_with_attempts(cost: f64, total: i64, rate_limited: i64) -> AppendEvent {
AppendEvent::new(
custom("llm_call"),
transcript_event(
"llm_call",
json!({
"cache_read_tokens": 0,
"cache_write_tokens": 0,
"cost_usd": cost,
"input_tokens": 100,
"model": "gpt-5.6-luna",
"output_tokens": 10,
"provider": "openai",
"provider_attempts": {
"total": total,
"retries": total - 1,
"rate_limited": rate_limited,
"empty_completion": 0,
"other": total - 1 - rate_limited,
},
}),
),
)
}
#[tokio::test]
async fn retried_provider_requests_are_visible_alongside_the_call_count() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
for event in [
llm_call_with_attempts(0.01, 3, 2),
llm_call_with_attempts(0.01, 1, 0),
llm_call_with_attempts(0.01, 2, 0),
] {
store.append(&meta.id, event).await.expect("append");
}
let run = project_run_record_from_session(&store, &meta.id)
.await
.expect("project");
assert_eq!(
run.usage.as_ref().expect("usage").call_count,
3,
"three logical calls"
);
let attempts = run
.metadata
.get("provider_attempts")
.expect("a run that retried must say so");
assert_eq!(attempts.get("total").and_then(|v| v.as_i64()), Some(6));
assert_eq!(attempts.get("retries").and_then(|v| v.as_i64()), Some(3));
assert_eq!(
attempts.get("rate_limited").and_then(|v| v.as_i64()),
Some(2),
"rate limiting is the signal that explains a slow or truncated run"
);
assert_eq!(attempts.get("other").and_then(|v| v.as_i64()), Some(1));
}
#[tokio::test]
async fn a_run_that_never_retried_reports_no_attempt_block() {
let (store, id) = capstone_like_store().await;
let run = project_run_record_from_session(&store, &id)
.await
.expect("project");
assert!(
!run.metadata.contains_key("provider_attempts"),
"no retries means nothing to report"
);
}
#[tokio::test]
async fn calls_recorded_before_attempts_existed_count_as_one_request_each() {
let store = MemorySessionStore::default();
let meta = store
.create(CreateSession::default())
.await
.expect("create session");
store
.append(&meta.id, llm_call(100, 10, 0.01))
.await
.expect("append");
store
.append(&meta.id, llm_call_with_attempts(0.01, 3, 3))
.await
.expect("append");
let run = project_run_record_from_session(&store, &meta.id)
.await
.expect("project");
let attempts = run.metadata.get("provider_attempts").expect("attempts");
assert_eq!(
attempts.get("total").and_then(|v| v.as_i64()),
Some(4),
"1 (unrecorded, floored to one request) + 3 (recorded)"
);
assert_eq!(
attempts.get("rate_limited").and_then(|v| v.as_i64()),
Some(3)
);
}