use super::*;
use crate::{
agent::AgentOutputSink, output::OutputEvent, sessions::SessionEvent,
tui::state::MissionControlState,
};
use std::sync::{Arc, Mutex};
fn usage(input: u64, output: u64, cache: Option<u64>) -> NormalizedUsageSnapshot {
NormalizedUsageSnapshot {
effective_input: input,
output,
cache_read: cache.unwrap_or_default(),
cache_known: cache.is_some(),
}
}
fn record(id: &str, sequence: u64, input: u64, final_usage: bool) -> SessionUsageRecord {
SessionUsageRecord {
id: id.into(),
request_sequence: sequence,
usage: usage(input, input / 2, Some(input / 4)),
final_usage,
}
}
fn event(kind: &str, payload: serde_json::Value) -> SessionEvent {
SessionEvent::new(kind, "session".into(), "/tmp".into(), payload)
}
#[test]
fn cumulative_subagents_replace_snapshots_and_ignore_stale_delivery() {
let mut state = MissionControlState::default();
let first = record("session-usage/run/activity/call/g1", 1, 20, true);
let second = record(&first.id, 2, 60, true);
state.session_usage.observe(first.clone());
state.session_usage.observe(record(&first.id, 2, 40, false));
state.session_usage.observe(second.clone());
state.session_usage.observe(first);
state
.session_usage
.observe(record(&second.id, 2, 40, false));
state.session_usage.observe(second);
state
.session_usage
.observe(record("session-usage/run/primary/1", 1, 100, true));
state
.session_usage
.observe(record("session-usage/run/primary/2", 2, 200, true));
assert_eq!(state.session_usage_totals(), usage(360, 180, Some(90)));
let outcome = state.session_usage.clone();
state.session_usage.merge(&outcome);
state.nodes.clear();
state.transcript.clear();
state.apply_output_event(&OutputEvent::CompactionCompleted {
current_tokens: 10,
max_tokens: 100,
summary: "summary".into(),
});
assert_eq!(state.session_usage_totals(), usage(360, 180, Some(90)));
state.reset_for_new_session();
assert_eq!(
state.session_usage_totals(),
NormalizedUsageSnapshot::default()
);
}
#[test]
fn unknown_cache_survives_other_agents_and_final_can_resolve_partial_unknown() {
let mut ledger = SessionUsageLedger::default();
let mut unknown = record("primary", 1, 100, false);
unknown.usage.cache_known = false;
ledger.observe(unknown);
ledger.observe(record("child", 1, 20, true));
assert!(!ledger.totals().cache_known);
ledger.observe(record("primary", 1, 120, true));
assert_eq!(ledger.totals(), usage(140, 70, Some(35)));
ledger.observe(SessionUsageRecord {
id: "failed child".into(),
usage: usage(10, 2, None),
request_sequence: 1,
final_usage: false,
});
assert_eq!(
ledger.totals(),
NormalizedUsageSnapshot {
cache_known: false,
..usage(150, 72, Some(35))
}
);
}
#[test]
fn replacing_saturated_snapshot_does_not_lose_other_requests() {
let mut ledger = SessionUsageLedger::default();
ledger.observe(record("one", 1, u64::MAX, false));
ledger.observe(record("two", 1, 20, true));
assert_eq!(ledger.totals().effective_input, u64::MAX);
ledger = SessionUsageLedger::from_checkpoint(serde_json::to_value(&ledger).unwrap()).unwrap();
ledger.observe(record("one", 1, 10, true));
assert_eq!(ledger.totals().effective_input, 30);
}
#[test]
fn persisted_usage_restores_all_requests_across_compaction_and_new_run_sequences() {
let events = vec![
event("user_input", serde_json::json!({"text":"one"})),
event(
"session_usage",
serde_json::to_value(record("session-usage/run-a/primary/1", 1, 20, false)).unwrap(),
),
event(
"session_usage",
serde_json::to_value(record("session-usage/run-a/primary/1", 1, 40, true)).unwrap(),
),
event(
"session_usage",
serde_json::to_value(record("session-usage/run-a/activity/call/g1", 2, 60, true))
.unwrap(),
),
event(
"assistant_output",
serde_json::json!({"text":"done", "usage":{"input":40, "output":20}}),
),
event("compaction", serde_json::json!({"summary":"summary"})),
event("user_input", serde_json::json!({"text":"two"})),
event(
"session_usage",
serde_json::to_value(record("session-usage/run-b/primary/1", 1, 100, true)).unwrap(),
),
];
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
let mut state = MissionControlState::default();
state
.session_usage
.observe(record("old session", 1, 999, true));
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(state.session_usage_totals(), usage(200, 100, Some(50)));
assert!(!state.session_usage.incomplete);
}
#[test]
fn disconnected_event_delivery_retains_durable_usage_and_completion_mailbox() {
let temp = tempfile::tempdir().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let ledger = Arc::new(Mutex::new(SessionUsageLedger::default()));
let recorder = Arc::new(SessionUsageRecorder::new(
Some(session.clone()),
temp.path().into(),
Arc::clone(&ledger),
));
let (sender, receiver) = crossbeam_channel::bounded(1);
drop(receiver);
let mut sink =
crate::tui::events::TuiOutputSink::new(sender, Arc::new(|_| {}), "run-1/".into());
sink.usage_recorder = Some(recorder);
assert!(
sink.output_event(OutputEvent::UsageSnapshot {
usage: usage(100, 40, Some(25)),
request_sequence: 1,
final_usage: true
})
.is_err()
);
let mut state = MissionControlState::default();
state.session_usage.merge(&ledger.lock().unwrap());
state.session_usage.merge(&ledger.lock().unwrap());
assert_eq!(state.session_usage_totals(), usage(100, 40, Some(25)));
let events = session
.read_events_tolerant_bounded(100, 1024 * 1024)
.unwrap()
.events;
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(state.session_usage_totals(), usage(100, 40, Some(25)));
}
#[test]
fn legacy_usage_is_a_lower_bound_and_manual_compaction_does_not_hide_it() {
let events = vec![
event("user_input", serde_json::json!({"text":"legacy"})),
event(
"assistant_output",
serde_json::json!({"text":"done", "usage":{"input":20,"output":4,"cache_read":2}}),
),
event(
"session_usage",
serde_json::to_value(record("session-usage/new/compaction/1", 1, 100, true)).unwrap(),
),
];
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
let mut state = MissionControlState::default();
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(state.session_usage_totals().effective_input, 120);
assert_eq!(state.session_usage_totals().output, 54);
assert_eq!(state.session_usage_totals().cache_read, 27);
assert!(!state.session_usage_totals().cache_known);
assert!(state.session_usage.incomplete);
}
#[test]
fn unrecorded_terminal_usage_remains_unknown_after_resuming() {
for status in ["cancelled", "failed"] {
let events = vec![
event("user_input", serde_json::json!({"text":"old request"})),
event(
"turn_status",
serde_json::json!({"status":status,"assistant_text":"partial response"}),
),
event("user_input", serde_json::json!({"text":"new request"})),
event(
"session_usage",
serde_json::to_value(record("session-usage/new/primary/1", 1, 100, true)).unwrap(),
),
];
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
let mut state = MissionControlState::default();
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(state.session_usage_totals().effective_input, 100);
assert!(!state.session_usage_is_complete(), "{status}");
assert!(!state.session_usage_totals().cache_known, "{status}");
}
}
#[test]
fn real_rotation_reopen_preserves_prefix_usage_and_reconciles_suffix() {
let temp = tempfile::tempdir().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let ledger = Arc::new(Mutex::new(SessionUsageLedger::default()));
let recorder = SessionUsageRecorder::new(Some(session.clone()), temp.path().into(), ledger);
recorder
.record("primary/1", usage(100, 20, Some(50)), 1, true)
.unwrap();
recorder
.record("activity/call/g1", usage(20, 5, Some(10)), 1, false)
.unwrap();
recorder
.record("activity/call/g1", usage(40, 10, Some(20)), 1, true)
.unwrap();
crate::sessions::record_session_event(
Some(&session), temp.path(), crate::sessions::SessionEventKind::AssistantOutput,
serde_json::json!({"text": "done", "usage": {"input": 100, "output": 20, "cache_read": 50}}),
).unwrap();
crate::sessions::record_session_compaction(
&session,
temp.path(),
"summary",
"openai",
"test",
2,
)
.unwrap();
for rotation in 0..2 {
let reopened = manager.open(session.id()).unwrap();
let events = reopened
.read_events_tolerant_bounded(100, 1024 * 1024)
.unwrap()
.events;
assert_eq!(
events[0].kind(),
Some(crate::sessions::SessionEventKind::Compaction)
);
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
let mut state = MissionControlState::default();
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(
state.session_usage_totals(),
usage(140, 30, Some(70)),
"rotation {rotation}"
);
assert!(state.session_usage_is_complete());
if rotation == 0 {
crate::sessions::record_session_compaction(
&reopened,
temp.path(),
"summary two",
"openai",
"test",
events.len(),
)
.unwrap();
}
}
}
#[test]
fn missing_primary_and_subagent_usage_stays_incomplete_after_rotation() {
for subagent in [false, true] {
let temp = tempfile::tempdir().unwrap();
let manager = crate::sessions::SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let ledger = Arc::new(Mutex::new(SessionUsageLedger::default()));
let recorder = Arc::new(SessionUsageRecorder::new(
Some(session.clone()),
temp.path().into(),
Arc::clone(&ledger),
));
let (sender, _receiver) = crossbeam_channel::bounded(32);
let mut sink =
crate::tui::events::TuiOutputSink::new(sender, Arc::new(|_| {}), "run-1/".into());
sink.usage_recorder = Some(Arc::clone(&recorder));
if subagent {
sink.activity_event(crate::output::ActivityEvent::UsageUpdate {
id: crate::output::ActivityId::new("call/g1"),
current_tokens: 100,
max_tokens: 1000,
reasoning_tokens: None,
source: crate::output::ContextUsageSource::FallbackEstimate,
request_sequence: 1,
})
.unwrap();
recorder.request_started("activity/call/g1", 2).unwrap();
recorder
.record("activity/call/g1", usage(20, 5, Some(10)), 2, true)
.unwrap();
} else {
sink.output_event(OutputEvent::ContextUsage {
current_tokens: 100,
max_tokens: 1000,
reasoning_tokens: None,
source: crate::output::ContextUsageSource::FallbackEstimate,
request_sequence: 1,
})
.unwrap();
sink.output_event(OutputEvent::UsageSnapshot {
usage: usage(20, 5, Some(10)),
request_sequence: 2,
final_usage: true,
})
.unwrap();
}
let mut state = MissionControlState::default();
state.session_usage.merge(&ledger.lock().unwrap());
assert!(!state.session_usage_is_complete());
crate::sessions::record_session_compaction(
&session,
temp.path(),
"summary",
"openai",
"test",
100,
)
.unwrap();
let reopened = manager.open(session.id()).unwrap();
let events = reopened
.read_events_tolerant_bounded(100, 1024 * 1024)
.unwrap()
.events;
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert!(!state.session_usage_is_complete());
assert_eq!(state.session_usage_totals().effective_input, 20);
}
}
#[test]
fn completed_runs_fold_to_constant_size_without_losing_totals() {
let mut ledger = SessionUsageLedger::default();
for run in 0..1000 {
let mut worker = SessionUsageLedger::default();
worker.observe(record(
&format!("session-usage/run-{run}/primary/1"),
1,
100,
true,
));
ledger.merge(&worker);
ledger.merge(&worker);
ledger.fold_completed_run();
}
assert_eq!(ledger.totals(), usage(100_000, 50_000, Some(25_000)));
let encoded = serde_json::to_value(&ledger).unwrap();
assert!(encoded["records"].as_object().unwrap().is_empty());
let restored: SessionUsageLedger = serde_json::from_value(encoded).unwrap();
assert_eq!(restored.totals(), ledger.totals());
assert!(restored.is_complete());
}
#[test]
fn compaction_request_start_without_events_leaves_an_unresolved_request() {
let ledger = Arc::new(Mutex::new(SessionUsageLedger::default()));
let recorder = Arc::new(SessionUsageRecorder::new(
None,
"/tmp".into(),
Arc::clone(&ledger),
));
let (sender, _receiver) = crossbeam_channel::bounded(8);
let mut sink =
crate::tui::events::TuiOutputSink::new(sender, Arc::new(|_| {}), "run-1/".into());
sink.usage_recorder = Some(recorder);
sink.usage_source = "compaction";
let mut usage = crate::agent::provider_stream::CompactionUsage::default();
usage
.observe(
"openai",
crate::compaction::CompactionObservation::RequestStarted,
1,
&mut Some(&mut sink),
)
.unwrap();
let mut state = MissionControlState::default();
state.session_usage.merge(&ledger.lock().unwrap());
assert!(!state.session_usage_is_complete());
assert_eq!(state.session_usage_totals().effective_input, 0);
}
fn apply_usage_snapshot(state: &mut MissionControlState, record: SessionUsageRecord) {
state.apply_activity_event(crate::output::ActivityEvent::UsageSnapshot {
id: crate::output::ActivityId::new(record.id),
usage: crate::output::NormalizedUsageAggregate {
whole_run: record.usage,
..Default::default()
},
request_sequence: record.request_sequence,
final_usage: record.final_usage,
});
}
#[test]
fn session_cache_percent_keeps_last_complete_value_until_new_usage_is_known() {
let mut state = MissionControlState::default();
let first_id = "session-usage/run/primary/1";
let second_id = "session-usage/run/primary/2";
assert_eq!(state.session_cache_percent(), None);
apply_usage_snapshot(&mut state, record(first_id, 1, 100, false));
assert_eq!(state.session_cache_percent(), None);
apply_usage_snapshot(&mut state, record(first_id, 1, 100, true));
assert_eq!(state.session_cache_percent(), Some(25));
let mut second = record(second_id, 2, 100, false);
second.usage.cache_read = 75;
apply_usage_snapshot(&mut state, second.clone());
assert!(!state.session_usage_is_complete());
assert!(!state.session_usage_totals().cache_known);
assert_eq!(state.session_cache_percent(), Some(25));
second.final_usage = true;
apply_usage_snapshot(&mut state, second);
assert_eq!(state.session_cache_percent(), Some(50));
state.session_usage.fold_completed_run();
let mut unknown = record("session-usage/next/primary/1", 1, 100, true);
unknown.usage.cache_known = false;
apply_usage_snapshot(&mut state, unknown);
assert_eq!(state.session_cache_percent(), Some(50));
state.reset_for_new_session();
assert_eq!(state.session_cache_percent(), None);
let mut uncached = record(first_id, 1, 100, true);
uncached.usage.cache_read = 0;
apply_usage_snapshot(&mut state, uncached);
assert_eq!(state.session_cache_percent(), Some(0));
}
#[test]
fn session_cache_percent_hydration_recovers_last_known_value_and_isolates_sessions() {
let known = record("session-usage/run/primary/1", 1, 100, true);
let pending = record("session-usage/run/primary/2", 2, 100, false);
let mut state = MissionControlState::default();
let mut old = record("session-usage/old/primary/1", 1, 100, true);
old.usage.cache_read = 90;
apply_usage_snapshot(&mut state, old);
assert_eq!(state.session_cache_percent(), Some(90));
let events = [
event("session_usage", serde_json::to_value(known).unwrap()),
event("session_usage", serde_json::to_value(&pending).unwrap()),
];
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert!(!state.session_usage_is_complete());
assert_eq!(state.session_cache_percent(), Some(25));
for events in [
vec![event(
"session_usage",
serde_json::to_value(&pending).unwrap(),
)],
vec![],
] {
let snapshot = crate::tui::sessions::commands::hydrate_session_history_snapshot(
&events,
&Default::default(),
);
crate::tui::sessions::commands::apply_session_hydration_snapshot(&mut state, snapshot);
assert_eq!(state.session_cache_percent(), None);
}
}
#[test]
fn session_cache_percent_handles_zero_input_and_large_counters() {
let mut state = MissionControlState::default();
let id = "session-usage/run/primary/1";
apply_usage_snapshot(&mut state, record(id, 1, 0, true));
assert_eq!(state.session_cache_percent(), None);
let mut large = record(id, 2, u64::MAX, true);
large.usage.cache_read = u64::MAX;
apply_usage_snapshot(&mut state, large);
assert_eq!(state.session_cache_percent(), Some(100));
let mut excessive = record(id, 3, 1, true);
excessive.usage.cache_read = 2;
apply_usage_snapshot(&mut state, excessive);
assert_eq!(state.session_cache_percent(), Some(100));
}