From fde9631cca555f6471e58f770809c1077d65bbe0 Mon Sep 17 00:00:00 2001 From: Tryanks Date: Tue, 8 Sep 2026 19:05:04 +0800 Subject: [PATCH 1/4] fix(usage): distinguish request context from traffic and account eligibility --- crates/agent/src/acp.rs | 2 + crates/agent/src/claude.rs | 325 ++++++++++++++++-- crates/agent/src/codex.rs | 5 +- crates/agent/src/lib.rs | 67 +++- crates/agent/src/opencode.rs | 6 +- crates/agent/src/pi.rs | 4 +- crates/agent/tests/fixtures/claude/USAGE.md | 44 +++ .../fixtures/claude/compaction_recorded.jsonl | 2 + .../fixtures/claude/usage_recorded.jsonl | 11 + .../tests/fixtures/claude/usage_resume.jsonl | 6 + .../tests/fixtures/claude/usage_scope.jsonl | 10 + crates/core/src/relay.rs | 18 +- crates/core/src/session.rs | 58 +++- crates/core/src/settings.rs | 69 ++++ crates/runtime/src/app/events.rs | 2 +- crates/runtime/src/app/providers.rs | 93 ++++- crates/runtime/src/app/sessions.rs | 11 + crates/services/src/import/mod.rs | 2 +- crates/ui/src/chat/components/dividers.rs | 36 +- crates/ui/src/chat/mod.rs | 7 +- crates/ui/src/chat/model.rs | 6 +- crates/ui/src/composer/components/pickers.rs | 17 +- crates/ui/src/context_meter.rs | 12 +- crates/ui/src/settings_page.rs | 7 +- crates/ui/src/store/snapshots.rs | 210 +++++++++++ docs/DESIGN.md | 21 ++ locales/en.yml | 8 + locales/zh-CN.yml | 8 + 28 files changed, 1007 insertions(+), 60 deletions(-) create mode 100644 crates/agent/tests/fixtures/claude/USAGE.md create mode 100644 crates/agent/tests/fixtures/claude/compaction_recorded.jsonl create mode 100644 crates/agent/tests/fixtures/claude/usage_recorded.jsonl create mode 100644 crates/agent/tests/fixtures/claude/usage_resume.jsonl create mode 100644 crates/agent/tests/fixtures/claude/usage_scope.jsonl diff --git a/crates/agent/src/acp.rs b/crates/agent/src/acp.rs index e3110dca..81b00302 100644 --- a/crates/agent/src/acp.rs +++ b/crates/agent/src/acp.rs @@ -1303,6 +1303,7 @@ async fn finish_turn( let (status, message, usage) = match outcome.result { Ok(response) => { let usage = response.usage.as_ref().map(|usage| TokenUsage { + freshness: crate::ContextFreshness::Current, total_processed_tokens: Some(usage.total_tokens), input_tokens: Some(usage.input_tokens), cached_input_tokens: usage.cached_read_tokens, @@ -1817,6 +1818,7 @@ impl State { } acp::SessionUpdate::UsageUpdate(usage) => { let usage = TokenUsage { + freshness: crate::ContextFreshness::Current, used_tokens: Some(usage.used), context_window: Some(usage.size), cost_usd: usage diff --git a/crates/agent/src/claude.rs b/crates/agent/src/claude.rs index 68f05a8a..9b6bbac4 100644 --- a/crates/agent/src/claude.rs +++ b/crates/agent/src/claude.rs @@ -1185,10 +1185,11 @@ pub(crate) struct Mapper { exit_plan_captured: bool, /// Control responses to write back (e.g. the auto-deny for `ExitPlanMode`). outgoing: Vec, - /// Cumulative tokens processed across every completed turn this session - /// (Claude reports only per-turn usage, so we accumulate it ourselves for - /// the "Total processed" display). - cumulative_processed: u64, + request_usage: Value, + latest_usage: TokenUsage, + usage_message_id: Option, + usage_message_ids: HashSet, + result_ids: HashSet, /// Successfully written steers awaiting the CLI's next input checkpoint. pending_steers: VecDeque, /// Once the CLI exposes request checkpoints, never use the legacy fallback. @@ -1265,7 +1266,11 @@ impl Mapper { pending_permission_modes: HashMap::new(), exit_plan_captured: false, outgoing: Vec::new(), - cumulative_processed: 0, + request_usage: json!({}), + latest_usage: TokenUsage::default(), + usage_message_id: None, + usage_message_ids: HashSet::new(), + result_ids: HashSet::new(), pending_steers: VecDeque::new(), saw_requesting: false, background_tasks: HashSet::new(), @@ -1293,6 +1298,9 @@ impl Mapper { /// Allocate the next synthesized turn id and mark it in-flight. fn start_turn(&mut self) -> String { + if self.latest_usage.freshness == crate::ContextFreshness::Current { + self.latest_usage.freshness = crate::ContextFreshness::LastKnown; + } self.turn_counter += 1; let id = format!("turn-{}", self.turn_counter); self.current_turn_id = Some(id.clone()); @@ -1563,7 +1571,44 @@ impl Mapper { Some("init") => {} // Claude compacted its context window (verified shape: // `{type:"system", subtype:"compact_boundary", compact_metadata:{…}}`). - Some("compact_boundary") => return vec![AgentEvent::ContextCompacted], + Some("compact_boundary") => { + self.request_usage = json!({}); + self.usage_message_id = None; + self.current_message_id = None; + self.latest_usage.used_tokens = None; + self.latest_usage.input_tokens = None; + self.latest_usage.cached_input_tokens = None; + self.latest_usage.output_tokens = None; + self.latest_usage.freshness = crate::ContextFreshness::AwaitingObservation; + return vec![AgentEvent::ContextCompacted(crate::Compaction { + in_progress: false, + trigger: msg + .pointer("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/compact_metadata/trigger") + .and_then(Value::as_str) + .map(str::to_owned), + pre_tokens: msg + .pointer("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/compact_metadata/pre_tokens") + .and_then(Value::as_u64), + post_tokens: msg + .pointer("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/compact_metadata/post_tokens") + .and_then(Value::as_u64), + dropped_tokens: msg + .pointer("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/compact_metadata/cumulative_dropped_tokens") + .and_then(Value::as_u64), + duration_ms: msg + .pointer("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/compact_metadata/duration_ms") + .and_then(Value::as_u64), + })]; + } + Some("status") if msg.get("status").and_then(Value::as_str) == Some("compacting") => { + self.latest_usage.freshness = crate::ContextFreshness::Compacting; + return vec![AgentEvent::ContextCompacted(crate::Compaction { + in_progress: true, + trigger: None, + pre_tokens: None, + ..Default::default() + })]; + } Some("task_started") => return self.on_task_started(msg), Some("task_updated") => return self.on_task_updated(msg), Some("task_notification") => return self.on_task_notification(msg), @@ -1851,7 +1896,12 @@ impl Mapper { .and_then(|m| m.get("id")) .and_then(Value::as_str) .map(str::to_string); - Vec::new() + let message = &event["message"]; + self.observe_request_usage( + message.get("id").and_then(Value::as_str), + message.get("usage"), + true, + ) } Some("content_block_delta") => { let index = event.get("index").and_then(Value::as_u64).unwrap_or(0); @@ -1890,8 +1940,8 @@ impl Mapper { event.pointer("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/delta/stop_reason").and_then(Value::as_str), ); if let Some(usage) = event.get("usage") { - let tu = map_usage(usage, None); - events.push(AgentEvent::TokenUsage(tu)); + let id = self.current_message_id.clone(); + events.extend(self.observe_request_usage(id.as_deref(), Some(usage), false)); } events } @@ -1899,6 +1949,69 @@ impl Mapper { } } + /// Stream deltas are cumulative fields within one request, never increments. + /// Assistant messages can repeat an ID for parallel tool blocks. + fn observe_request_usage( + &mut self, + id: Option<&str>, + usage: Option<&Value>, + start: bool, + ) -> Vec { + let Some(id) = id else { + return Vec::new(); + }; + if self.usage_message_id.as_deref() != Some(id) { + if !self.usage_message_ids.insert(id.to_owned()) { + return Vec::new(); + } + self.usage_message_id = Some(id.to_owned()); + self.request_usage = json!({}); + } + if let Some(fields) = usage.and_then(Value::as_object) { + for key in [ + "input_tokens", + "cache_read_input_tokens", + "cache_creation_input_tokens", + "output_tokens", + ] { + if let Some(value) = fields.get(key).filter(|v| v.is_u64()) { + // Non-streaming assistant output is a placeholder; never regress + // the cumulative output already observed in the stream. + if key != "output_tokens" || value.as_u64() >= self.request_usage[key].as_u64() + { + self.request_usage[key] = value.clone(); + } + } + } + } + self.latest_usage.turn_processed_tokens = None; + let mapped = map_usage(&self.request_usage, None); + self.latest_usage.input_tokens = mapped.input_tokens; + self.latest_usage.cached_input_tokens = mapped.cached_input_tokens; + self.latest_usage.output_tokens = mapped.output_tokens; + // Claude's context observation counts the latest request's input and caches. + // Output deltas describe generated traffic, not a new input observation. + self.latest_usage.used_tokens = mapped.input_tokens.map(|input| { + input + .saturating_add(mapped.cached_input_tokens.unwrap_or(0)) + .saturating_add( + self.request_usage["cache_creation_input_tokens"] + .as_u64() + .unwrap_or(0), + ) + }); + self.latest_usage.freshness = if self.latest_usage.used_tokens.is_some() { + crate::ContextFreshness::Current + } else { + crate::ContextFreshness::AwaitingObservation + }; + if usage.is_some() || start { + vec![AgentEvent::TokenUsage(self.latest_usage)] + } else { + Vec::new() + } + } + fn block_item_id(&self, index: u64) -> String { match &self.current_message_id { Some(id) => format!("{id}:{index}"), @@ -1911,7 +2024,11 @@ impl Mapper { Some(m) => m, None => return Vec::new(), }; - let mut out = Vec::new(); + let mut out = self.observe_request_usage( + message.get("id").and_then(Value::as_str), + message.get("usage"), + false, + ); if message .pointer("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/stop_details/type") .and_then(Value::as_str) @@ -2492,6 +2609,11 @@ impl Mapper { } fn on_result(&mut self, msg: &Value) -> Vec { + if let Some(id) = msg.get("uuid").and_then(Value::as_str) + && !self.result_ids.insert(id.to_owned()) + { + return Vec::new(); + } self.awaiting_turn_checkpoint = false; let mut events = self.observe_stop_reason(msg.get("stop_reason").and_then(Value::as_str)); let turn_id = self @@ -2505,12 +2627,50 @@ impl Mapper { status = TurnStatus::Interrupted; } let usage = msg.get("usage").map(|u| { - let mut usage = map_usage(u, msg.get("modelUsage")); - // Accumulate this turn's processed tokens into the session total. - self.cumulative_processed += crate::processed_tokens(usage); - usage.total_processed_tokens = Some(self.cumulative_processed); + let aggregate = map_usage(u, msg.get("modelUsage")); + let processed = [ + "input_tokens", + "cache_read_input_tokens", + "cache_creation_input_tokens", + "output_tokens", + ] + .into_iter() + .any(|key| u.get(key).and_then(Value::as_u64).is_some()) + .then(|| crate::processed_tokens(aggregate)); + let mut usage = self.latest_usage; + // Select capacity for the main model; a child's larger window is irrelevant. + usage.context_window = self + .last_served_model + .as_deref() + .and_then(|model| { + msg.get("modelUsage")? + .as_object()? + .iter() + .find(|(key, value)| { + key.split('[').next() == Some(model) + || value.get("canonicalModel").and_then(Value::as_str) + == Some(model) + })? + .1 + .get("contextWindow")? + .as_u64() + }) + .or_else(|| { + msg.get("modelUsage")? + .as_object() + .filter(|m| m.len() == 1)? + .values() + .next()? + .get("contextWindow")? + .as_u64() + }) + .or(usage.context_window); + usage.turn_processed_tokens = processed; + // The timeline owns lifetime accumulation, including restored history. + usage.total_processed_tokens = None; usage.cost_usd = msg.get("total_cost_usd").and_then(Value::as_f64); usage.duration_ms = msg.get("duration_ms").and_then(Value::as_u64); + self.latest_usage = usage; usage }); if status == TurnStatus::Failed { @@ -4023,6 +4183,119 @@ mod tests { } } + #[test] + fn usage_scope_fixture_deduplicates_results_and_excludes_child_context() { + let mut m = Mapper::new(); + let mut observations = Vec::new(); + let mut completions = 0; + for line in include_str!("../tests/fixtures/claude/usage_scope.jsonl").lines() { + for event in feed(&mut m, line) { + match event { + AgentEvent::TokenUsage(u) => observations.push(u.used_tokens), + AgentEvent::TurnCompleted { usage: Some(u), .. } => { + completions += 1; + assert_eq!(u.used_tokens, Some(500260)); + assert_eq!(u.turn_processed_tokens, Some(4001300)); + } + _ => {} + } + } + } + assert_eq!(completions, 1); + assert_eq!( + observations, + vec![ + Some(400150), + Some(400150), + Some(400150), + Some(500260), + Some(500260) + ] + ); + } + + #[test] + fn recorded_compaction_retains_metadata_without_claiming_post_context() { + let mut m = Mapper::new(); + let mut events = Vec::new(); + for line in include_str!("../tests/fixtures/claude/compaction_recorded.jsonl").lines() { + events.extend(feed(&mut m, line)); + } + assert!(matches!(&events[0], AgentEvent::ContextCompacted(c) if c.in_progress)); + assert!( + matches!(&events[1], AgentEvent::ContextCompacted(c) if !c.in_progress && c.pre_tokens == Some(19555) && c.post_tokens == Some(2948) && c.dropped_tokens == Some(16607) && c.duration_ms == Some(24455)) + ); + assert_eq!(m.latest_usage.used_tokens, None); + assert_eq!( + m.latest_usage.freshness, + crate::ContextFreshness::AwaitingObservation + ); + } + + #[test] + fn recorded_multi_request_usage_keeps_latest_context() { + let mut m = Mapper::new(); + m.start_turn(); + let mut observations = Vec::new(); + let mut completion = None; + for line in include_str!("../tests/fixtures/claude/usage_recorded.jsonl").lines() { + for event in feed(&mut m, line) { + match event { + AgentEvent::TokenUsage(u) => observations.push(u.used_tokens), + AgentEvent::TurnCompleted { usage, .. } => completion = usage, + _ => {} + } + } + } + assert_eq!( + observations, + vec![ + Some(20250), + Some(20250), + Some(20250), + Some(20250), + Some(20635), + Some(20635), + Some(20635), + Some(20764), + Some(20764), + Some(20764) + ] + ); + let usage = completion.unwrap(); + assert_eq!(usage.used_tokens, Some(20764)); + assert_eq!(usage.output_tokens, Some(5)); + assert_eq!(usage.turn_processed_tokens, Some(62084)); + assert_eq!(usage.context_window, Some(1000000)); + } + + #[test] + fn streaming_usage_keeps_request_scope_at_completion() { + let mut m = Mapper::new(); + m.start_turn(); + let start = feed( + &mut m, + r#"{"type":"stream_event","event":{"type":"message_start","message":{"id":"usage-1","usage":{"input_tokens":100,"cache_read_input_tokens":400,"cache_creation_input_tokens":50,"output_tokens":0}}}}"#, + ); + assert!( + matches!(start.last(), Some(AgentEvent::TokenUsage(u)) if u.used_tokens == Some(550)), + "message start must expose the measured request context: {start:?}" + ); + let delta = feed( + &mut m, + r#"{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":20}}}"#, + ); + assert!( + matches!(delta.last(), Some(AgentEvent::TokenUsage(u)) if u.used_tokens == Some(550)), + "output-only delta must retain input and cache: {delta:?}" + ); + let result = feed( + &mut m, + r#"{"type":"result","subtype":"success","usage":{"input_tokens":1000,"cache_read_input_tokens":4000000,"output_tokens":200},"modelUsage":{"claude":{"contextWindow":1000000}}}"#, + ); + assert!(result.iter().any(|e| matches!(e, AgentEvent::TurnCompleted { usage: Some(u), .. } if u.used_tokens == Some(550) && u.turn_processed_tokens == Some(4001200))), "result accounting must not replace occupancy: {result:?}"); + } + #[test] fn compact_boundary_maps_to_context_compacted() { let mut m = Mapper::new(); @@ -4030,11 +4303,13 @@ mod tests { &mut m, r#"{"type":"system","subtype":"compact_boundary","session_id":"s1","compact_metadata":{"trigger":"manual","pre_tokens":500,"post_tokens":10}}"#, ); - assert!(matches!(evs.as_slice(), [AgentEvent::ContextCompacted])); + assert!( + matches!(evs.as_slice(), [AgentEvent::ContextCompacted(crate::Compaction { in_progress: false, trigger: Some(trigger), pre_tokens: Some(500), post_tokens: Some(10), .. })] if trigger == "manual") + ); } #[test] - fn result_accumulates_total_processed_tokens() { + fn result_reports_each_turn_traffic_for_timeline_accumulation() { let mut m = Mapper::new(); m.start_turn(); let evs = feed( @@ -4045,8 +4320,9 @@ mod tests { AgentEvent::TurnCompleted { usage, .. } => usage.unwrap(), other => panic!("expected TurnCompleted, got {other:?}"), }; - assert_eq!(first.total_processed_tokens, Some(120)); - // A second turn accumulates on top of the first. + assert_eq!(first.turn_processed_tokens, Some(120)); + assert_eq!(first.total_processed_tokens, None); + // Each result is a turn delta, independent of the mapper process lifetime. m.start_turn(); let evs = feed( &mut m, @@ -4056,7 +4332,8 @@ mod tests { AgentEvent::TurnCompleted { usage, .. } => usage.unwrap(), other => panic!("expected TurnCompleted, got {other:?}"), }; - assert_eq!(second.total_processed_tokens, Some(155)); + assert_eq!(second.turn_processed_tokens, Some(35)); + assert_eq!(second.total_processed_tokens, None); } #[test] @@ -4872,10 +5149,12 @@ mod tests { assert_eq!(tid, &turn_id); assert_eq!(*status, TurnStatus::Completed); let usage = usage.as_ref().expect("usage present"); - assert_eq!(usage.input_tokens, Some(100)); - assert_eq!(usage.cached_input_tokens, Some(50)); - assert_eq!(usage.output_tokens, Some(20)); - assert_eq!(usage.used_tokens, Some(180)); + assert_eq!(usage.input_tokens, None); + assert_eq!(usage.cached_input_tokens, None); + assert_eq!(usage.output_tokens, None); + assert_eq!(usage.used_tokens, None); + assert_eq!(usage.turn_processed_tokens, Some(180)); + assert_eq!(usage.total_processed_tokens, None); assert_eq!(usage.context_window, Some(1_000_000)); assert_eq!(usage.cost_usd, Some(0.125)); assert_eq!(usage.duration_ms, Some(4321)); diff --git a/crates/agent/src/codex.rs b/crates/agent/src/codex.rs index 42f3303f..e2daf223 100644 --- a/crates/agent/src/codex.rs +++ b/crates/agent/src/codex.rs @@ -1562,7 +1562,9 @@ impl Actor { } Some("contextCompaction") => { if method == "item/completed" { - self.events.emit(AgentEvent::ContextCompacted).await; + self.events + .emit(AgentEvent::ContextCompacted(Default::default())) + .await; } return; } @@ -2466,6 +2468,7 @@ fn map_usage(value: &Value) -> Option { // The session-cumulative running total lives in a sibling `total` object. let total_processed_tokens = value.pointer("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/total/totalTokens").and_then(Value::as_u64); Some(TokenUsage { + freshness: crate::ContextFreshness::Current, context_window: value.get("modelContextWindow").and_then(Value::as_u64), total_processed_tokens, input_tokens: last.get("inputTokens").and_then(Value::as_u64), diff --git a/crates/agent/src/lib.rs b/crates/agent/src/lib.rs index 9729d14e..e998dd49 100644 --- a/crates/agent/src/lib.rs +++ b/crates/agent/src/lib.rs @@ -1085,7 +1085,7 @@ pub enum AgentEvent { }, /// The provider compacted its context window (Claude `system/compact_boundary`; /// Codex `contextCompaction` item). Rendered as a "Context compacted" work-log row. - ContextCompacted, + ContextCompacted(Compaction), /// A tcode-level context-window change. This is never emitted by an adapter; /// the runtime persists it after the user message that selected the window. ContextWindowChanged { @@ -1451,9 +1451,36 @@ pub enum ApprovalDecision { Option(String), } +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ContextFreshness { + /// Older logs lack provenance; their counts may be turn aggregates. + #[default] + Unknown, + Current, + LastKnown, + Compacting, + AwaitingObservation, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(default)] +pub struct Compaction { + pub in_progress: bool, + pub trigger: Option, + pub pre_tokens: Option, + /// Provider-reported compaction output size; not a new request observation. + pub post_tokens: Option, + pub dropped_tokens: Option, + pub duration_ms: Option, +} + #[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)] #[serde(default)] pub struct TokenUsage { + pub freshness: ContextFreshness, + /// Main-loop traffic for the completed turn, separate from occupancy. + pub turn_processed_tokens: Option, pub input_tokens: Option, pub cached_input_tokens: Option, pub output_tokens: Option, @@ -1461,7 +1488,7 @@ pub struct TokenUsage { pub used_tokens: Option, pub context_window: Option, /// Cumulative tokens processed over the session's lifetime, if known (Codex - /// `thread/tokenUsage` running total; Claude accumulated per-turn usage). + /// `thread/tokenUsage` running total; timeline accumulation of Claude turns). /// Shown as "Total processed" in the context-meter popover. pub total_processed_tokens: Option, /// Provider-reported cost for this turn/message, in US dollars. @@ -1473,6 +1500,7 @@ pub struct TokenUsage { impl TokenUsage { #[cfg(feature = "process")] fn merge(&mut self, usage: Self) { + self.freshness = usage.freshness; self.input_tokens = add_token_counts(self.input_tokens, usage.input_tokens); self.cached_input_tokens = add_token_counts(self.cached_input_tokens, usage.cached_input_tokens); @@ -1868,3 +1896,38 @@ mod thread_item_serde_tests { ); } } + +#[cfg(test)] +mod usage_compatibility_tests { + use super::*; + + #[test] + fn compaction_and_usage_read_old_literal_records_without_inventing_provenance() { + let event: AgentEvent = serde_json::from_str(r#"{"type":"context_compacted"}"#).unwrap(); + assert!( + matches!(event, AgentEvent::ContextCompacted(c) if !c.in_progress && c.pre_tokens.is_none()) + ); + let usage: TokenUsage = serde_json::from_str( + r#"{"used_tokens":4100000,"context_window":1000000,"total_processed_tokens":4100000}"#, + ) + .unwrap(); + assert_eq!(usage.freshness, ContextFreshness::Unknown); + assert_eq!(usage.total_processed_tokens, Some(4100000)); + // Existing readers use a unit variant with this tag. Extra fields are + // additive and do not introduce an unknown protocol event tag. + #[derive(Deserialize)] + #[serde(tag = "type", rename_all = "snake_case")] + enum LegacyEvent { + ContextCompacted, + } + let _: LegacyEvent = serde_json::from_str(r#"{"type":"context_compacted","in_progress":false,"trigger":"manual","pre_tokens":500}"#).unwrap(); + let event = AgentEvent::ContextCompacted(Compaction { + trigger: Some("manual".into()), + pre_tokens: Some(500), + ..Default::default() + }); + let wire = serde_json::to_value(event).unwrap(); + assert_eq!(wire["type"], "context_compacted"); + assert_eq!(wire["pre_tokens"], 500); + } +} diff --git a/crates/agent/src/opencode.rs b/crates/agent/src/opencode.rs index 75b94c3d..4ca61120 100644 --- a/crates/agent/src/opencode.rs +++ b/crates/agent/src/opencode.rs @@ -725,7 +725,9 @@ impl OpenCodeMapper { .push((request_id.to_owned(), None)); } } - "session.compacted" => mapped.events.push(AgentEvent::ContextCompacted), + "session.compacted" => mapped + .events + .push(AgentEvent::ContextCompacted(Default::default())), _ => {} } mapped @@ -903,6 +905,7 @@ impl OpenCodeMapper { aggregate.merge(usage); } let aggregate = TokenUsage { + freshness: crate::ContextFreshness::Current, total_processed_tokens: Some(self.cumulative_processed), ..aggregate }; @@ -1123,6 +1126,7 @@ fn usage_from_tokens(tokens: Option<&Value>) -> Option { let cache_read = crate::json_u64(tokens.pointer("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/cache/read")); let cache_write = crate::json_u64(tokens.pointer("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/cache/write")).unwrap_or(0); (input.is_some() || output.is_some() || cache_read.is_some()).then_some(TokenUsage { + freshness: crate::ContextFreshness::Current, input_tokens: input, cached_input_tokens: cache_read, output_tokens: output, diff --git a/crates/agent/src/pi.rs b/crates/agent/src/pi.rs index 58f9b231..6d875450 100644 --- a/crates/agent/src/pi.rs +++ b/crates/agent/src/pi.rs @@ -866,7 +866,7 @@ impl PiMapper { }], "compaction_end" => { if message.get("result").is_some_and(|value| !value.is_null()) { - vec![AgentEvent::ContextCompacted] + vec![AgentEvent::ContextCompacted(Default::default())] } else { vec![AgentEvent::Warning { message: format!( @@ -1086,6 +1086,7 @@ impl PiMapper { let processed = crate::processed_tokens(usage); self.cumulative_processed = self.cumulative_processed.saturating_add(processed); let usage = TokenUsage { + freshness: crate::ContextFreshness::Current, total_processed_tokens: Some(self.cumulative_processed), ..usage }; @@ -1409,6 +1410,7 @@ fn map_usage(usage: Option<&Value>) -> Option { let output = crate::json_u64(usage.get("output")); let cache_read = crate::json_u64(usage.get("cacheRead")); (input.is_some() || output.is_some() || cache_read.is_some()).then_some(TokenUsage { + freshness: crate::ContextFreshness::Current, input_tokens: input, cached_input_tokens: cache_read, output_tokens: output, diff --git a/crates/agent/tests/fixtures/claude/USAGE.md b/crates/agent/tests/fixtures/claude/USAGE.md new file mode 100644 index 00000000..62a717f0 --- /dev/null +++ b/crates/agent/tests/fixtures/claude/USAGE.md @@ -0,0 +1,44 @@ +# Usage fixtures + +`usage_recorded.jsonl` and `compaction_recorded.jsonl` were captured on +2026-09-08 from new Claude Code 2.1.263 sessions in `/tmp/tcode-usage-provider-evidence`, +with `TCODE_DATA_DIR` pointing to a throwaway profile. No existing session was +opened or resumed. The probe used `--no-session-persistence`, `--verbose`, +`--output-format stream-json`, and `--include-partial-messages`. + +The first prompt asked for two separate Bash `printf` calls and a DONE response. +Only message IDs (replaced with stable fixture IDs), model names, usage and result +accounting fields were retained. Prompts, content, initialization payloads, +environment, stderr, account identifiers and credentials were not saved. +The independently calculated latest input/cache contexts are 20,250, 20,635 and +20,764. Turn traffic is 6 + 10,840 + 50,803 + 435 = 62,084. Repeated assistant +blocks share usage; their output values are placeholders, while stream deltas +are cumulative within each request. + +The second probe used `--input-format stream-json`, sent a new harmless printf +turn, then `/compact` in that same newly created process. The recorded lifecycle +is compacting → compact boundary, with manual trigger, 19,555 pre-tokens, 2,948 +post-tokens, 16,607 cumulative dropped tokens and 24,455 ms duration. Transcript +UUIDs in preserved-segment metadata were omitted. Post-tokens describe the +compaction result; they are not a fresh request context observation. + +`usage_scope.jsonl` and `usage_resume.jsonl` are constructed protocol fixtures, +not live captures. They protect >capacity turn traffic, output-only deltas, +repeated message/result IDs, main/subagent routing, a fresh process on an existing +timeline, compaction and interrupted completion. Their expected counts are +literal assertions in the production adapter → timeline → composer replay. + +Contract references: +- https://code.claude.com/docs/en/agent-sdk/cost-tracking +- https://platform.claude.com/docs/en/build-with-claude/streaming +- https://code.claude.com/docs/en/statusline#context-window-fields + +The exact reporter's 2.5M/4.1M run and automatic compaction were not captured. +The fixtures deliberately distinguish observed evidence from constructed cases. + +A separate read-only native account probe ran the existing ignored +`provider_usage::tests::live_provider_usage` test explicitly. Both Codex and +Claude returned non-empty windows with no error; only normalized usage output +was captured outside the repository. The unsupported endpoint and temporary +failure cases remain deterministic local fixtures, not claims about a live +custom account. diff --git a/crates/agent/tests/fixtures/claude/compaction_recorded.jsonl b/crates/agent/tests/fixtures/claude/compaction_recorded.jsonl new file mode 100644 index 00000000..e7705efb --- /dev/null +++ b/crates/agent/tests/fixtures/claude/compaction_recorded.jsonl @@ -0,0 +1,2 @@ +{"type":"system","subtype":"status","status":"compacting"} +{"type":"system","subtype":"compact_boundary","compact_metadata":{"trigger":"manual","pre_tokens":19555,"post_tokens":2948,"cumulative_dropped_tokens":16607,"duration_ms":24455}} diff --git a/crates/agent/tests/fixtures/claude/usage_recorded.jsonl b/crates/agent/tests/fixtures/claude/usage_recorded.jsonl new file mode 100644 index 00000000..ca0dd8d0 --- /dev/null +++ b/crates/agent/tests/fixtures/claude/usage_recorded.jsonl @@ -0,0 +1,11 @@ +{"type":"stream_event","event":{"type":"message_start","message":{"id":"recorded-request-1","usage":{"input_tokens":2,"cache_creation_input_tokens":10326,"cache_read_input_tokens":9922,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":10326},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"}}}} +{"type":"assistant","message":{"id":"recorded-request-1","usage":{"input_tokens":2,"cache_creation_input_tokens":10326,"cache_read_input_tokens":9922,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":10326},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"},"model":"claude-opus-5"}} +{"type":"assistant","message":{"id":"recorded-request-1","usage":{"input_tokens":2,"cache_creation_input_tokens":10326,"cache_read_input_tokens":9922,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":10326},"output_tokens":2,"service_tier":"standard","inference_geo":"not_available"},"model":"claude-opus-5"}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"input_tokens":2,"cache_creation_input_tokens":10326,"cache_read_input_tokens":9922,"output_tokens":343,"output_tokens_details":{"thinking_tokens":255},"iterations":[{"input_tokens":2,"output_tokens":343,"cache_read_input_tokens":9922,"cache_creation_input_tokens":10326,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":10326},"type":"message"}]}}} +{"type":"stream_event","event":{"type":"message_start","message":{"id":"recorded-request-2","usage":{"input_tokens":2,"cache_creation_input_tokens":385,"cache_read_input_tokens":20248,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":385},"output_tokens":17,"service_tier":"standard","inference_geo":"not_available"}}}} +{"type":"assistant","message":{"id":"recorded-request-2","usage":{"input_tokens":2,"cache_creation_input_tokens":385,"cache_read_input_tokens":20248,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":385},"output_tokens":17,"service_tier":"standard","inference_geo":"not_available"},"model":"claude-opus-5"}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"input_tokens":2,"cache_creation_input_tokens":385,"cache_read_input_tokens":20248,"output_tokens":87,"output_tokens_details":{"thinking_tokens":0},"iterations":[{"input_tokens":2,"output_tokens":87,"cache_read_input_tokens":20248,"cache_creation_input_tokens":385,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":385},"type":"message"}]}}} +{"type":"stream_event","event":{"type":"message_start","message":{"id":"recorded-request-3","usage":{"input_tokens":2,"cache_creation_input_tokens":129,"cache_read_input_tokens":20633,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":129},"output_tokens":1,"service_tier":"standard","inference_geo":"not_available"}}}} +{"type":"assistant","message":{"id":"recorded-request-3","usage":{"input_tokens":2,"cache_creation_input_tokens":129,"cache_read_input_tokens":20633,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":129},"output_tokens":1,"service_tier":"standard","inference_geo":"not_available"},"model":"claude-opus-5"}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"input_tokens":2,"cache_creation_input_tokens":129,"cache_read_input_tokens":20633,"output_tokens":5,"output_tokens_details":{"thinking_tokens":0},"iterations":[{"input_tokens":2,"output_tokens":5,"cache_read_input_tokens":20633,"cache_creation_input_tokens":129,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":129},"type":"message"}]}}} +{"type":"result","subtype":"success","usage":{"input_tokens":6,"cache_creation_input_tokens":10840,"cache_read_input_tokens":50803,"output_tokens":435,"output_tokens_details":{"thinking_tokens":255},"server_tool_use":{"web_search_requests":0,"web_fetch_requests":0},"service_tier":"standard","cache_creation":{"ephemeral_1h_input_tokens":10840,"ephemeral_5m_input_tokens":0},"inference_geo":"not_available","iterations":[{"input_tokens":2,"output_tokens":5,"cache_read_input_tokens":20633,"cache_creation_input_tokens":129,"cache_creation":{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":129},"type":"message"}],"speed":"standard"},"modelUsage":{"claude-opus-5[1m]":{"inputTokens":6,"outputTokens":435,"cacheReadInputTokens":50803,"cacheCreationInputTokens":10840,"webSearchRequests":0,"costUSD":0.1447065,"contextWindow":1000000,"maxOutputTokens":64000,"thinkingTokens":255,"canonicalModel":"claude-opus-5","provider":"firstParty","costBasis":"list"}},"duration_ms":10822} diff --git a/crates/agent/tests/fixtures/claude/usage_resume.jsonl b/crates/agent/tests/fixtures/claude/usage_resume.jsonl new file mode 100644 index 00000000..3efac15c --- /dev/null +++ b/crates/agent/tests/fixtures/claude/usage_resume.jsonl @@ -0,0 +1,6 @@ +{"type":"system","subtype":"init","session_id":"fixture-session","model":"claude-opus-5"} +{"type":"system","subtype":"status","status":"compacting"} +{"type":"system","subtype":"compact_boundary","compact_metadata":{"trigger":"manual","pre_tokens":500260}} +{"type":"stream_event","event":{"type":"message_start","message":{"id":"request-three","usage":{"input_tokens":20,"cache_read_input_tokens":30,"cache_creation_input_tokens":10,"output_tokens":0}}}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":5}}} +{"type":"result","uuid":"result-two","subtype":"error_during_execution","is_error":true,"result":"Interrupted test turn","usage":{"input_tokens":20,"cache_read_input_tokens":30,"cache_creation_input_tokens":10,"output_tokens":5}} diff --git a/crates/agent/tests/fixtures/claude/usage_scope.jsonl b/crates/agent/tests/fixtures/claude/usage_scope.jsonl new file mode 100644 index 00000000..3037f1ee --- /dev/null +++ b/crates/agent/tests/fixtures/claude/usage_scope.jsonl @@ -0,0 +1,10 @@ +{"type":"system","subtype":"init","session_id":"fixture-session","model":"claude-opus-5"} +{"type":"stream_event","event":{"type":"message_start","message":{"id":"request-one","usage":{"input_tokens":100,"cache_read_input_tokens":400000,"cache_creation_input_tokens":50,"output_tokens":0}}}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":20}}} +{"type":"stream_event","parent_tool_use_id":"child-task","event":{"type":"message_start","message":{"id":"child-request","usage":{"input_tokens":900000,"output_tokens":500}}}} +{"type":"assistant","message":{"id":"request-one","content":[],"usage":{"input_tokens":100,"cache_read_input_tokens":400000,"cache_creation_input_tokens":50,"output_tokens":1}}} +{"type":"stream_event","event":{"type":"message_start","message":{"id":"request-two","usage":{"input_tokens":200,"cache_read_input_tokens":500000,"cache_creation_input_tokens":60,"output_tokens":0}}}} +{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":30}}} +{"type":"assistant","message":{"id":"request-one","content":[],"usage":{"input_tokens":100,"cache_read_input_tokens":400000,"cache_creation_input_tokens":50,"output_tokens":1}}} +{"type":"result","uuid":"result-one","subtype":"success","usage":{"input_tokens":1000,"cache_read_input_tokens":4000000,"cache_creation_input_tokens":100,"output_tokens":200},"modelUsage":{"claude-opus-5":{"contextWindow":1000000}}} +{"type":"result","uuid":"result-one","subtype":"success","usage":{"input_tokens":1000,"cache_read_input_tokens":4000000,"cache_creation_input_tokens":100,"output_tokens":200},"modelUsage":{"claude-opus-5":{"contextWindow":1000000}}} diff --git a/crates/core/src/relay.rs b/crates/core/src/relay.rs index 4bb967b9..c0edd2f9 100644 --- a/crates/core/src/relay.rs +++ b/crates/core/src/relay.rs @@ -233,9 +233,21 @@ fn render_turn(number: usize, entries: &[&TimelineEntry], timeline: &Timeline) - to_provider.display_name() )); } - EntryContent::ContextCompacted => { - activity(&mut body, "context", "provider", "compacted") - } + EntryContent::ContextCompacted(c) => activity( + &mut body, + "context", + "provider", + &format!( + "{}; trigger={:?}; pre_tokens={:?}", + if c.in_progress { + "compacting" + } else { + "compacted" + }, + c.trigger, + c.pre_tokens + ), + ), EntryContent::ContextWindowChanged { window } => activity( &mut body, "context", diff --git a/crates/core/src/session.rs b/crates/core/src/session.rs index 049a115c..6bddeae0 100644 --- a/crates/core/src/session.rs +++ b/crates/core/src/session.rs @@ -433,7 +433,7 @@ pub enum EntryContent { reason: Option, }, /// The provider compacted its context window (a "Context compacted" work-log row). - ContextCompacted, + ContextCompacted(agent::Compaction), /// The user changed the context window for the next provider turn. ContextWindowChanged { window: u64, @@ -616,6 +616,11 @@ impl Timeline { } } AgentEvent::TurnStarted { turn_id } => { + if let Some(usage) = self.usage.as_mut() + && usage.freshness == agent::ContextFreshness::Current + { + usage.freshness = agent::ContextFreshness::LastKnown; + } // Reuse the open turn (typically opened by the user message); // otherwise begin a fresh one. let turn = match self.current_turn { @@ -681,6 +686,24 @@ impl Timeline { } AgentEvent::RewindFailed { .. } => {} AgentEvent::TurnCompleted { status, usage, .. } => { + let newly_completed = self + .current_turn + .is_none_or(|turn| self.turns[turn].status.is_none()); + let usage = usage.map(|mut usage| { + if let Some(processed) = usage.turn_processed_tokens { + let previous = self + .usage + .and_then(|u| u.total_processed_tokens) + .unwrap_or(0); + usage.total_processed_tokens = + Some(previous.saturating_add(if newly_completed { + processed + } else { + 0 + })); + } + usage + }); self.turn_running = false; self.last_turn_status = Some(*status); if let Some(turn) = self.current_turn { @@ -701,7 +724,7 @@ impl Timeline { } } if usage.is_some() { - self.usage = *usage; + self.usage = usage; } // A finished turn can no longer be waiting on approvals or input. self.pending_approvals.clear(); @@ -766,7 +789,19 @@ impl Timeline { self.pending_user_input = None; } } - AgentEvent::TokenUsage(usage) => self.usage = Some(*usage), + AgentEvent::TokenUsage(usage) => { + let mut usage = *usage; + usage.context_window = usage + .context_window + .or(self.usage.and_then(|u| u.context_window)); + if usage.turn_processed_tokens.is_some() || usage.total_processed_tokens.is_none() { + usage.total_processed_tokens = self + .usage + .and_then(|u| u.total_processed_tokens) + .or(usage.total_processed_tokens); + } + self.usage = Some(usage); + } AgentEvent::Warning { message } => log::warn!("provider warning: {message}"), AgentEvent::ProviderStartFailed { error } => { let turn = self.ensure_turn(ts); @@ -880,12 +915,25 @@ impl Timeline { turn, }); } - AgentEvent::ContextCompacted => { + AgentEvent::ContextCompacted(compaction) => { + let in_progress = compaction.in_progress; + let usage = self.usage.get_or_insert_with(Default::default); + usage.freshness = if in_progress { + agent::ContextFreshness::Compacting + } else { + agent::ContextFreshness::AwaitingObservation + }; + if !in_progress { + usage.used_tokens = None; + usage.input_tokens = None; + usage.cached_input_tokens = None; + usage.output_tokens = None; + } let turn = self.ensure_turn(ts); let id = self.synthetic_id("compacted"); self.entries.push(Arc::new(TimelineEntry { id, - content: EntryContent::ContextCompacted, + content: EntryContent::ContextCompacted(compaction.clone()), ts, turn, })); diff --git a/crates/core/src/settings.rs b/crates/core/src/settings.rs index a5325f49..0ce7e2da 100644 --- a/crates/core/src/settings.rs +++ b/crates/core/src/settings.rs @@ -227,6 +227,53 @@ pub struct ResolvedProfile { pub settings: ProviderSettings, } +impl ResolvedProfile { + /// Whether this endpoint/auth configuration can query native account limits. + pub fn supports_account_usage(&self) -> bool { + let (endpoint_key, native_endpoint, credentials, backends): (&str, &str, &[&str], &[&str]) = + match self.kind { + ProviderKind::ClaudeCode => ( + "ANTHROPIC_BASE_URL", + "https://api.anthropic.com", + &["ANTHROPIC_API_KEY", "ANTHROPIC_AUTH_TOKEN"], + &[ + "CLAUDE_CODE_USE_BEDROCK", + "CLAUDE_CODE_USE_VERTEX", + "CLAUDE_CODE_USE_FOUNDRY", + ], + ), + ProviderKind::Codex => ( + "OPENAI_BASE_URL", + "https://api.openai.com/v1", + &["OPENAI_API_KEY", "CODEX_API_KEY"], + &[], + ), + _ => return false, + }; + // Secret presence is enough to classify API-key auth; never inspect or + // replicate secret values. Missing native sign-in remains a probe error. + !self.settings.env.iter().enumerate().any(|(index, env)| { + // LaunchEnv uses the last occurrence of a repeated environment key. + if self.settings.env[index + 1..] + .iter() + .any(|later| later.name == env.name) + { + return false; + } + let configured = env.sensitive || !env.value.trim().is_empty(); + configured + && (credentials.contains(&env.name.as_str()) + || (backends.contains(&env.name.as_str()) + && env.value != "0" + && env.value != "false") + || (env.name == endpoint_key + && (env.sensitive + || env.value.trim().trim_end_matches('/').to_ascii_lowercase() + != native_endpoint))) + }) + } +} + /// One configured model, unique by provider and model ID within its role. /// Reasoning effort is selected per tool call from the provider's capabilities. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -1965,3 +2012,25 @@ mod tests { assert_eq!(back.get("another"), Some(&serde_json::json!("value"))); } } + +#[cfg(test)] +mod account_usage_tests { + use super::*; + + #[test] + fn resolved_custom_endpoint_is_not_a_native_account() { + let settings: Settings = serde_json::from_str(r#"{"profiles":{"custom":{"kind":"claude_code","display_name":"Kimi","env":[{"name":"ANTHROPIC_BASE_URL","value":"https://api.example.com/anthropic"},{"name":"ANTHROPIC_API_KEY","sensitive":true}]}}}"#).unwrap(); + let mut custom = settings.resolved_profile("custom").unwrap(); + assert!( + !custom.supports_account_usage(), + "custom protocol compatibility does not imply native subscription support" + ); + custom.settings.display_name = Some("Claude".into()); + assert!(!custom.supports_account_usage()); + custom.settings.env.clear(); + assert!( + custom.supports_account_usage(), + "a custom profile using native account auth remains eligible, even signed out" + ); + } +} diff --git a/crates/runtime/src/app/events.rs b/crates/runtime/src/app/events.rs index 5c296e51..67c6c5e7 100644 --- a/crates/runtime/src/app/events.rs +++ b/crates/runtime/src/app/events.rs @@ -408,7 +408,7 @@ impl AppState { | AgentEvent::UserInputRequested { .. } | AgentEvent::UserInputResolved { .. } | AgentEvent::TokenUsage(_) - | AgentEvent::ContextCompacted + | AgentEvent::ContextCompacted(_) | AgentEvent::ContextWindowChanged { .. } | AgentEvent::PlanUpdated { .. } | AgentEvent::ProposedPlanDelta { .. } diff --git a/crates/runtime/src/app/providers.rs b/crates/runtime/src/app/providers.rs index 5aa73840..521e2399 100644 --- a/crates/runtime/src/app/providers.rs +++ b/crates/runtime/src/app/providers.rs @@ -10,6 +10,7 @@ pub struct ProviderCatalog { pub provider_usage: HashMap, /// Profile ids with a usage fetch in flight. pub usage_checking: HashSet, + usage_revisions: HashMap, pub(super) provider_secret_names: HashMap>, } @@ -26,10 +27,32 @@ impl ProviderCatalog { provider_snapshots: HashMap::new(), provider_usage: HashMap::new(), usage_checking: HashSet::new(), + usage_revisions: HashMap::new(), provider_secret_names, } } + pub(super) fn invalidate_usage(&mut self, id: &str) { + self.provider_usage.remove(id); + self.usage_checking.remove(id); + *self.usage_revisions.entry(id.to_owned()).or_default() += 1; + } + + fn complete_usage( + &mut self, + id: String, + revision: u64, + usage: Option, + ) { + if self.usage_revisions.get(&id).copied().unwrap_or(0) != revision { + return; + } + self.usage_checking.remove(&id); + if let Some(usage) = usage { + self.provider_usage.insert(id, usage); + } + } + pub(super) fn status_snapshot( &self, acp_marketplace_items: Vec, @@ -212,6 +235,7 @@ impl AppState { value: Option<&str>, cx: &mut HostCx, ) { + self.providers.invalidate_usage(id); self.enqueue_store_write( StoreWrite::SetProfileSecret { profile_id: id.to_string(), @@ -399,9 +423,10 @@ impl AppState { fn refresh_provider_usage_inner(&mut self, cx: &mut HostCx, stale_only: bool) { let now = now_secs(); for profile in self.all_profiles() { - let profile_id = profile.id; + let supported = profile.supports_account_usage(); + let profile_id = profile.id.clone(); let provider = profile.kind; - if !matches!(provider, ProviderKind::Codex | ProviderKind::ClaudeCode) + if !supported || self .providers .provider_snapshots @@ -418,6 +443,12 @@ impl AppState { continue; } self.providers.usage_checking.insert(profile_id.clone()); + let revision = self + .providers + .usage_revisions + .get(&profile_id) + .copied() + .unwrap_or(0); let binary = self.resolve_profile_binary(&profile_id); let settings = self.settings.clone(); let settings_store = self.settings_store.clone(); @@ -435,9 +466,8 @@ impl AppState { ) .await; host_cx.enqueue(move |state, _cx| { - state.providers.usage_checking.remove(&profile_id); - if let Some(usage) = usage { - state.providers.provider_usage.insert(profile_id, usage); + if state.settings.resolved_profile(&profile_id).as_ref() == Some(&profile) { + state.providers.complete_usage(profile_id, revision, usage); } }); }); @@ -831,3 +861,56 @@ pub(super) fn session_options( }), } } + +#[cfg(test)] +mod usage_lifecycle_tests { + use super::*; + + #[test] + fn configuration_invalidation_rejects_old_probe_without_finishing_new_probe() { + let mut catalog = ProviderCatalog::new(HashMap::new(), HashMap::new()); + let error = tcode_core::usage::ProviderUsage { + error: Some("temporarily unreachable".into()), + ..Default::default() + }; + catalog.usage_checking.insert("claude".into()); + catalog.complete_usage("claude".into(), 0, Some(error.clone())); + assert_eq!( + catalog.provider_usage["claude"], error, + "supported failures stay visible" + ); + // Exercise secret replacement through its production entry point while + // a response with the previous revision is still outstanding. + use crate::app::test_support::*; + let cx = &mut TestAppContext::default(); + let test_store = TestStore::new("tcode-usage-secret-race"); + let state = cx.new_entity(TestClientState::new((*test_store).clone())); + state.update(cx, |state, cx| { + std::mem::swap(&mut state.providers, &mut catalog); + state.set_profile_secret( + "claude", + "CLAUDE_CODE_OAUTH_TOKEN", + Some("fixture-replacement"), + cx, + ); + std::mem::swap(&mut state.providers, &mut catalog); + }); + assert!(!catalog.provider_usage.contains_key("claude")); + catalog.usage_checking.insert("claude".into()); + catalog.complete_usage("claude".into(), 0, Some(error)); + assert!(catalog.usage_checking.contains("claude")); + assert!(!catalog.provider_usage.contains_key("claude")); + let windows = tcode_core::usage::ProviderUsage { + windows: vec![tcode_core::usage::UsageWindow { + kind: tcode_core::usage::UsageWindowKind::FiveHour, + scope: None, + used_percent: 42., + resets_at: Some(1800000000), + }], + ..Default::default() + }; + catalog.complete_usage("claude".into(), 1, Some(windows.clone())); + assert_eq!(catalog.provider_usage["claude"], windows); + assert!(!catalog.usage_checking.contains("claude")); + } +} diff --git a/crates/runtime/src/app/sessions.rs b/crates/runtime/src/app/sessions.rs index 4471f410..3b004bf8 100644 --- a/crates/runtime/src/app/sessions.rs +++ b/crates/runtime/src/app/sessions.rs @@ -527,6 +527,17 @@ impl AppState { pub fn update_settings(&mut self, settings: Settings, cx: &mut HostCx) { self.enqueue_settings(&settings, cx); let language = settings.language.clone(); + let changed: HashSet<_> = self + .providers + .provider_usage + .keys() + .chain(self.providers.usage_checking.iter()) + .filter(|id| self.settings.resolved_profile(id) != settings.resolved_profile(id)) + .cloned() + .collect(); + for id in changed { + self.providers.invalidate_usage(&id); + } self.settings = settings; self.providers.provider_secret_names = provider_secret_names(&self.settings, &self.settings_store); diff --git a/crates/services/src/import/mod.rs b/crates/services/src/import/mod.rs index d13a097e..820f4610 100644 --- a/crates/services/src/import/mod.rs +++ b/crates/services/src/import/mod.rs @@ -218,7 +218,7 @@ pub fn import_thread( } } ConvertedEntry::Compacted { ts } => { - events.push((ts, AgentEvent::ContextCompacted)); + events.push((ts, AgentEvent::ContextCompacted(Default::default()))); } } } diff --git a/crates/ui/src/chat/components/dividers.rs b/crates/ui/src/chat/components/dividers.rs index 8020dd33..b9f16317 100644 --- a/crates/ui/src/chat/components/dividers.rs +++ b/crates/ui/src/chat/components/dividers.rs @@ -30,7 +30,9 @@ fn divider(id: SharedString, label: String, tint: Hsla, cx: &App) -> AnyElement .child(stub()) .child( div() - .flex_none() + .min_w_0() + .flex_shrink_1() + .text_center() .text_size(px(11.)) .text_color(tint) .child(label), @@ -85,10 +87,38 @@ pub(crate) fn model_change_divider( /// Context compaction rewrites what the model remembers, so it announces /// itself at model-swap prominence rather than hiding in the work log. -pub(crate) fn context_compacted_divider(id: &str, cx: &App) -> AnyElement { +pub(crate) fn context_compacted_divider( + id: &str, + metadata: Option<&agent::Compaction>, + cx: &App, +) -> AnyElement { + let mut label = if metadata.is_some_and(|c| c.in_progress) { + crate::tr!("chat.context_compacting").into_owned() + } else { + crate::tr!("chat.context_compacted").into_owned() + }; + if let Some(c) = metadata { + if let Some(trigger) = &c.trigger { + let trigger = match trigger.as_str() { + "manual" => crate::tr!("chat.context_manual").into_owned(), + "auto" => crate::tr!("chat.context_auto").into_owned(), + other => other.to_owned(), + }; + label.push_str(&format!(" · {trigger}")); + } + if let Some(tokens) = c.pre_tokens { + label.push_str(&format!( + " · {}", + crate::tr!( + "chat.context_pre_tokens", + tokens = crate::context_meter::format_tokens(Some(tokens)) + ) + )); + } + } divider( SharedString::from(format!("context-compacted-{id}")), - crate::tr!("chat.context_compacted").into_owned(), + label, cx.theme().warning, cx, ) diff --git a/crates/ui/src/chat/mod.rs b/crates/ui/src/chat/mod.rs index 4658fa6f..6409d185 100644 --- a/crates/ui/src/chat/mod.rs +++ b/crates/ui/src/chat/mod.rs @@ -1015,7 +1015,12 @@ impl ChatView { } Segment::ContextCompacted(entry) => { column = column.child(components::dividers::context_compacted_divider( - &entry.id, cx, + &entry.id, + match &entry.content { + EntryContent::ContextCompacted(c) => Some(c), + _ => None, + }, + cx, )); } Segment::ContextWindowChanged(entry) => { diff --git a/crates/ui/src/chat/model.rs b/crates/ui/src/chat/model.rs index 71f346d0..fb0b7711 100644 --- a/crates/ui/src/chat/model.rs +++ b/crates/ui/src/chat/model.rs @@ -159,7 +159,7 @@ pub(crate) fn segment_entries<'a>( flush_activities(&mut segments, &mut activities); segments.push(Segment::ModelChange(entry)); } - EntryContent::ContextCompacted => { + EntryContent::ContextCompacted(_) => { flush_activities(&mut segments, &mut activities); segments.push(Segment::ContextCompacted(entry)); } @@ -219,7 +219,7 @@ pub(crate) fn work_log_counts(entries: &[&TimelineEntry]) -> WorkLogCounts { | EntryContent::Item(ItemContent::WebSearch { .. }) | EntryContent::Item(ItemContent::Other { .. }) => counts.tools += 1, EntryContent::Item(ItemContent::Subagent { .. }) => counts.subagents += 1, - EntryContent::ContextCompacted + EntryContent::ContextCompacted(_) | EntryContent::ContextWindowChanged { .. } | EntryContent::Steer { .. } | EntryContent::Item(ItemContent::UserMessage { .. }) @@ -1116,7 +1116,7 @@ fn hash_entry_shape(content: &EntryContent, hash: &mut DefaultHasher) { to.hash(hash); reason.hash(hash); } - EntryContent::ContextCompacted => {} + EntryContent::ContextCompacted(_) => {} EntryContent::ContextWindowChanged { window } => window.hash(hash), EntryContent::Item(ItemContent::WebSearch { query }) => { "web_search".len().hash(hash); diff --git a/crates/ui/src/composer/components/pickers.rs b/crates/ui/src/composer/components/pickers.rs index 5e59f065..3de82105 100644 --- a/crates/ui/src/composer/components/pickers.rs +++ b/crates/ui/src/composer/components/pickers.rs @@ -1582,8 +1582,8 @@ fn render_context_meter_pane( .child(stat), ); - if max.is_some() { - let fraction = pct.unwrap_or(0.0).clamp(0.0, 100.0) / 100.0; + if let Some(pct) = pct { + let fraction = pct.clamp(0.0, 100.0) / 100.0; pane = pane.child( div() .w_full() @@ -1600,6 +1600,19 @@ fn render_context_meter_pane( ); } + let freshness = match usage.map(|u| u.freshness) { + Some(agent::ContextFreshness::Current) if used.is_some() => "composer.context_latest", + Some(agent::ContextFreshness::LastKnown) => "composer.context_last_known", + Some(agent::ContextFreshness::Compacting) => "chat.context_compacting", + _ => "composer.context_updating", + }; + pane = pane.child( + div() + .text_size(px(11.)) + .text_color(muted) + .child(crate::tr!(freshness)), + ); + // "Total processed" — the session-cumulative token count, when the provider // reports it (a native running total or adapter-side accumulation). if let Some(total) = usage.and_then(|u| u.total_processed_tokens) { diff --git a/crates/ui/src/context_meter.rs b/crates/ui/src/context_meter.rs index 908cfd6c..c45f6491 100644 --- a/crates/ui/src/context_meter.rs +++ b/crates/ui/src/context_meter.rs @@ -5,6 +5,12 @@ use agent::TokenUsage; /// The used tokens a meter reflects: the provider's reported total-in-use, else /// the input-token count. pub fn used_tokens(usage: &TokenUsage) -> Option { + if matches!( + usage.freshness, + agent::ContextFreshness::Unknown | agent::ContextFreshness::AwaitingObservation + ) { + return None; + } usage.used_tokens.or(usage.input_tokens) } @@ -44,7 +50,7 @@ pub fn format_percentage(percentage: Option) -> Option { /// `<1_000_000` as `Nk`, otherwise `x.ym`; trailing `.0` is omitted. pub fn format_tokens(value: Option) -> String { let Some(v) = value else { - return "0".to_string(); + return crate::tr!("composer.context_unknown").into_owned(); }; let v = v as f64; if v < 1_000.0 { @@ -66,6 +72,7 @@ mod tests { fn usage(used: Option, window: Option) -> TokenUsage { TokenUsage { + freshness: agent::ContextFreshness::Current, used_tokens: used, context_window: window, ..Default::default() @@ -90,6 +97,7 @@ mod tests { #[test] fn used_falls_back_to_input_tokens() { let u = TokenUsage { + freshness: agent::ContextFreshness::Current, input_tokens: Some(1_234), ..Default::default() }; @@ -115,6 +123,6 @@ mod tests { assert_eq!(format_tokens(Some(200_000)), "200k"); assert_eq!(format_tokens(Some(1_000_000)), "1m"); assert_eq!(format_tokens(Some(1_250_000)), "1.2m"); - assert_eq!(format_tokens(None), "0"); + assert_eq!(format_tokens(None), crate::tr!("composer.context_unknown")); } } diff --git a/crates/ui/src/settings_page.rs b/crates/ui/src/settings_page.rs index 68799000..530175d7 100644 --- a/crates/ui/src/settings_page.rs +++ b/crates/ui/src/settings_page.rs @@ -1638,12 +1638,7 @@ impl SettingsPage { let profiles: Vec<_> = store .enabled_profiles() .into_iter() - .filter(|profile| { - matches!( - profile.kind, - agent::ProviderKind::Codex | agent::ProviderKind::ClaudeCode - ) - }) + .filter(|profile| profile.supports_account_usage()) .collect(); let rows: Vec<(String, String, agent::ProviderKind, _, bool)> = profiles .iter() diff --git a/crates/ui/src/store/snapshots.rs b/crates/ui/src/store/snapshots.rs index a0207294..7a6e79a1 100644 --- a/crates/ui/src/store/snapshots.rs +++ b/crates/ui/src/store/snapshots.rs @@ -158,6 +158,9 @@ pub(crate) fn composer_state( .requested_profile_id .clone() .unwrap_or_else(|| Settings::builtin_profile_id(status.provider).to_owned()); + settings + .resolved_profile(&profile_id) + .filter(|profile| profile.supports_account_usage())?; providers.provider_usage.get(&profile_id).cloned() }); @@ -341,4 +344,211 @@ mod tests { Some(500_000) ); } + #[test] + fn account_usage_eligibility_does_not_hide_session_context_or_supported_errors() { + let mut settings: Settings = serde_json::from_str(r#"{"profiles":{"custom":{"kind":"claude_code","env":[{"name":"ANTHROPIC_BASE_URL","value":"https://api.example.com/anthropic"}]}}}"#).unwrap(); + let mut status = session_status(); + status.provider = agent::ProviderKind::ClaudeCode; + status.requested_profile_id = Some("custom".into()); + let mut timeline = Timeline::default(); + timeline.apply_at( + None, + &agent::AgentEvent::TokenUsage(agent::TokenUsage { + freshness: agent::ContextFreshness::Current, + used_tokens: Some(1234), + ..Default::default() + }), + ); + let mut providers = ProvidersStatus::default(); + providers.provider_usage.insert( + "custom".into(), + tcode_core::usage::ProviderUsage { + error: Some("temporarily unreachable".into()), + ..Default::default() + }, + ); + let custom = composer_state(Some(&status), Some(&timeline), &settings, &providers); + assert_eq!(custom.token_usage.unwrap().used_tokens, Some(1234)); + assert!(custom.usage.is_none()); + settings + .profiles + .get_mut("custom") + .unwrap() + .settings + .env + .clear(); + let native = composer_state(Some(&status), Some(&timeline), &settings, &providers); + assert_eq!( + native.usage.unwrap().error.as_deref(), + Some("temporarily unreachable") + ); + assert_eq!(native.token_usage.unwrap().used_tokens, Some(1234)); + } + + #[test] + #[cfg(unix)] + fn claude_usage_replays_adapter_timeline_and_composer_across_resume() { + use std::os::unix::fs::PermissionsExt; + let root = std::env::temp_dir().join(format!( + "tcode-usage-replay-{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + )); + std::fs::create_dir_all(&root).unwrap(); + let binary = root.join("claude-fixture"); + std::fs::write(&binary, "#!/bin/sh\ncase \"$*\" in *--version*) echo '2.1.200'; exit;; esac\nIFS= read -r request\nif [ \"$TCODE_USAGE_INTERRUPTED\" = 1 ]; then IFS= read -r request; fi\ncat \"$TCODE_USAGE_FIXTURE\"\nwhile IFS= read -r request; do :; done\n").unwrap(); + std::fs::set_permissions(&binary, std::fs::Permissions::from_mode(0o700)).unwrap(); + let mut timeline = Timeline::default(); + let mut replay = Timeline::default(); + let mut status = session_status(); + status.provider = agent::ProviderKind::ClaudeCode; + status.requested_model = Some("claude-opus-5".into()); + status.provider_option_selections = vec![agent::OptionSelection { + id: "contextWindow".into(), + value: serde_json::json!(300000), + }]; + let settings = Settings::default(); + let providers = ProvidersStatus::default(); + let mut recorded_events = Vec::new(); + for (index, fixture) in [ + include_str!("../../../agent/tests/fixtures/claude/usage_scope.jsonl"), + include_str!("../../../agent/tests/fixtures/claude/usage_resume.jsonl"), + include_str!("../../../agent/tests/fixtures/claude/usage_recorded.jsonl"), + ] + .into_iter() + .enumerate() + { + let path = root.join(format!("fixture-{index}.jsonl")); + std::fs::write(&path, fixture).unwrap(); + smol::block_on(async { + let handle = agent::claude::start(agent::SessionOptions { + cwd: root.clone(), + model: Some("claude-opus-5".into()), + resume: (index > 0).then(|| { + agent::ResumeCursor(serde_json::json!({"session_id":"fixture-session"})) + }), + fork: false, + binary_path: Some(binary.clone()), + approval_mode: agent::ApprovalMode::Supervised, + option_selections: vec![], + interaction_mode: agent::InteractionMode::Build, + mcp_servers: vec![], + launch_env: agent::LaunchEnv { + home: Some(root.clone()), + env: vec![ + ( + "TCODE_USAGE_FIXTURE".into(), + path.to_string_lossy().into_owned(), + ), + ( + "TCODE_USAGE_INTERRUPTED".into(), + if index == 1 { "1" } else { "0" }.into(), + ), + ], + }, + extra_args: vec![], + acp: None, + }) + .await + .unwrap(); + handle + .commands + .send(agent::SessionCommand::SendTurn { + delivery_id: 1, + text: "fixture".into(), + options: None, + attachments: vec![], + }) + .await + .unwrap(); + if index == 1 { + handle + .commands + .send(agent::SessionCommand::Interrupt) + .await + .unwrap(); + } + loop { + let event = + smol::future::race(async { handle.events.recv().await.unwrap() }, async { + smol::Timer::after(std::time::Duration::from_secs(10)).await; + panic!("fixture adapter timed out") + }) + .await; + timeline.apply_at(Some(1 + recorded_events.len() as u64), &event); + let snapshot = + composer_state(Some(&status), Some(&timeline), &settings, &providers); + if let agent::AgentEvent::ContextCompacted(c) = &event { + let u = snapshot.token_usage.unwrap(); + if c.in_progress { + assert_eq!(u.freshness, agent::ContextFreshness::Compacting); + } else { + assert_eq!(u.used_tokens, None); + assert_eq!(c.pre_tokens, Some(500260)); + assert_eq!(c.trigger.as_deref(), Some("manual")); + assert_eq!(crate::context_meter::used_tokens(&u), None); + } + } + if let agent::AgentEvent::TokenUsage(u) = &event { + assert_ne!( + u.used_tokens, + Some(900000), + "subagent context must not enter main usage" + ); + if u.output_tokens == Some(20) { + assert_eq!(u.used_tokens, Some(400150)); + } + } + let completed = matches!(event, agent::AgentEvent::TurnCompleted { .. }); + recorded_events.push(event); + if completed { + break; + } + } + handle + .commands + .send(agent::SessionCommand::Shutdown) + .await + .unwrap(); + }); + let snapshot = composer_state(Some(&status), Some(&timeline), &settings, &providers); + let usage = snapshot.token_usage.unwrap(); + assert_eq!(usage.used_tokens, Some([500260, 60, 20764][index])); + assert_eq!( + usage.total_processed_tokens, + Some([4001300, 4001365, 4063449][index]) + ); + if index == 1 { + assert_eq!( + timeline.last_turn_status, + Some(agent::TurnStatus::Interrupted) + ); + } + // A repeated persisted completion is idempotent, including a cancelled result. + timeline.apply_at(Some(99), recorded_events.last().unwrap()); + assert_eq!( + timeline.usage.unwrap().total_processed_tokens, + usage.total_processed_tokens + ); + } + for (index, event) in recorded_events.iter().enumerate() { + replay.apply_at(Some(1 + index as u64), event); + } + assert_eq!(replay.usage, timeline.usage); + let old: agent::AgentEvent = serde_json::from_str(r#"{"type":"token_usage","used_tokens":4100000,"input_tokens":1000000,"context_window":1000000,"total_processed_tokens":4100000}"#).unwrap(); + replay.apply_at(None, &old); + let old = composer_state(Some(&status), Some(&replay), &settings, &providers) + .token_usage + .unwrap(); + assert_eq!(old.freshness, agent::ContextFreshness::Unknown); + assert_eq!( + crate::context_meter::used_tokens(&old), + None, + "legacy aggregate has no occupancy provenance" + ); + std::fs::remove_dir_all(root).unwrap(); + } } diff --git a/docs/DESIGN.md b/docs/DESIGN.md index 74b7ba9b..a795e0e4 100644 --- a/docs/DESIGN.md +++ b/docs/DESIGN.md @@ -679,6 +679,27 @@ Send above; approval mode and Build/Plan below, with the context ring at the trailing edge of the second row. The ring keeps a 44pt hit target opening its details sheet. Option labels truncate when space is tight; both rows fit at 360pt in English and Simplified Chinese. +The context details distinguish the latest main-conversation request from +processed traffic. Claude occupancy is the latest input plus cache-read and +cache-creation tokens; generated output and repeated requests do not inflate it. +Capacity is separate and respects the selected model limit. Unknown observations +say **Unknown**, with no measured percentage or empty progress bar. During a new +turn the previous observation says **Last known context · updating** until replaced. +Older saved events without provenance stay unverified until a new observation; +their aggregate counts are never relabeled as measured occupancy. **Total processed** +is a separate lifetime main-loop total, reconstructed from completed turn traffic. +Compaction has separate in-progress and completed dividers, with supplied trigger +and pre-compaction count. Labels wrap at narrow widths. Completion invalidates +occupancy until another request observation, even when the harness reports a +post-compaction size. Warning color is not a claim about the harness's trigger. + +Settings → Usage and the composer's account limits use the resolved profile's +account-usage capability. Unsupported custom endpoints/API-key configurations +are omitted independently of session context usage. Native account profiles, +including custom-named profiles, retain sign-in/network errors and retry. +Changing profile configuration or secrets invalidates cached limits and rejects +older in-flight results. + Sending during a turn queues the message; the secondary send action steers when the provider supports it. Stop interrupts the current turn. Queue/steer guidance belongs in the send tooltip. diff --git a/locales/en.yml b/locales/en.yml index 7502e9f9..c6f3a73e 100644 --- a/locales/en.yml +++ b/locales/en.yml @@ -23,6 +23,10 @@ chat: work_log_failed: "Failed" work_log_interrupted: "Interrupted" served_model_tooltip: "Serving model differs from the requested model" + context_manual: "Manual" + context_auto: "Automatic" + context_compacting: "Compacting context…" + context_pre_tokens: "%{tokens} tokens before compaction" context_compacted: "Context compacted" context_window_changed: "Context window set to %{window}" previous_logs: "+%{count} previous log entries" @@ -615,6 +619,10 @@ composer: context_window: "Context window" no_usage: "No usage yet this session." context_window_title: "Context Window" + context_unknown: "Unknown" + context_latest: "Latest request context" + context_last_known: "Last known context · updating" + context_updating: "Waiting for context observation" total_processed: "Total processed" steer_tooltip: "Steer the current turn" compacts_automatically: "%{provider} automatically compacts its context when needed." diff --git a/locales/zh-CN.yml b/locales/zh-CN.yml index 3fbd5608..d30a0496 100644 --- a/locales/zh-CN.yml +++ b/locales/zh-CN.yml @@ -23,6 +23,10 @@ chat: work_log_failed: "失败" work_log_interrupted: "已中断" served_model_tooltip: "实际服务模型与请求的模型不一致" + context_manual: "手动" + context_auto: "自动" + context_compacting: "正在压缩上下文…" + context_pre_tokens: "压缩前 %{tokens} 个令牌" context_compacted: "上下文已压缩" context_window_changed: "上下文长度已修改为 %{window}" previous_logs: "前面还有 %{count} 条日志" @@ -609,6 +613,10 @@ composer: context_window: "上下文窗口" no_usage: "本会话暂无用量信息。" context_window_title: "上下文窗口" + context_unknown: "未知" + context_latest: "最近请求的上下文" + context_last_known: "上次已知上下文 · 更新中" + context_updating: "等待上下文观测值" total_processed: "累计处理" steer_tooltip: "引导当前轮次" compacts_automatically: "%{provider} 会在需要时自动压缩其上下文。" From e100f66662732e48bc9a959d3462b385c6438e8b Mon Sep 17 00:00:00 2001 From: Tryanks Date: Tue, 8 Sep 2026 19:12:44 +0800 Subject: [PATCH 2/4] fix(usage): clear request identity before a new turn --- crates/agent/src/claude.rs | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/crates/agent/src/claude.rs b/crates/agent/src/claude.rs index 9b6bbac4..99428b9c 100644 --- a/crates/agent/src/claude.rs +++ b/crates/agent/src/claude.rs @@ -1298,6 +1298,9 @@ impl Mapper { /// Allocate the next synthesized turn id and mark it in-flight. fn start_turn(&mut self) -> String { + self.current_message_id = None; + self.usage_message_id = None; + self.request_usage = json!({}); if self.latest_usage.freshness == crate::ContextFreshness::Current { self.latest_usage.freshness = crate::ContextFreshness::LastKnown; } @@ -4294,6 +4297,30 @@ mod tests { r#"{"type":"result","subtype":"success","usage":{"input_tokens":1000,"cache_read_input_tokens":4000000,"output_tokens":200},"modelUsage":{"claude":{"contextWindow":1000000}}}"#, ); assert!(result.iter().any(|e| matches!(e, AgentEvent::TurnCompleted { usage: Some(u), .. } if u.used_tokens == Some(550) && u.turn_processed_tokens == Some(4001200))), "result accounting must not replace occupancy: {result:?}"); + m.start_turn(); + let missing_start = feed( + &mut m, + r#"{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":1}}}"#, + ); + assert!( + missing_start.is_empty(), + "a new turn cannot reuse the previous request identity: {missing_start:?}" + ); + assert_eq!(m.latest_usage.freshness, crate::ContextFreshness::LastKnown); + let unknown = feed( + &mut m, + r#"{"type":"stream_event","event":{"type":"message_start","message":{"id":"usage-2"}}}"#, + ); + assert!( + matches!(unknown.last(), Some(AgentEvent::TokenUsage(u)) if u.used_tokens.is_none()) + ); + let output_only = feed( + &mut m, + r#"{"type":"stream_event","event":{"type":"message_delta","usage":{"output_tokens":1}}}"#, + ); + assert!( + matches!(output_only.last(), Some(AgentEvent::TokenUsage(u)) if u.used_tokens.is_none()) + ); } #[test] From 28091a121cd76cb1b52814e86dfaa4600e43cd26 Mon Sep 17 00:00:00 2001 From: Tryanks Date: Tue, 8 Sep 2026 19:23:00 +0800 Subject: [PATCH 3/4] fix(usage): retain totals and capacity in timestamp-free replay --- crates/core/src/session.rs | 56 +++++++++++++++++++++++++++----- crates/ui/src/store/snapshots.rs | 3 ++ 2 files changed, 50 insertions(+), 9 deletions(-) diff --git a/crates/core/src/session.rs b/crates/core/src/session.rs index 6bddeae0..8af014b8 100644 --- a/crates/core/src/session.rs +++ b/crates/core/src/session.rs @@ -624,7 +624,7 @@ impl Timeline { // Reuse the open turn (typically opened by the user message); // otherwise begin a fresh one. let turn = match self.current_turn { - Some(t) if self.turns[t].end_ts.is_none() => t, + Some(t) if self.turn_is_open() => t, _ => self.push_turn(ts), }; // TurnStarted is the authoritative turn start; prefer it over @@ -690,6 +690,9 @@ impl Timeline { .current_turn .is_none_or(|turn| self.turns[turn].status.is_none()); let usage = usage.map(|mut usage| { + usage.context_window = usage + .context_window + .or(self.usage.and_then(|u| u.context_window)); if let Some(processed) = usage.turn_processed_tokens { let previous = self .usage @@ -957,7 +960,6 @@ impl Timeline { /// a `TurnCompleted` has been folded, which records a status even when the /// event carried no timestamp to store as `end_ts`; both must be checked or /// a stray later transition would leak into the next turn's accounting - /// (`TurnStarted` reuses a turn whose `end_ts` is unset). fn turn_is_open(&self) -> bool { self.current_turn.is_some_and(|turn| { let turn = &self.turns[turn]; @@ -1477,6 +1479,43 @@ mod tests { }) } + #[test] + fn usage_replay_without_timestamps_keeps_turn_totals_and_known_capacity() { + let usage = |processed| TokenUsage { + freshness: agent::ContextFreshness::Current, + used_tokens: Some(550), + turn_processed_tokens: processed, + ..Default::default() + }; + let timeline = Timeline::fold_events([ + AgentEvent::TurnStarted { + turn_id: "one".into(), + }, + AgentEvent::TokenUsage(TokenUsage { + context_window: Some(1000000), + ..usage(None) + }), + AgentEvent::TurnCompleted { + turn_id: "one".into(), + status: TurnStatus::Completed, + usage: Some(usage(Some(100))), + }, + AgentEvent::TurnStarted { + turn_id: "two".into(), + }, + AgentEvent::TokenUsage(usage(None)), + AgentEvent::TurnCompleted { + turn_id: "two".into(), + status: TurnStatus::Interrupted, + usage: Some(usage(Some(20))), + }, + ]); + let usage = timeline.usage.unwrap(); + assert_eq!(usage.total_processed_tokens, Some(120)); + assert_eq!(usage.context_window, Some(1000000)); + assert_eq!(timeline.turns.len(), 2); + } + #[test] fn provider_relay_marker_folds_before_the_next_user_message() { let timeline = Timeline::fold_events([ @@ -3270,11 +3309,8 @@ mod tests { #[test] fn a_turn_finalized_without_a_timestamp_rejects_later_tool_transitions() { - // An untimestamped TurnCompleted records a status but no end_ts, and a - // following TurnStarted therefore *reuses* that turn. A stray tool - // transition in between must not survive into the reopened turn's - // accounting — otherwise the ghost item stays open and eats the whole - // next turn. + // A terminal status closes the turn even without an end timestamp. + // A stray transition between turns must not enter the next turn's clock. let timeline = Timeline::fold_events(vec![ at(1_000, turn_started()), turn_completed().into(), @@ -3283,9 +3319,11 @@ mod tests { at(10_000, turn_completed()), ]); - let timing = timeline.turns[0] + assert_eq!(timeline.turns.len(), 2); + assert_eq!(timeline.turns[0].timing, None); + let timing = timeline.turns[1] .timing - .expect("the reopened turn is fully timestamped"); + .expect("the new turn is fully timestamped"); assert_eq!(timing.total_ms, 6_000); assert_eq!(timing.tool_ms, 0); } diff --git a/crates/ui/src/store/snapshots.rs b/crates/ui/src/store/snapshots.rs index 7a6e79a1..b16e238d 100644 --- a/crates/ui/src/store/snapshots.rs +++ b/crates/ui/src/store/snapshots.rs @@ -517,6 +517,7 @@ mod tests { let snapshot = composer_state(Some(&status), Some(&timeline), &settings, &providers); let usage = snapshot.token_usage.unwrap(); assert_eq!(usage.used_tokens, Some([500260, 60, 20764][index])); + assert_eq!(usage.context_window, Some(300000)); assert_eq!( usage.total_processed_tokens, Some([4001300, 4001365, 4063449][index]) @@ -538,6 +539,8 @@ mod tests { replay.apply_at(Some(1 + index as u64), event); } assert_eq!(replay.usage, timeline.usage); + let untimed = Timeline::fold_events(recorded_events.clone()); + assert_eq!(untimed.usage, timeline.usage); let old: agent::AgentEvent = serde_json::from_str(r#"{"type":"token_usage","used_tokens":4100000,"input_tokens":1000000,"context_window":1000000,"total_processed_tokens":4100000}"#).unwrap(); replay.apply_at(None, &old); let old = composer_state(Some(&status), Some(&replay), &settings, &providers) From c491be947e9921b1d39d186de1db9ffc9a5e7e84 Mon Sep 17 00:00:00 2001 From: Tryanks Date: Tue, 8 Sep 2026 19:37:10 +0800 Subject: [PATCH 4/4] Keep known context capacity visible while occupancy is unknown --- crates/ui/src/composer/components/pickers.rs | 9 ++++----- docs/DESIGN.md | 2 +- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/crates/ui/src/composer/components/pickers.rs b/crates/ui/src/composer/components/pickers.rs index 3de82105..958ae1dc 100644 --- a/crates/ui/src/composer/components/pickers.rs +++ b/crates/ui/src/composer/components/pickers.rs @@ -1545,14 +1545,13 @@ fn render_context_meter_pane( let used = usage.as_ref().and_then(context_meter::used_tokens); let max = usage.and_then(|u| u.context_window); let pct_label = context_meter::format_percentage(pct); - let stat: AnyElement = match (max, pct_label.clone()) { - (Some(max), Some(pct_label)) => h_flex() + let stat: AnyElement = match max { + Some(max) => h_flex() .gap_1() .text_size(px(11.)) .font_family(cx.theme().mono_font_family.clone()) .text_color(muted) - .child(pct_label) - .child("·") + .when_some(pct_label, |row, label| row.child(label).child("·")) .child(format!( "{}/{}", context_meter::format_tokens(used), @@ -1614,7 +1613,7 @@ fn render_context_meter_pane( ); // "Total processed" — the session-cumulative token count, when the provider - // reports it (a native running total or adapter-side accumulation). + // reports it (a native running total or timeline accumulation). if let Some(total) = usage.and_then(|u| u.total_processed_tokens) { pane = pane.child( h_flex() diff --git a/docs/DESIGN.md b/docs/DESIGN.md index a795e0e4..f5805922 100644 --- a/docs/DESIGN.md +++ b/docs/DESIGN.md @@ -682,7 +682,7 @@ both rows fit at 360pt in English and Simplified Chinese. The context details distinguish the latest main-conversation request from processed traffic. Claude occupancy is the latest input plus cache-read and cache-creation tokens; generated output and repeated requests do not inflate it. -Capacity is separate and respects the selected model limit. Unknown observations +Capacity stays visible independently and respects the selected model limit. Unknown observations say **Unknown**, with no measured percentage or empty progress bar. During a new turn the previous observation says **Last known context · updating** until replaced. Older saved events without provenance stay unverified until a new observation;