diff --git a/src/apps/desktop/src/api/agentic_api.rs b/src/apps/desktop/src/api/agentic_api.rs index 52d4ec353..a870ffb92 100644 --- a/src/apps/desktop/src/api/agentic_api.rs +++ b/src/apps/desktop/src/api/agentic_api.rs @@ -17,8 +17,8 @@ use crate::startup_trace::DesktopStartupTrace; use bitfun_agent_runtime::deep_review::sanitize_focused_review_public_metadata; use bitfun_agent_runtime::sdk::{ AgentDialogSteerRequest, AgentDialogTurnExecution, AgentDialogTurnRequest, - AgentInputAttachment, AgentSessionCreateResult, AgentSessionModelSelection, - AgentSessionModeUpdateRequest, AgentSessionModelSelectionUpdateRequest, + AgentInputAttachment, AgentSessionCreateResult, AgentSessionModeUpdateRequest, + AgentSessionModelSelection, AgentSessionModelSelectionUpdateRequest, AgentSessionModelUpdateRequest, AgentSubmissionSource, AgentTurnCancellationRequest, DialogSteerOutcome, PermissionAuditRecord, PermissionGrant, PermissionGrantKey, PermissionReply, PermissionRequest, @@ -54,7 +54,7 @@ use bitfun_core::service::config::project_permission_store::{ use bitfun_core::service::remote_ssh::workspace_state::is_remote_path; use bitfun_core::service::remote_ssh::workspace_state::resolve_workspace_session_identity; use bitfun_core::service::session::{ - DialogTurnData, SessionMemoryMode, SessionMetadata, SessionRelationship, + DialogTurnData, SessionContextUsage, SessionMemoryMode, SessionMetadata, SessionRelationship, SessionRelationshipKind, SessionTurnCatalog, SessionTurnWindowResponse, }; use bitfun_core::service::workspace::WorkspaceKind; @@ -502,6 +502,7 @@ pub struct RestoreSessionWithTurnsResponse { pub struct RestoreSessionViewResponse { pub session: SessionResponse, pub turns: Vec, + pub current_context_usage: Option, pub turn_catalog: SessionTurnCatalog, pub context_restore_state: String, pub is_partial: bool, @@ -3072,6 +3073,7 @@ pub async fn restore_session_view( .map_err(|error| format!("Failed to restore session view: {error}"))?; let session = restored.session; let mut turns = restored.turns; + let current_context_usage = restored.current_context_usage; let total_turn_count = restored.total_turn_count; let turn_catalog = restored.turn_catalog; let timings = restored.timings; @@ -3124,6 +3126,7 @@ pub async fn restore_session_view( Ok(RestoreSessionViewResponse { session: session_to_response_with_turn_count(session, total_turn_count), turns, + current_context_usage, turn_catalog, context_restore_state: "pending".to_string(), is_partial, diff --git a/src/apps/desktop/src/lib.rs b/src/apps/desktop/src/lib.rs index efb6381dc..b05a671fd 100644 --- a/src/apps/desktop/src/lib.rs +++ b/src/apps/desktop/src/lib.rs @@ -1944,6 +1944,14 @@ async fn init_agentic_system() -> anyhow::Result<( bitfun_core::service::token_usage::TokenUsageSubscriber::new(token_usage_service.clone()), ); event_router.subscribe_internal("token_usage".to_string(), token_usage_subscriber); + event_router.subscribe_internal( + "session_context_usage".to_string(), + Arc::new( + bitfun_core::agentic::session::SessionContextUsageSubscriber::new( + session_manager.clone(), + ), + ), + ); event_router.subscribe_internal( "thread_goal_tokens".to_string(), Arc::new(bitfun_core::agentic::goal_mode::ThreadGoalTokenSubscriber), diff --git a/src/apps/desktop/src/runtime/session_application.rs b/src/apps/desktop/src/runtime/session_application.rs index 17997b7ac..671bbc3ad 100644 --- a/src/apps/desktop/src/runtime/session_application.rs +++ b/src/apps/desktop/src/runtime/session_application.rs @@ -25,8 +25,9 @@ use bitfun_core::service::remote_ssh::workspace_state::{ }; use bitfun_core::service::remote_ssh::SSHConnectionManager; use bitfun_core::service::session::{ - DialogTurnData, DialogTurnKind, SessionMetadata, SessionStatus, SessionTranscriptExport, - SessionTranscriptExportOptions, SessionTurnCatalog, SessionTurnWindowResponse, + DialogTurnData, DialogTurnKind, SessionContextUsage, SessionMetadata, SessionStatus, + SessionTranscriptExport, SessionTranscriptExportOptions, SessionTurnCatalog, + SessionTurnWindowResponse, }; use bitfun_core::service::session_usage::SessionUsageReport; use bitfun_core::service::token_usage::TokenUsageService; @@ -36,12 +37,7 @@ use bitfun_runtime_ports::{AgentContextReloadRequest, SessionTurnWindowRequest}; use serde::{Deserialize, Serialize}; use tokio::sync::RwLock; -const UI_CUSTOM_METADATA_KEYS: [&str; 4] = [ - "titleSource", - "titleKey", - "titleParams", - "lastRequestTokenUsage", -]; +const UI_CUSTOM_METADATA_KEYS: [&str; 3] = ["titleSource", "titleKey", "titleParams"]; #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] @@ -129,6 +125,7 @@ fn local_command_turn_record_request( pub(crate) struct DesktopSessionViewRestore { pub session: Session, pub turns: Vec, + pub current_context_usage: Option, pub total_turn_count: usize, pub turn_catalog: SessionTurnCatalog, pub timings: SessionViewRestoreTiming, @@ -827,10 +824,17 @@ impl DesktopSessionApplication { .loaded_session_snapshot(session_id) .map_err(|error| DesktopSessionApplicationError::Core(error.to_string()))?; overlay_live_session_state(&mut session, live_session); + let current_context_usage = self + .compatibility + .load_persisted_session_metadata(&storage_path, session_id) + .await + .map_err(desktop_core_session_error)? + .and_then(|metadata| metadata.current_context_usage); timings.resolve_storage_path_duration_ms = resolve_storage_path_duration_ms; Ok(DesktopSessionViewRestore { session, turns, + current_context_usage, total_turn_count, turn_catalog, timings, @@ -940,7 +944,9 @@ mod tests { AgentSessionWorkspaceBinding, AgentSessionWorkspaceRequest, AgentSubmissionPort, AgentSubmissionRequest, AgentSubmissionResult, PortError, PortErrorKind, PortResult, }; - use bitfun_core::service::session::{SessionKind, SessionMemoryMode}; + use bitfun_core::service::session::{ + SessionContextUsage, SessionContextUsageSource, SessionKind, SessionMemoryMode, + }; use serde_json::json; use std::sync::Mutex; @@ -1364,6 +1370,14 @@ mod tests { "titleSource": "i18n", "titleKey": "old" })); + current.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-7".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }); let mut incoming = current.clone(); incoming.session_name = "Renamed".to_string(); @@ -1373,6 +1387,14 @@ mod tests { incoming.status = SessionStatus::Active; incoming.turn_count = 1; incoming.review_action_state = Some(json!({ "phase": "fixing" })); + incoming.current_context_usage = Some(SessionContextUsage { + turn_id: "stale-turn".to_string(), + input_tokens: 1, + output_tokens: None, + total_tokens: 1, + timestamp: 1, + source: SessionContextUsageSource::ContextCompression, + }); incoming.custom_metadata = Some(json!({ "titleSource": "i18n", "titleKey": "new", @@ -1400,6 +1422,17 @@ mod tests { assert_eq!(current.status, SessionStatus::Archived); assert_eq!(current.turn_count, 7); assert_eq!(current.review_action_state, incoming.review_action_state); + assert_eq!( + current.current_context_usage, + Some(SessionContextUsage { + turn_id: "turn-7".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }) + ); let custom = current.custom_metadata.unwrap(); assert_eq!(custom["threadGoal"]["objective"], "preserve"); assert_eq!(custom["titleKey"], "new"); diff --git a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs index e95fefc0b..aaad99038 100644 --- a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs +++ b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs @@ -2665,6 +2665,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet snapshot_session_id: None, tags: Vec::new(), custom_metadata: None, + current_context_usage: None, relationship: None, todos: None, review_action_state: None, diff --git a/src/crates/assembly/core/src/agentic/session/context_usage.rs b/src/crates/assembly/core/src/agentic/session/context_usage.rs new file mode 100644 index 000000000..f2bf44b82 --- /dev/null +++ b/src/crates/assembly/core/src/agentic/session/context_usage.rs @@ -0,0 +1,157 @@ +use super::SessionManager; +use crate::agentic::events::{AgenticEvent, EventSubscriber}; +use bitfun_agent_runtime::event_bus::EventSubscriberResult; +use bitfun_services_core::session::{SessionContextUsage, SessionContextUsageSource}; +use log::warn; +use std::sync::Arc; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// Persists the runtime-owned context usage used by every product surface. +pub struct SessionContextUsageSubscriber { + session_manager: Arc, +} + +impl SessionContextUsageSubscriber { + pub fn new(session_manager: Arc) -> Self { + Self { session_manager } + } +} + +fn now_ms() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis() + .min(u128::from(u64::MAX)) as u64 +} + +fn usage_from_event(event: &AgenticEvent) -> Option<(&str, SessionContextUsage)> { + match event { + AgenticEvent::TokenUsageUpdated { + session_id, + turn_id, + input_tokens, + output_tokens, + total_tokens, + .. + } => Some(( + session_id, + SessionContextUsage { + turn_id: turn_id.clone(), + input_tokens: *input_tokens as u64, + output_tokens: output_tokens.map(|value| value as u64), + total_tokens: *total_tokens as u64, + timestamp: now_ms(), + source: SessionContextUsageSource::ModelRequest, + }, + )), + AgenticEvent::ContextCompressionCompleted { + session_id, + turn_id, + tokens_after, + applied: true, + .. + } => Some(( + session_id, + SessionContextUsage { + turn_id: turn_id.clone(), + input_tokens: *tokens_after as u64, + output_tokens: None, + total_tokens: *tokens_after as u64, + timestamp: now_ms(), + source: SessionContextUsageSource::ContextCompression, + }, + )), + _ => None, + } +} + +#[async_trait::async_trait] +impl EventSubscriber for SessionContextUsageSubscriber { + async fn on_event(&self, event: &AgenticEvent) -> EventSubscriberResult { + let Some((session_id, usage)) = usage_from_event(event) else { + return Ok(()); + }; + + if let Err(error) = self + .session_manager + .persist_current_context_usage(session_id, usage) + .await + { + warn!( + "Failed to persist session context usage: session_id={}, error={}", + session_id, error + ); + } + + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn maps_model_request_usage() { + let event = AgenticEvent::TokenUsageUpdated { + session_id: "session-1".to_string(), + turn_id: "turn-1".to_string(), + model_config_id: "model-config".to_string(), + effective_model_name: "model".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + max_context_tokens: Some(128_000), + is_subagent: false, + cached_tokens: None, + token_details: None, + }; + + let (session_id, usage) = usage_from_event(&event).expect("usage event"); + assert_eq!(session_id, "session-1"); + assert_eq!(usage.turn_id, "turn-1"); + assert_eq!(usage.input_tokens, 42_000); + assert_eq!(usage.output_tokens, Some(1_500)); + assert_eq!(usage.total_tokens, 43_500); + assert_eq!(usage.source, SessionContextUsageSource::ModelRequest); + } + + #[test] + fn maps_only_applied_context_compression() { + let event = AgenticEvent::ContextCompressionCompleted { + session_id: "session-1".to_string(), + turn_id: "turn-1".to_string(), + compression_id: "compression-1".to_string(), + compression_count: 1, + tokens_before: 90_000, + tokens_after: 15_000, + compression_ratio: 0.17, + duration_ms: 500, + has_summary: true, + summary_source: "model".to_string(), + applied: true, + }; + + let (_, usage) = usage_from_event(&event).expect("applied compression"); + assert_eq!(usage.input_tokens, 15_000); + assert_eq!(usage.output_tokens, None); + assert_eq!(usage.total_tokens, 15_000); + assert_eq!(usage.source, SessionContextUsageSource::ContextCompression); + + let not_applied = AgenticEvent::ContextCompressionCompleted { + session_id: "session-1".to_string(), + turn_id: "turn-1".to_string(), + compression_id: "compression-1".to_string(), + compression_count: 1, + tokens_before: 90_000, + tokens_after: 15_000, + compression_ratio: 0.17, + duration_ms: 500, + has_summary: true, + summary_source: "model".to_string(), + applied: false, + }; + assert!(usage_from_event(¬_applied).is_none()); + } +} diff --git a/src/crates/assembly/core/src/agentic/session/mod.rs b/src/crates/assembly/core/src/agentic/session/mod.rs index 5d8219553..2a9db7374 100644 --- a/src/crates/assembly/core/src/agentic/session/mod.rs +++ b/src/crates/assembly/core/src/agentic/session/mod.rs @@ -4,6 +4,7 @@ pub mod compression; pub mod context_store; +mod context_usage; pub mod evidence_ledger; pub mod file_read_state; pub mod prompt_cache; @@ -16,6 +17,7 @@ pub mod turn_skill_agent_snapshot_store; pub use compression::*; pub use context_store::*; +pub use context_usage::*; pub use evidence_ledger::*; pub use file_read_state::*; pub use prompt_cache::*; diff --git a/src/crates/assembly/core/src/agentic/session/session_manager.rs b/src/crates/assembly/core/src/agentic/session/session_manager.rs index 0d3b41c28..aba2a8b11 100644 --- a/src/crates/assembly/core/src/agentic/session/session_manager.rs +++ b/src/crates/assembly/core/src/agentic/session/session_manager.rs @@ -38,9 +38,9 @@ use crate::service::config::{ }; use crate::service::remote_ssh::workspace_state::LOCAL_WORKSPACE_SSH_HOST; use crate::service::session::{ - DialogTurnData, DialogTurnKind, ModelRoundData, SessionMemoryMode, SessionMetadata, - SessionRelationship, SessionStatus, TextItemData, ThinkingItemData, ToolCallData, ToolItemData, - ToolResultData, TranscriptLineRange, TurnStatus, UserMessageData, + DialogTurnData, DialogTurnKind, ModelRoundData, SessionContextUsage, SessionMemoryMode, + SessionMetadata, SessionRelationship, SessionStatus, TextItemData, ThinkingItemData, + ToolCallData, ToolItemData, ToolResultData, TranscriptLineRange, TurnStatus, UserMessageData, }; use crate::service::snapshot::{ ensure_snapshot_manager_for_workspace, get_or_create_snapshot_manager, @@ -6252,6 +6252,30 @@ impl SessionManager { .await } + pub(crate) async fn persist_current_context_usage( + &self, + session_id: &str, + usage: SessionContextUsage, + ) -> BitFunResult<()> { + let _mutation_guard = self.acquire_session_mutation(session_id).await?; + let should_persist_usage = self.sessions.get(session_id).is_some_and(|session| { + !session.agent_type.starts_with("acp:") + && session + .dialog_turn_ids + .iter() + .any(|turn_id| turn_id == &usage.turn_id) + && self.should_persist_session(&session) + }); + if !should_persist_usage || !self.config.enable_persistence { + return Ok(()); + } + + self.update_persisted_session_metadata(session_id, |metadata| { + metadata.current_context_usage = Some(usage); + }) + .await + } + pub async fn merge_session_relationship( &self, session_id: &str, @@ -8007,9 +8031,10 @@ mod tests { AIConfig as ServiceAIConfig, AIModelConfig as ServiceAIModelConfig, }; use crate::service::session::{ - DialogTurnData, DialogTurnKind, ModelRoundData, SessionKind, SessionMetadata, - SessionRelationship, SessionRelationshipKind, ToolCallData, ToolItemData, ToolResultData, - TurnStatus, UserMessageData, + DialogTurnData, DialogTurnKind, ModelRoundData, SessionContextUsage, + SessionContextUsageSource, SessionKind, SessionMetadata, SessionRelationship, + SessionRelationshipKind, ToolCallData, ToolItemData, ToolResultData, TurnStatus, + UserMessageData, }; use crate::util::errors::BitFunError; use bitfun_core_types::{ @@ -8228,6 +8253,163 @@ mod tests { ) } + #[tokio::test] + async fn current_context_usage_is_persisted_by_the_session_owner() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = test_manager(persistence_manager.clone()); + let session = manager + .create_session( + "Usage persistence".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..SessionConfig::default() + }, + ) + .await + .expect("session should create"); + let usage = SessionContextUsage { + turn_id: "turn-1".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }; + manager + .sessions + .get_mut(&session.session_id) + .expect("session should be active") + .dialog_turn_ids + .push(usage.turn_id.clone()); + + manager + .persist_current_context_usage(&session.session_id, usage.clone()) + .await + .expect("usage should persist"); + + let metadata = persistence_manager + .load_session_metadata(workspace.path(), &session.session_id) + .await + .expect("metadata should load") + .expect("metadata should exist"); + assert_eq!(metadata.current_context_usage, Some(usage)); + } + + #[tokio::test] + async fn acp_context_usage_is_not_persisted_as_native_prompt_usage() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = test_manager(persistence_manager.clone()); + let session = manager + .create_session( + "ACP usage".to_string(), + "acp:codex".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..SessionConfig::default() + }, + ) + .await + .expect("session should create"); + + manager + .persist_current_context_usage( + &session.session_id, + SessionContextUsage { + turn_id: "turn-1".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }, + ) + .await + .expect("ACP usage should be ignored"); + + let metadata = persistence_manager + .load_session_metadata(workspace.path(), &session.session_id) + .await + .expect("metadata should load") + .expect("metadata should exist"); + assert!(metadata.current_context_usage.is_none()); + } + + #[tokio::test] + async fn delayed_context_usage_does_not_restore_a_removed_turn() { + let workspace = TestWorkspace::new(); + let persistence_manager = Arc::new( + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"), + ); + let manager = Arc::new(test_manager(persistence_manager.clone())); + let session = manager + .create_session( + "Delayed usage".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().to_string()), + ..SessionConfig::default() + }, + ) + .await + .expect("session should create"); + manager + .sessions + .get_mut(&session.session_id) + .expect("session should be active") + .dialog_turn_ids + .push("turn-1".to_string()); + + let mutation_guard = manager + .acquire_session_mutation(&session.session_id) + .await + .expect("mutation guard"); + let delayed_manager = manager.clone(); + let delayed_session_id = session.session_id.clone(); + let delayed_write = tokio::spawn(async move { + delayed_manager + .persist_current_context_usage( + &delayed_session_id, + SessionContextUsage { + turn_id: "turn-1".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }, + ) + .await + }); + tokio::task::yield_now().await; + assert!(!delayed_write.is_finished()); + + manager + .sessions + .get_mut(&session.session_id) + .expect("session should remain active") + .dialog_turn_ids + .clear(); + drop(mutation_guard); + delayed_write + .await + .expect("delayed write should join") + .expect("delayed write should be ignored"); + + let metadata = persistence_manager + .load_session_metadata(workspace.path(), &session.session_id) + .await + .expect("metadata should load") + .expect("metadata should exist"); + assert!(metadata.current_context_usage.is_none()); + } + #[tokio::test] async fn execution_binding_rejects_a_session_after_its_first_turn() { let manager = in_memory_test_manager(); diff --git a/src/crates/assembly/core/src/agentic/system.rs b/src/crates/assembly/core/src/agentic/system.rs index 71ae18866..8a34433d1 100644 --- a/src/crates/assembly/core/src/agentic/system.rs +++ b/src/crates/assembly/core/src/agentic/system.rs @@ -74,12 +74,6 @@ pub async fn init_agentic_system_for_profile_with_runtime_ownership( let path_manager = try_get_path_manager_arc()?; let persistence_manager = Arc::new(persistence::PersistenceManager::new(path_manager.clone())?); let token_usage_service = Arc::new(TokenUsageService::new(path_manager.clone()).await?); - let token_usage_subscriber = Arc::new(TokenUsageSubscriber::new(token_usage_service.clone())); - event_router.subscribe_internal("token_usage".to_string(), token_usage_subscriber); - event_router.subscribe_internal( - "thread_goal_tokens".to_string(), - Arc::new(ThreadGoalTokenSubscriber), - ); let context_store = Arc::new(session::SessionContextStore::new()); let context_compressor = Arc::new(session::ContextCompressor::new(Default::default())); @@ -90,6 +84,21 @@ pub async fn init_agentic_system_for_profile_with_runtime_ownership( Default::default(), )); + event_router.subscribe_internal( + "token_usage".to_string(), + Arc::new(TokenUsageSubscriber::new(token_usage_service.clone())), + ); + event_router.subscribe_internal( + "session_context_usage".to_string(), + Arc::new(session::SessionContextUsageSubscriber::new( + session_manager.clone(), + )), + ); + event_router.subscribe_internal( + "thread_goal_tokens".to_string(), + Arc::new(ThreadGoalTokenSubscriber), + ); + let tool_registry = tools::registry::get_global_tool_registry(); let tool_state_manager = Arc::new(tools::pipeline::ToolStateManager::new(event_queue.clone())); let permission_request_manager = diff --git a/src/crates/services/services-core/src/session/lineage.rs b/src/crates/services/services-core/src/session/lineage.rs index a051578a7..3deeed8e2 100644 --- a/src/crates/services/services-core/src/session/lineage.rs +++ b/src/crates/services/services-core/src/session/lineage.rs @@ -398,6 +398,18 @@ pub fn build_branched_session_metadata(facts: BranchSessionMetadataFacts<'_>) -> facts.boundary, facts.branch_lineage, ); + if metadata + .current_context_usage + .as_ref() + .is_some_and(|usage| { + !facts + .branched_turns + .iter() + .any(|turn| turn.turn_id == usage.turn_id) + }) + { + metadata.current_context_usage = None; + } metadata.relationship = None; metadata.todos = None; metadata.review_action_state = None; @@ -480,8 +492,9 @@ fn normalize_nonempty(value: &str) -> Option { mod tests { use super::*; use crate::session::{ - ModelRoundData, SessionMetadata, SessionRelationship, SessionRelationshipKind, - TextItemData, ToolCallData, ToolItemData, UserMessageData, + ModelRoundData, SessionContextUsage, SessionContextUsageSource, SessionMetadata, + SessionRelationship, SessionRelationshipKind, TextItemData, ToolCallData, ToolItemData, + UserMessageData, }; use serde_json::json; @@ -768,6 +781,14 @@ mod tests { "parentDialogTurnId": "legacy-turn", "preserved": "value" })); + source.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-3".to_string(), + input_tokens: 50_000, + output_tokens: Some(500), + total_tokens: 50_500, + timestamp: 40, + source: SessionContextUsageSource::ModelRequest, + }); source.relationship = Some(SessionRelationship { kind: Some(SessionRelationshipKind::Subagent), parent_session_id: Some("parent".to_string()), @@ -818,6 +839,7 @@ mod tests { assert!(branched.deep_review_run_manifest.is_none()); assert!(branched.unread_completion.is_none()); assert!(branched.needs_user_attention.is_none()); + assert!(branched.current_context_usage.is_none()); let custom_metadata = branched .custom_metadata @@ -835,6 +857,40 @@ mod tests { ); } + #[test] + fn build_branched_session_metadata_keeps_usage_for_a_copied_turn() { + let mut source = metadata("source"); + source.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-2".to_string(), + input_tokens: 2_000, + output_tokens: Some(200), + total_tokens: 2_200, + timestamp: 40, + source: SessionContextUsageSource::ModelRequest, + }); + let turns = vec![turn("target", "turn-1", 0), turn("target", "turn-2", 1)]; + let branch_lineage = BranchSessionLineage { + base_session_name: "Source".to_string(), + ordinal: 1, + }; + + let branched = build_branched_session_metadata(BranchSessionMetadataFacts { + source_metadata: &source, + target_session_id: "target".to_string(), + target_session_name: "Target".to_string(), + target_agent_type: "agentic".to_string(), + source_session_id: "source", + source_turn_id: "turn-2", + source_turn_index: 1, + boundary: SessionBranchBoundary::ThroughTurn, + branched_turns: &turns, + branch_lineage: &branch_lineage, + now_ms: 42, + }); + + assert_eq!(branched.current_context_usage, source.current_context_usage); + } + #[test] fn branch_lineage_uses_the_inherited_title_namespace_for_renamed_suffixes() { let mut root = metadata("root"); diff --git a/src/crates/services/services-core/src/session/metadata.rs b/src/crates/services/services-core/src/session/metadata.rs index e7b9f6bbe..69bfa7ea1 100644 --- a/src/crates/services/services-core/src/session/metadata.rs +++ b/src/crates/services/services-core/src/session/metadata.rs @@ -71,6 +71,7 @@ pub fn build_session_metadata(facts: SessionMetadataBuildFacts<'_>) -> SessionMe .or_else(|| existing.and_then(|value| value.snapshot_session_id.clone())), tags: existing.map(|value| value.tags.clone()).unwrap_or_default(), custom_metadata: existing.and_then(|value| value.custom_metadata.clone()), + current_context_usage: existing.and_then(|value| value.current_context_usage.clone()), relationship: build_session_relationship(facts.session_kind, existing), todos: existing.and_then(|value| value.todos.clone()), review_action_state: existing.and_then(|value| value.review_action_state.clone()), @@ -227,6 +228,13 @@ pub fn refresh_session_metadata_from_turns( metadata.message_count = turns.iter().map(estimate_turn_message_count).sum(); metadata.tool_call_count = turns.iter().map(DialogTurnData::count_tool_calls).sum(); metadata.last_finished_at = turns.iter().filter_map(dialog_turn_finished_at).max(); + if metadata + .current_context_usage + .as_ref() + .is_some_and(|usage| !turns.iter().any(|turn| turn.turn_id == usage.turn_id)) + { + metadata.current_context_usage = None; + } metadata.last_active_at = last_active_at; fill_workspace_path_if_missing(metadata, workspace_path); } diff --git a/src/crates/services/services-core/src/session/types.rs b/src/crates/services/services-core/src/session/types.rs index 0d4029e00..7da2e9839 100644 --- a/src/crates/services/services-core/src/session/types.rs +++ b/src/crates/services/services-core/src/session/types.rs @@ -182,6 +182,15 @@ pub struct SessionMetadata { #[serde(skip_serializing_if = "Option::is_none", alias = "custom_metadata")] pub custom_metadata: Option, + /// Latest authoritative context-usage value for restoring context display + /// state across hosts and process restarts. + #[serde( + default, + skip_serializing_if = "Option::is_none", + alias = "current_context_usage" + )] + pub current_context_usage: Option, + /// Structured child-session relationship metadata. #[serde( default, @@ -551,6 +560,44 @@ pub struct DialogTurnTokenUsageData { pub timestamp: u64, } +/// Source of a persisted session context-usage value. +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum SessionContextUsageSource { + ModelRequest, + ContextCompression, +} + +/// Exact context-usage value owned by the Agent Session runtime. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "camelCase")] +pub struct SessionContextUsage { + /// Dialog turn that produced this value. + #[serde(alias = "turn_id")] + pub turn_id: String, + + /// Input/prompt tokens for a model request, or the compacted context size. + #[serde(alias = "input_tokens")] + pub input_tokens: u64, + + /// Output/completion tokens when this value came from a model request. + #[serde( + default, + skip_serializing_if = "Option::is_none", + alias = "output_tokens" + )] + pub output_tokens: Option, + + /// Provider total for a model request, or the compacted context size. + #[serde(alias = "total_tokens")] + pub total_tokens: u64, + + /// Runtime event timestamp in milliseconds since epoch. + pub timestamp: u64, + + pub source: SessionContextUsageSource, +} + /// Persisted dialog turn kind. #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] @@ -996,6 +1043,7 @@ impl SessionMetadata { snapshot_session_id: None, tags: Vec::new(), custom_metadata: None, + current_context_usage: None, relationship: None, todos: None, review_action_state: None, diff --git a/src/crates/services/services-core/tests/session_metadata_contracts.rs b/src/crates/services/services-core/tests/session_metadata_contracts.rs index e633495c9..629a83167 100644 --- a/src/crates/services/services-core/tests/session_metadata_contracts.rs +++ b/src/crates/services/services-core/tests/session_metadata_contracts.rs @@ -3,8 +3,9 @@ use bitfun_services_core::session::{ build_session_index_snapshot, refresh_session_metadata_from_turns, remove_session_index_entry, try_refresh_session_metadata_for_saved_turn, upsert_session_index_entry, DialogTurnData, - DialogTurnKind, ModelRoundData, SessionKind, SessionMetadata, StoredSessionIndexFile, - TextItemData, ToolCallData, ToolItemData, TurnStatus, UserMessageData, + DialogTurnKind, ModelRoundData, SessionContextUsage, SessionContextUsageSource, SessionKind, + SessionMetadata, StoredSessionIndexFile, TextItemData, ToolCallData, ToolItemData, TurnStatus, + UserMessageData, }; fn metadata(session_id: &str) -> SessionMetadata { @@ -172,6 +173,73 @@ fn full_refresh_recomputes_metadata_counters_from_turns() { ); } +#[test] +fn session_context_usage_round_trips_as_top_level_metadata() { + let mut metadata = metadata("session-1"); + metadata.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-3".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }); + + let value = serde_json::to_value(&metadata).expect("serialize metadata"); + assert_eq!(value["currentContextUsage"]["turnId"], "turn-3"); + assert_eq!(value["currentContextUsage"]["source"], "model_request"); + + let restored: SessionMetadata = serde_json::from_value(value).expect("deserialize metadata"); + assert_eq!( + restored.current_context_usage, + metadata.current_context_usage + ); +} + +#[test] +fn full_refresh_drops_context_usage_for_a_removed_turn() { + let mut metadata = metadata("session-1"); + metadata.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-1".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }); + + refresh_session_metadata_from_turns( + &mut metadata, + "D:/workspace/project", + &[turn("session-1", 0, 1, 0)], + 42, + ); + + assert!(metadata.current_context_usage.is_none()); +} + +#[test] +fn full_refresh_keeps_context_usage_for_a_surviving_turn() { + let mut metadata = metadata("session-1"); + metadata.current_context_usage = Some(SessionContextUsage { + turn_id: "turn-0".to_string(), + input_tokens: 42_000, + output_tokens: Some(1_500), + total_tokens: 43_500, + timestamp: 123, + source: SessionContextUsageSource::ModelRequest, + }); + + refresh_session_metadata_from_turns( + &mut metadata, + "D:/workspace/project", + &[turn("session-1", 0, 1, 0)], + 42, + ); + + assert!(metadata.current_context_usage.is_some()); +} + #[test] fn saved_turn_refresh_updates_incrementally_for_append_and_replace() { let mut metadata = metadata("session-1"); diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.test.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.test.ts index ecef4c457..5462ae2c1 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.test.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.test.ts @@ -15,7 +15,7 @@ import type { DialogTurn, FlowToolItem, FlowUserSteeringItem, ModelRound, Sessio import type { FlowChatContext } from './types'; import { markOptimisticDispatchTurnMetadata } from '@/features/dispatch/optimisticDispatchTurn'; -const { handleCompressionCompleted } = __test_only__; +const { handleCompressionCompleted, handleTokenUsageUpdate } = __test_only__; vi.mock('../../../shared/notification-system/services/NotificationService', () => ({ notificationService: { @@ -1166,6 +1166,8 @@ describe('handleCompressionCompleted', () => { outputTokens: undefined, totalTokens: 15_000, timestamp: expect.any(Number), + turnId: 'turn-1', + source: 'context_compression', }); }); @@ -1201,4 +1203,78 @@ describe('handleCompressionCompleted', () => { const session = FlowChatStore.getInstance().getState().sessions.get('session-1'); expect(session?.currentTokenUsage).toBeUndefined(); }); + + it('ignores a delayed successful compression after its source turn was removed', () => { + putFinishingSessionInStore(); + FlowChatStore.getInstance().deleteDialogTurn('session-1', 'turn-1'); + + handleCompressionCompleted(createFlowChatContext(), { + sessionId: 'session-1', + turnId: 'turn-1', + compressionId: 'compression-1', + applied: true, + tokensBefore: 90_000, + tokensAfter: 15_000, + }); + + const session = FlowChatStore.getInstance().getState().sessions.get('session-1'); + expect(session?.currentTokenUsage).toBeUndefined(); + }); +}); + +describe('handleTokenUsageUpdate', () => { + beforeEach(() => { + vi.restoreAllMocks(); + resetFlowChatStore(); + stateMachineManager.clear(); + }); + + afterEach(() => { + resetFlowChatStore(); + stateMachineManager.clear(); + }); + + it('tracks the source turn on current usage without adding provenance to accumulated turn usage', () => { + putFinishingSessionInStore(); + + handleTokenUsageUpdate(createFlowChatContext(), { + sessionId: 'session-1', + turnId: 'turn-1', + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + }); + + const session = FlowChatStore.getInstance().getState().sessions.get('session-1'); + expect(session?.currentTokenUsage).toMatchObject({ + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + turnId: 'turn-1', + source: 'model_request', + }); + expect(session?.dialogTurns[0].tokenUsage).toMatchObject({ + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + }); + expect(session?.dialogTurns[0].tokenUsage).not.toHaveProperty('turnId'); + expect(session?.dialogTurns[0].tokenUsage).not.toHaveProperty('source'); + }); + + it('ignores a delayed model usage update after its source turn was removed', () => { + putFinishingSessionInStore(); + FlowChatStore.getInstance().deleteDialogTurn('session-1', 'turn-1'); + + handleTokenUsageUpdate(createFlowChatContext(), { + sessionId: 'session-1', + turnId: 'turn-1', + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + }); + + const session = FlowChatStore.getInstance().getState().sessions.get('session-1'); + expect(session?.currentTokenUsage).toBeUndefined(); + }); }); diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.ts index 12e7ca676..f39592b03 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/EventHandlerModule.ts @@ -57,7 +57,6 @@ const pendingImageAnalysisTurns = new Map(); import { debouncedSaveDialogTurn, immediateSaveDialogTurn, - persistLastRequestTokenUsage, saveDialogTurnToDisk, cleanupSaveState, } from './PersistenceModule'; @@ -164,6 +163,7 @@ export const __test_only__ = { handleDialogTurnFailed, handleSubagentSessionLinked, handleModelRoundStart, + handleTokenUsageUpdate, handleCompressionCompleted, }; @@ -2127,6 +2127,13 @@ function handleTokenUsageUpdate(context: FlowChatContext, event: any): void { log.debug('Session not found (token usage update)', { sessionId }); return; } + if ( + typeof turnId !== 'string' + || !session.dialogTurns.some(turn => turn.id === turnId) + ) { + log.debug('Dropped token usage update for non-visible turn', { sessionId, turnId }); + return; + } if (typeof inputTokens !== 'number' || typeof totalTokens !== 'number') { log.debug('Dropped invalid token usage update', { event }); return; @@ -2135,20 +2142,11 @@ function handleTokenUsageUpdate(context: FlowChatContext, event: any): void { store.updateTokenUsage(sessionId, { inputTokens, outputTokens: typeof outputTokens === 'number' ? outputTokens : undefined, - totalTokens + totalTokens, + turnId, + source: 'model_request', }, turnId); - // Persist the exact last request usage so the context display survives a - // restart. Skip ACP sessions: their display is driven by - // currentAcpContextUsage instead. - if (!session.mode?.startsWith('acp:') && !session.config.agentType?.startsWith('acp:')) { - persistLastRequestTokenUsage(context, sessionId, { - inputTokens, - outputTokens: typeof outputTokens === 'number' ? outputTokens : undefined, - totalTokens, - }); - } - if (maxContextTokens !== undefined && maxContextTokens !== null) { store.updateSessionMaxContextTokens(sessionId, maxContextTokens); } @@ -2264,6 +2262,14 @@ function handleCompressionCompleted(context: FlowChatContext, event: any): void }); const store = FlowChatStore.getInstance(); + const session = store.getState().sessions.get(sessionId); + if ( + typeof turnId !== 'string' + || !session?.dialogTurns.some(turn => turn.id === turnId) + ) { + log.debug('Dropped compression completion for non-visible turn', { sessionId, turnId }); + return; + } store.updateModelRoundItem(sessionId, turnId, compressionId, { toolResult: { @@ -2294,6 +2300,8 @@ function handleCompressionCompleted(context: FlowChatContext, event: any): void inputTokens: tokensAfter, totalTokens: tokensAfter, outputTokens: undefined, + turnId, + source: 'context_compression', }); } diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/PersistenceModule.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/PersistenceModule.ts index daac1d78e..fc6bffd29 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/PersistenceModule.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/PersistenceModule.ts @@ -604,83 +604,3 @@ export async function touchSessionActivity( log.debug('Failed to touch session activity', { sessionId, error }); } } - -const lastRequestTokenUsageDebouncers = new Map< - string, - ReturnType ->(); - -/** - * Persist the exact last model request usage into session metadata so the - * input-box context display can be restored after an app restart. - * - * Only the session-level last request value is stored; the dialog turn usage - * stays accumulated per turn and must not be reused as a single-request - * approximation. The write is trailing-throttled because agentic sessions - * can emit one TokenUsageUpdated per model round. - */ -export function persistLastRequestTokenUsage( - context: FlowChatContext, - sessionId: string, - usage: { inputTokens: number; outputTokens?: number; totalTokens: number }, -): void { - const existingTimer = lastRequestTokenUsageDebouncers.get(sessionId); - if (existingTimer) { - clearTimeout(existingTimer); - } - const timer = setTimeout(() => { - lastRequestTokenUsageDebouncers.delete(sessionId); - void persistLastRequestTokenUsageNow(context, sessionId, usage).catch(error => { - log.warn('Failed to persist last request token usage', { sessionId, error }); - }); - }, COALESCED_IMMEDIATE_SAVE_DELAY_MS); - lastRequestTokenUsageDebouncers.set(sessionId, timer); -} - -async function persistLastRequestTokenUsageNow( - context: FlowChatContext, - sessionId: string, - usage: { inputTokens: number; outputTokens?: number; totalTokens: number }, -): Promise { - const { sessionAPI } = await import('@/infrastructure/api/service-api/SessionAPI'); - - const session = context.flowChatStore.getState().sessions.get(sessionId); - if (!session) return; - if (isTransientSession(session) || isObserverOnlyDispatchSession(sessionId, session)) return; - - const workspacePath = requireSessionProjectWorkspacePath(session, sessionId); - - let existingMetadata: any = null; - try { - existingMetadata = await sessionAPI.loadSessionMetadata( - sessionId, - workspacePath, - session.remoteConnectionId, - session.remoteSshHost - ); - } catch { - // Metadata may not exist yet for a fresh session; the patch below still works. - } - - const metadata = { - ...existingMetadata, - sessionId, - customMetadata: { - ...(existingMetadata?.customMetadata ?? {}), - lastRequestTokenUsage: { - inputTokens: usage.inputTokens, - outputTokens: usage.outputTokens, - totalTokens: usage.totalTokens, - timestamp: Date.now(), - }, - }, - }; - - await sessionAPI.saveSessionMetadata( - metadata, - workspacePath, - ['titleMetadata'], - session.remoteConnectionId, - session.remoteSshHost - ); -} diff --git a/src/web-ui/src/flow_chat/store/FlowChatStore.test.ts b/src/web-ui/src/flow_chat/store/FlowChatStore.test.ts index 1e45a2ff5..e75c0e6ea 100644 --- a/src/web-ui/src/flow_chat/store/FlowChatStore.test.ts +++ b/src/web-ui/src/flow_chat/store/FlowChatStore.test.ts @@ -683,6 +683,140 @@ describe('FlowChatStore token usage', () => { totalTokens: 425, }); }); + + it('falls back safely when deleting the turn that sourced current usage', () => { + const previousTurn = { + id: 'turn-1', + sessionId: 'session-1', + userMessage: { id: 'user-1', content: 'first', timestamp: 1_000 }, + modelRounds: [{ id: 'round-1' }], + tokenUsage: { + inputTokens: 600, + outputTokens: 100, + totalTokens: 700, + timestamp: 1_500, + }, + status: 'completed' as const, + startTime: 1_000, + }; + const sourceTurn = { + id: 'turn-2', + sessionId: 'session-1', + userMessage: { id: 'user-2', content: 'second', timestamp: 2_000 }, + modelRounds: [{ id: 'round-2' }], + status: 'completed' as const, + startTime: 2_000, + }; + const session = createSession({ + dialogTurns: [previousTurn, sourceTurn], + currentTokenUsage: { + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + timestamp: 2_500, + turnId: 'turn-2', + source: 'model_request', + }, + }); + flowChatStore.setState(() => ({ + sessions: new Map([[session.sessionId, session]]), + activeSessionId: session.sessionId, + })); + + flowChatStore.deleteDialogTurn(session.sessionId, 'turn-2'); + + expect(flowChatStore.getState().sessions.get(session.sessionId)?.currentTokenUsage).toEqual({ + ...previousTurn.tokenUsage, + turnId: 'turn-1', + }); + }); + + it('does not derive a stale fallback from partial history after deleting the usage source', () => { + const previousTurn = { + id: 'turn-1', + sessionId: 'session-1', + userMessage: { id: 'user-1', content: 'partial older turn', timestamp: 1_000 }, + modelRounds: [{ id: 'round-1' }], + tokenUsage: { + inputTokens: 600, + outputTokens: 100, + totalTokens: 700, + timestamp: 1_500, + }, + status: 'completed' as const, + startTime: 1_000, + }; + const sourceTurn = { + id: 'turn-2', + sessionId: 'session-1', + userMessage: { id: 'user-2', content: 'source', timestamp: 2_000 }, + modelRounds: [{ id: 'round-2' }], + status: 'completed' as const, + startTime: 2_000, + }; + const session = createSession({ + dialogTurns: [previousTurn, sourceTurn], + isPartial: true, + currentTokenUsage: { + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + timestamp: 2_500, + turnId: 'turn-2', + source: 'model_request', + }, + }); + flowChatStore.setState(() => ({ + sessions: new Map([[session.sessionId, session]]), + activeSessionId: session.sessionId, + })); + + flowChatStore.deleteDialogTurn(session.sessionId, 'turn-2'); + + expect( + flowChatStore.getState().sessions.get(session.sessionId)?.currentTokenUsage, + ).toBeUndefined(); + }); + + it('clears usage when truncation removes its source and no safe fallback exists', () => { + const retainedTurn = { + id: 'turn-1', + sessionId: 'session-1', + userMessage: { id: 'user-1', content: 'first', timestamp: 1_000 }, + modelRounds: [], + status: 'completed' as const, + startTime: 1_000, + }; + const sourceTurn = { + id: 'turn-2', + sessionId: 'session-1', + userMessage: { id: 'user-2', content: 'second', timestamp: 2_000 }, + modelRounds: [{ id: 'round-2' }], + status: 'completed' as const, + startTime: 2_000, + }; + const session = createSession({ + dialogTurns: [retainedTurn, sourceTurn], + currentTokenUsage: { + inputTokens: 1_200, + outputTokens: 320, + totalTokens: 1_520, + timestamp: 2_500, + turnId: 'turn-2', + source: 'model_request', + }, + }); + flowChatStore.setState(() => ({ + sessions: new Map([[session.sessionId, session]]), + activeSessionId: session.sessionId, + })); + + flowChatStore.truncateDialogTurnsFrom(session.sessionId, 1); + + expect( + flowChatStore.getState().sessions.get(session.sessionId)?.currentTokenUsage, + ).toBeUndefined(); + }); }); describe('FlowChatStore round attempts', () => { @@ -1847,6 +1981,224 @@ describe('FlowChatStore historical session hydration state', () => { }); }); + it('clears Peer usage whose source turn is absent even when the running snapshot is unchanged', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'agentic', + state: 'Processing { current_turn_id: "turn-live", phase: Streaming }', + turnCount: 1, + createdAt: 1, + }, + turns: [{ + turnId: 'turn-live', + turnIndex: 0, + sessionId: 'history-1', + timestamp: 1, + userMessage: { id: 'user-live', content: 'continue', timestamp: 1 }, + modelRounds: [{ + id: 'round-live', + turnId: 'turn-live', + roundIndex: 0, + timestamp: 1, + textItems: [{ + id: 'host-text-id', + content: 'partial answer', + isStreaming: true, + timestamp: 2, + status: 'streaming', + }], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'streaming', + }], + startTime: 1, + status: 'inprogress', + }], + contextRestoreState: 'pending', + isPartial: false, + loadedTurnCount: 1, + totalTurnCount: 1, + }); + const localTurn = { + id: 'turn-live', + sessionId: 'history-1', + userMessage: { id: 'user-live', content: 'continue', timestamp: 1 }, + modelRounds: [{ + id: 'round-live', + index: 0, + items: [{ + id: 'controller-text-id', + type: 'text' as const, + content: 'partial answer plus live data', + isStreaming: true, + isMarkdown: true, + timestamp: 3, + status: 'streaming' as const, + }], + isStreaming: true, + isComplete: false, + status: 'streaming' as const, + startTime: 1, + }], + status: 'processing' as const, + startTime: 1, + backendTurnIndex: 0, + }; + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + historyState: 'ready', + dialogTurns: [localTurn], + currentTokenUsage: { + inputTokens: 42_000, + outputTokens: 1_000, + totalTokens: 43_000, + timestamp: 3, + turnId: 'turn-no-longer-visible', + source: 'model_request', + }, + })], + ]), + activeSessionId: 'history-1', + })); + + const result = await flowChatStore.refreshPeerSessionSnapshot( + 'history-1', + '/Users/host/project', + { replaceRunningSnapshot: false }, + ); + + expect(result.applied).toBe(true); + const refreshedSession = flowChatStore.getState().sessions.get('history-1'); + expect(refreshedSession?.currentTokenUsage).toBeUndefined(); + expect(refreshedSession?.dialogTurns[0].modelRounds[0].items[0]).toMatchObject({ + id: 'controller-text-id', + content: 'partial answer plus live data', + }); + }); + + it('replaces stale local usage with authoritative Peer usage for a multi-round turn', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'agentic', + state: 'Idle', + turnCount: 1, + createdAt: 1, + }, + turns: [{ + turnId: 'turn-live', + turnIndex: 0, + sessionId: 'history-1', + timestamp: 1, + userMessage: { id: 'user-live', content: 'continue', timestamp: 1 }, + modelRounds: [ + { + id: 'round-1', + turnId: 'turn-live', + roundIndex: 0, + timestamp: 1, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'completed', + }, + { + id: 'round-2', + turnId: 'turn-live', + roundIndex: 1, + timestamp: 2, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 2, + status: 'completed', + }, + ], + tokenUsage: { + inputTokens: 90_000, + outputTokens: 2_000, + totalTokens: 92_000, + timestamp: 3, + }, + startTime: 1, + endTime: 3, + status: 'completed', + }], + currentContextUsage: { + inputTokens: 42_000, + outputTokens: 1_500, + totalTokens: 43_500, + timestamp: 4, + turnId: 'turn-live', + source: 'model_request', + }, + contextRestoreState: 'ready', + isPartial: false, + loadedTurnCount: 1, + totalTurnCount: 1, + }); + const localTurn = { + id: 'turn-live', + sessionId: 'history-1', + userMessage: { id: 'user-live', content: 'continue', timestamp: 1 }, + modelRounds: [{ + id: 'round-1', + index: 0, + items: [], + isStreaming: false, + isComplete: true, + status: 'completed' as const, + startTime: 1, + }], + status: 'completed' as const, + startTime: 1, + backendTurnIndex: 0, + }; + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + historyState: 'ready', + dialogTurns: [localTurn], + currentTokenUsage: { + inputTokens: 12_000, + outputTokens: 500, + totalTokens: 12_500, + timestamp: 2, + turnId: 'turn-live', + source: 'model_request', + }, + })], + ]), + activeSessionId: 'history-1', + })); + + const result = await flowChatStore.refreshPeerSessionSnapshot( + 'history-1', + '/Users/host/project', + { replaceRunningSnapshot: false }, + ); + + expect(result.applied).toBe(true); + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toEqual({ + inputTokens: 42_000, + outputTokens: 1_500, + totalTokens: 43_500, + timestamp: 4, + turnId: 'turn-live', + source: 'model_request', + }); + }); + it('replaces a stale running projection after the Peer Host has completed', async () => { peerModeFlagMock.active = true; apiMocks.restoreSessionView.mockResolvedValueOnce({ @@ -5305,6 +5657,147 @@ describe('FlowChatStore historical session hydration state', () => { }); }); + it('does not backfill through the latest terminal turn when its usage spans multiple rounds', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'agentic', + state: 'Idle', + turnCount: 2, + createdAt: 1, + }, + turns: [ + { + ...createPersistedTurn(0), + modelRounds: [{ + id: 'round-0', + turnId: 'turn-0', + roundIndex: 0, + timestamp: 1, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'completed', + }], + endTime: 2, + tokenUsage: { + inputTokens: 1000, + outputTokens: 100, + totalTokens: 1100, + timestamp: 2, + }, + }, + { + ...createPersistedTurn(1), + modelRounds: [ + { + id: 'round-1', + turnId: 'turn-1', + roundIndex: 0, + timestamp: 3, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 3, + status: 'completed', + }, + { + id: 'round-2', + turnId: 'turn-1', + roundIndex: 1, + timestamp: 4, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 4, + status: 'completed', + }, + ], + endTime: 5, + tokenUsage: { + inputTokens: 8_900_000, + outputTokens: 300, + totalTokens: 8_900_300, + timestamp: 5, + }, + }, + ], + contextRestoreState: 'ready', + }); + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + isHistorical: true, + historyState: 'metadata-only', + })], + ]), + activeSessionId: 'history-1', + })); + + await flowChatStore.loadSessionHistory('history-1', 'D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toBeUndefined(); + }); + + it('uses the restored agent type to suppress native usage for ACP hydration', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'acp:test', + state: 'Idle', + turnCount: 1, + createdAt: 1, + }, + turns: [{ + ...createPersistedTurn(0), + modelRounds: [{ + id: 'round-0', + turnId: 'turn-0', + roundIndex: 0, + timestamp: 1, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'completed', + }], + endTime: 2, + tokenUsage: { + inputTokens: 2400, + outputTokens: 300, + totalTokens: 2700, + timestamp: 2, + }, + }], + contextRestoreState: 'ready', + }); + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + isHistorical: true, + historyState: 'metadata-only', + mode: 'agentic', + config: { agentType: 'agentic' }, + })], + ]), + activeSessionId: 'history-1', + })); + + await flowChatStore.loadSessionHistory('history-1', 'D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')).toMatchObject({ + mode: 'acp:test', + currentTokenUsage: undefined, + }); + }); + it('keeps an existing currentTokenUsage when hydrating historical turns', async () => { peerModeFlagMock.active = true; apiMocks.restoreSessionView.mockResolvedValueOnce({ @@ -5352,6 +5845,8 @@ describe('FlowChatStore historical session hydration state', () => { outputTokens: 1, totalTokens: 1000, timestamp: 5, + turnId: 'turn-0', + source: 'model_request', }, })], ]), @@ -5364,10 +5859,138 @@ describe('FlowChatStore historical session hydration state', () => { inputTokens: 999, outputTokens: 1, totalTokens: 1000, + turnId: 'turn-0', + source: 'model_request', }); }); - it('restores the exact last request token usage from persisted metadata', async () => { + it('discards stale exact usage but keeps a safe fallback when restore reports no exact usage', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'agentic', + state: 'Idle', + turnCount: 1, + createdAt: 1, + }, + turns: [{ + ...createPersistedTurn(0), + modelRounds: [{ + id: 'round-0', + turnId: 'turn-0', + roundIndex: 0, + timestamp: 1, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'completed', + }], + endTime: 2, + tokenUsage: { + inputTokens: 2400, + outputTokens: 300, + totalTokens: 2700, + timestamp: 2, + }, + }], + currentContextUsage: null, + contextRestoreState: 'ready', + }); + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + isHistorical: true, + historyState: 'metadata-only', + currentTokenUsage: { + inputTokens: 999, + outputTokens: 1, + totalTokens: 1000, + timestamp: 5, + turnId: 'turn-0', + source: 'model_request', + }, + })], + ]), + activeSessionId: 'history-1', + })); + + await flowChatStore.loadSessionHistory('history-1', 'D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toMatchObject({ + inputTokens: 2400, + outputTokens: 300, + totalTokens: 2700, + turnId: 'turn-0', + }); + }); + + it('invalidates restored context usage when its source turn is not visible after hydration', async () => { + peerModeFlagMock.active = true; + apiMocks.restoreSessionView.mockResolvedValueOnce({ + session: { + sessionId: 'history-1', + sessionName: 'History 1', + agentType: 'agentic', + state: 'Idle', + turnCount: 1, + createdAt: 1, + }, + turns: [{ + ...createPersistedTurn(0), + modelRounds: [{ + id: 'round-0', + turnId: 'turn-0', + roundIndex: 0, + timestamp: 1, + textItems: [], + toolItems: [], + thinkingItems: [], + startTime: 1, + status: 'completed', + }], + endTime: 2, + tokenUsage: { + inputTokens: 2400, + outputTokens: 300, + totalTokens: 2700, + timestamp: 2, + }, + }], + contextRestoreState: 'ready', + }); + flowChatStore.setState(() => ({ + sessions: new Map([ + ['history-1', createSession({ + sessionId: 'history-1', + isHistorical: true, + historyState: 'metadata-only', + currentTokenUsage: { + inputTokens: 42000, + outputTokens: 1500, + totalTokens: 43500, + timestamp: 5, + turnId: 'deleted-turn', + source: 'model_request', + }, + })], + ]), + activeSessionId: 'history-1', + })); + + await flowChatStore.loadSessionHistory('history-1', 'D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toMatchObject({ + inputTokens: 2400, + outputTokens: 300, + totalTokens: 2700, + }); + }); + + it('restores current token usage from top-level persisted context metadata', async () => { apiMocks.listSessions.mockResolvedValueOnce([ { sessionId: 'history-1', @@ -5376,13 +5999,13 @@ describe('FlowChatStore historical session hydration state', () => { modelName: 'auto', createdAt: 10, lastActiveAt: 20, - customMetadata: { - lastRequestTokenUsage: { - inputTokens: 42000, - outputTokens: 1500, - totalTokens: 43500, - timestamp: 21, - }, + currentContextUsage: { + inputTokens: 42000, + outputTokens: 1500, + totalTokens: 43500, + timestamp: 21, + turnId: 'turn-7', + source: 'model_request', }, }, ]); @@ -5393,10 +6016,37 @@ describe('FlowChatStore historical session hydration state', () => { inputTokens: 42000, outputTokens: 1500, totalTokens: 43500, + turnId: 'turn-7', + source: 'model_request', }); }); - it('ignores invalid persisted last request token usage', async () => { + it('does not restore native context metadata for ACP sessions', async () => { + apiMocks.listSessions.mockResolvedValueOnce([ + { + sessionId: 'history-1', + title: 'Saved ACP session', + agentType: 'acp:test', + modelName: 'auto', + createdAt: 10, + lastActiveAt: 20, + currentContextUsage: { + inputTokens: 42000, + outputTokens: 1500, + totalTokens: 43500, + timestamp: 21, + turnId: 'turn-7', + source: 'model_request', + }, + }, + ]); + + await flowChatStore.initializeFromDisk('D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toBeUndefined(); + }); + + it('ignores invalid persisted current context usage', async () => { apiMocks.listSessions.mockResolvedValueOnce([ { sessionId: 'history-1', @@ -5405,13 +6055,38 @@ describe('FlowChatStore historical session hydration state', () => { modelName: 'auto', createdAt: 10, lastActiveAt: 20, - customMetadata: { - lastRequestTokenUsage: { - inputTokens: 0, - outputTokens: 0, - totalTokens: 0, - timestamp: 21, - }, + currentContextUsage: { + inputTokens: 0, + outputTokens: 0, + totalTokens: 0, + timestamp: 21, + turnId: 'turn-7', + source: 'model_request', + }, + }, + ]); + + await flowChatStore.initializeFromDisk('D:/workspace/BitFun'); + + expect(flowChatStore.getState().sessions.get('history-1')?.currentTokenUsage).toBeUndefined(); + }); + + it('ignores persisted context usage without valid provenance', async () => { + apiMocks.listSessions.mockResolvedValueOnce([ + { + sessionId: 'history-1', + title: 'Saved session', + agentType: 'agentic', + modelName: 'auto', + createdAt: 10, + lastActiveAt: 20, + currentContextUsage: { + inputTokens: 42000, + outputTokens: 1500, + totalTokens: 43500, + timestamp: 21, + turnId: ' ', + source: 'unknown_source', }, }, ]); diff --git a/src/web-ui/src/flow_chat/store/FlowChatStore.ts b/src/web-ui/src/flow_chat/store/FlowChatStore.ts index bf3f9108b..05be069b2 100644 --- a/src/web-ui/src/flow_chat/store/FlowChatStore.ts +++ b/src/web-ui/src/flow_chat/store/FlowChatStore.ts @@ -39,6 +39,7 @@ import { i18nService } from '@/infrastructure/i18n/core/I18nService'; import type { DialogTurnData, LocalCommandMetadata, + SessionContextUsage, SessionKind, SessionTurnCatalog, } from '@/shared/types/session-history'; @@ -57,6 +58,7 @@ import { import { sessionProjectWorkspacePath } from '../utils/sessionWorkspace'; import type { SessionTitleDescriptor } from '../utils/sessionTitle'; import { deriveContextUsageFromTurns } from '../utils/tokenUsageDisplay'; +import { isAcpAgentType } from '../utils/acpSession'; import { deriveSessionTitleState, deriveSessionTitleStateFromMetadata, @@ -107,11 +109,10 @@ function firstNonEmptyString(...values: unknown[]): string | undefined { return undefined; } -function isAcpSessionForContextUsage(session: Session): boolean { - return Boolean( - session.mode?.startsWith('acp:') - || session.config.agentType?.startsWith('acp:'), - ); +function persistedCurrentContextUsageValue( + metadata: { currentContextUsage?: SessionContextUsage }, +): unknown { + return metadata.currentContextUsage; } function deriveRestoredCurrentTokenUsage(value: unknown): TokenUsage | undefined { @@ -124,14 +125,98 @@ function deriveRestoredCurrentTokenUsage(value: unknown): TokenUsage | undefined return undefined; } const totalTokens = record.totalTokens; + const turnId = typeof record.turnId === 'string' && record.turnId.trim() + ? record.turnId.trim() + : undefined; + const source = record.source === 'model_request' || record.source === 'context_compression' + ? record.source + : undefined; + const outputTokens = record.outputTokens; + const timestamp = record.timestamp; + if ( + !turnId + || !source + || typeof totalTokens !== 'number' + || !Number.isFinite(totalTokens) + || totalTokens < 0 + || ( + outputTokens !== undefined + && ( + typeof outputTokens !== 'number' + || !Number.isFinite(outputTokens) + || outputTokens < 0 + ) + ) + || typeof timestamp !== 'number' + || !Number.isFinite(timestamp) + ) { + return undefined; + } return { inputTokens, - outputTokens: typeof record.outputTokens === 'number' ? record.outputTokens : undefined, - totalTokens: typeof totalTokens === 'number' && Number.isFinite(totalTokens) ? totalTokens : inputTokens, - timestamp: typeof record.timestamp === 'number' ? record.timestamp : Date.now(), + outputTokens, + totalTokens, + timestamp, + turnId, + source, }; } +function reconcileHydratedCurrentTokenUsage( + currentTokenUsage: TokenUsage | undefined, + dialogTurns: DialogTurn[], + restoredAgentType: string | undefined, + sourceVisibilityTurns: DialogTurn[] = dialogTurns, +): TokenUsage | undefined { + if (isAcpAgentType(restoredAgentType)) { + return undefined; + } + + const sourceTurnId = currentTokenUsage?.turnId; + const retainedUsage = sourceTurnId + && !sourceVisibilityTurns.some(turn => turn.id === sourceTurnId) + ? undefined + : currentTokenUsage; + + return retainedUsage ?? deriveContextUsageFromTurns(dialogTurns); +} + +function reconcileRestoreViewCurrentTokenUsage( + currentTokenUsage: TokenUsage | undefined, + authoritativeUsage: SessionContextUsage | null | undefined, + dialogTurns: DialogTurn[], + restoredAgentType: string | undefined, + sourceVisibilityTurns: DialogTurn[] = dialogTurns, +): TokenUsage | undefined { + if (isAcpAgentType(restoredAgentType)) { + return undefined; + } + const candidateUsage = authoritativeUsage === undefined + ? currentTokenUsage + : authoritativeUsage === null + ? undefined + : deriveRestoredCurrentTokenUsage(authoritativeUsage); + return reconcileHydratedCurrentTokenUsage( + candidateUsage, + dialogTurns, + restoredAgentType, + sourceVisibilityTurns, + ); +} + +function currentTokenUsageAfterSourceRemoval( + session: Pick, + dialogTurns: DialogTurn[], + sourceRemoved: boolean, +): TokenUsage | undefined { + if (!sourceRemoved) { + return session.currentTokenUsage; + } + return session.isPartial === true + ? undefined + : deriveContextUsageFromTurns(dialogTurns); +} + function persistedSessionRemoteScope( metadata: { remoteConnectionId?: unknown; @@ -4968,10 +5053,16 @@ export class FlowChatStore { && !isProvisionalUsageReportTurn(deletedTurn); shiftLaterOrdinals = countedOptimisticTurn; const nextCanonicalTurnCount = canonicalSessionTurns({ dialogTurns: updatedDialogTurns }).length; + const currentTokenUsage = currentTokenUsageAfterSourceRemoval( + session, + updatedDialogTurns, + session.currentTokenUsage?.turnId === dialogTurnId, + ); const updatedSession = { ...session, dialogTurns: updatedDialogTurns, + currentTokenUsage, loadedTurnCount: nextCanonicalTurnCount, totalTurnCount: countedOptimisticTurn ? Math.max(nextCanonicalTurnCount, projectedSessionTurnCount(session) - 1) @@ -5082,6 +5173,12 @@ export class FlowChatStore { const clampedIndex = Math.max(0, Math.min(turnIndex, session.dialogTurns.length)); const updatedDialogTurns = session.dialogTurns.slice(0, clampedIndex); + const removedCurrentUsageSource = Boolean( + session.currentTokenUsage?.turnId + && session.dialogTurns.slice(clampedIndex).some( + turn => turn.id === session.currentTokenUsage?.turnId, + ), + ); const hasCompleteHistory = session.isPartial !== true; const historyView = this.sessionHistoryViews.get(sessionId); const currentCatalog = session.turnCatalog?.sessionId === sessionId @@ -5099,6 +5196,11 @@ export class FlowChatStore { const updatedSession = { ...session, dialogTurns: updatedDialogTurns, + currentTokenUsage: currentTokenUsageAfterSourceRemoval( + session, + updatedDialogTurns, + removedCurrentUsageSource, + ), ...(hasCompleteHistory ? { loadedTurnCount: updatedDialogTurns.length, totalTurnCount: updatedDialogTurns.length, @@ -5609,18 +5711,23 @@ export class FlowChatStore { public updateTokenUsage( sessionId: string, - tokenUsage: { inputTokens: number; outputTokens?: number; totalTokens: number }, + tokenUsage: Pick< + TokenUsage, + 'inputTokens' | 'outputTokens' | 'totalTokens' | 'turnId' | 'source' + >, dialogTurnId?: string ): void { this.setState(prev => { const session = prev.sessions.get(sessionId); if (!session) return prev; - const nextTokenUsage = { + const nextTokenUsage: TokenUsage = { inputTokens: tokenUsage.inputTokens, outputTokens: tokenUsage.outputTokens, totalTokens: tokenUsage.totalTokens, - timestamp: Date.now() + timestamp: Date.now(), + ...(tokenUsage.turnId ? { turnId: tokenUsage.turnId } : {}), + ...(tokenUsage.source ? { source: tokenUsage.source } : {}), }; let dialogTurns = session.dialogTurns; if (dialogTurnId) { @@ -6269,9 +6376,7 @@ export class FlowChatStore { remoteConnectionId, remoteSshHost, ); - const restoredCurrentTokenUsage = deriveRestoredCurrentTokenUsage( - metadata.customMetadata?.lastRequestTokenUsage, - ); + const persistedCurrentContextUsage = persistedCurrentContextUsageValue(metadata); this.setState(prev => { if (surfaceGeneration !== this.surfaceGeneration) { @@ -6283,6 +6388,9 @@ export class FlowChatStore { const rawAgentType = metadata.agentType || 'agentic'; const validatedAgentType = isValidPersistedAgentType(rawAgentType) ? rawAgentType : 'agentic'; + const restoredCurrentTokenUsage = isAcpAgentType(validatedAgentType) + ? undefined + : deriveRestoredCurrentTokenUsage(persistedCurrentContextUsage); if (rawAgentType !== validatedAgentType) { log.warn('Invalid agentType, falling back to agentic', { sessionId: metadata.sessionId, rawAgentType, validatedAgentType }); @@ -6651,9 +6759,7 @@ export class FlowChatStore { remoteConnectionId, remoteSshHost, ); - const restoredCurrentTokenUsage = deriveRestoredCurrentTokenUsage( - metadata.customMetadata?.lastRequestTokenUsage, - ); + const persistedCurrentContextUsage = persistedCurrentContextUsageValue(metadata); this.setState(prev => { if (prev.sessions.has(metadata.sessionId)) { @@ -6662,6 +6768,9 @@ export class FlowChatStore { const rawAgentType = metadata.agentType || 'agentic'; const validatedAgentType = isValidPersistedAgentType(rawAgentType) ? rawAgentType : 'agentic'; + const restoredCurrentTokenUsage = isAcpAgentType(validatedAgentType) + ? undefined + : deriveRestoredCurrentTokenUsage(persistedCurrentContextUsage); if (rawAgentType !== validatedAgentType) { log.warn('Invalid agentType, falling back to agentic', { sessionId: metadata.sessionId, rawAgentType, validatedAgentType }); @@ -6872,14 +6981,28 @@ export class FlowChatStore { } } + mergedTurns.sort(compareDialogTurnOrder); const turnCatalog = restored.turnCatalog?.sessionId === sessionId ? selectPreferredTurnCatalog(session.turnCatalog, restored.turnCatalog) : session.turnCatalog; - if (!turnsChanged && turnCatalog === session.turnCatalog) { + const restoredAgentType = + restored.session.agentType || session.mode || session.config.agentType; + const currentTokenUsage = reconcileRestoreViewCurrentTokenUsage( + session.currentTokenUsage, + restored.currentContextUsage, + mergedTurns, + restoredAgentType, + snapshotTurns, + ); + if ( + !turnsChanged + && turnCatalog === session.turnCatalog + && restoredAgentType === session.mode + && currentTokenUsage === session.currentTokenUsage + ) { return prev; } - mergedTurns.sort(compareDialogTurnOrder); const newSessions = new Map(prev.sessions); newSessions.set(sessionId, { ...session, @@ -6907,16 +7030,12 @@ export class FlowChatStore { ? { reasoningPreset: restored.session.reasoningPreset?.trim() || undefined } : {}), }, - mode: restored.session.agentType || session.mode, + mode: restoredAgentType, lastUserDialogMode: restored.session.lastUserDialogAgentType || session.lastUserDialogMode, lastSubmittedMode: restored.session.lastSubmittedAgentType ?? session.lastSubmittedMode, - currentTokenUsage: - session.currentTokenUsage - ?? (!isAcpSessionForContextUsage(session) - ? deriveContextUsageFromTurns(mergedTurns) - : undefined), + currentTokenUsage, }); applied = true; @@ -7047,6 +7166,7 @@ export class FlowChatStore { let restoredTotalTurnCount: number | undefined; let restoredTurnCatalog: SessionTurnCatalog | undefined; let restoredTiming: SessionViewRestoreTiming | undefined; + let restoredCurrentContextUsage: SessionContextUsage | null | undefined; // Finish or resume relay history import before Core restores its model // context. Ordinary local sessions return after one metadata read, while @@ -7187,6 +7307,7 @@ export class FlowChatStore { ? restored.turnCatalog : undefined; restoredTiming = restored.timings; + restoredCurrentContextUsage = restored.currentContextUsage; } catch (error) { if (!isUnsupportedTauriCommandError(error, 'restore_session_view')) { throw error; @@ -7345,6 +7466,9 @@ export class FlowChatStore { const session = prev.sessions.get(sessionId); if (!session) return prev; + const restoredAgentType = + restoredSessionInfo?.agentType || session.mode || session.config.agentType; + const updatedSession = { ...session, dialogTurns, @@ -7365,15 +7489,16 @@ export class FlowChatStore { ? { reasoningPreset: restoredSessionInfo.reasoningPreset?.trim() || undefined } : {}), }, - mode: restoredSessionInfo?.agentType || session.mode, + mode: restoredAgentType, lastUserDialogMode: restoredLastUserDialogMode, lastSubmittedMode: restoredSessionInfo?.lastSubmittedAgentType ?? session.lastSubmittedMode, - currentTokenUsage: - session.currentTokenUsage - ?? (!isAcpSessionForContextUsage(session) - ? deriveContextUsageFromTurns(dialogTurns) - : undefined), + currentTokenUsage: reconcileRestoreViewCurrentTokenUsage( + session.currentTokenUsage, + restoredCurrentContextUsage, + dialogTurns, + restoredAgentType, + ), }; const newSessions = new Map(prev.sessions); diff --git a/src/web-ui/src/flow_chat/types/flow-chat.ts b/src/web-ui/src/flow_chat/types/flow-chat.ts index 893bbbeae..55ac8b56e 100644 --- a/src/web-ui/src/flow_chat/types/flow-chat.ts +++ b/src/web-ui/src/flow_chat/types/flow-chat.ts @@ -6,6 +6,7 @@ import type { DialogTurnKind, SessionKind, + SessionContextUsageSource, SessionTitleSource, SessionTurnCatalog, } from '@/shared/types/session-history'; @@ -199,6 +200,10 @@ export interface TokenUsage { outputTokens?: number; totalTokens: number; timestamp: number; + /** Persisted source turn used to invalidate usage after history rewrites. */ + turnId?: string; + /** Runtime provenance for restored session-level context usage. */ + source?: SessionContextUsageSource; } export interface AcpContextUsage { diff --git a/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.test.ts b/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.test.ts index 7c50c66f0..156b585bc 100644 --- a/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.test.ts +++ b/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.test.ts @@ -179,7 +179,10 @@ describe('deriveContextUsageFromTurns', () => { makeTurn({ id: 'turn-2', status: 'completed', tokenUsage: usage(2000) }), ]; - expect(deriveContextUsageFromTurns(turns)).toEqual(usage(2000)); + expect(deriveContextUsageFromTurns(turns)).toEqual({ + ...usage(2000), + turnId: 'turn-2', + }); }); it('skips unfinished turns and falls back to the last completed turn', () => { @@ -189,7 +192,10 @@ describe('deriveContextUsageFromTurns', () => { makeTurn({ id: 'turn-3', status: 'pending', tokenUsage: usage(300) }), ]; - expect(deriveContextUsageFromTurns(turns)).toEqual(usage(1000)); + expect(deriveContextUsageFromTurns(turns)).toEqual({ + ...usage(1000), + turnId: 'turn-1', + }); }); it('skips turns without usage and returns the last completed one that has it', () => { @@ -198,10 +204,13 @@ describe('deriveContextUsageFromTurns', () => { makeTurn({ id: 'turn-2', status: 'error', tokenUsage: usage(2500) }), ]; - expect(deriveContextUsageFromTurns(turns)).toEqual(usage(2500)); + expect(deriveContextUsageFromTurns(turns)).toEqual({ + ...usage(2500), + turnId: 'turn-2', + }); }); - it('skips completed turns with zero or invalid input tokens', () => { + it('uses the latest valid terminal usage even when an older turn is invalid', () => { const turns = [ makeTurn({ id: 'turn-1', @@ -215,10 +224,26 @@ describe('deriveContextUsageFromTurns', () => { }), ]; - expect(deriveContextUsageFromTurns(turns)).toMatchObject({ inputTokens: 420 }); + expect(deriveContextUsageFromTurns(turns)).toMatchObject({ + inputTokens: 420, + turnId: 'turn-2', + }); + }); + + it('does not scan past the latest terminal turn when its usage is invalid', () => { + const turns = [ + makeTurn({ id: 'turn-1', status: 'completed', tokenUsage: usage(1000) }), + makeTurn({ + id: 'turn-2', + status: 'completed', + tokenUsage: { inputTokens: 0, totalTokens: 0, timestamp: 3000 }, + }), + ]; + + expect(deriveContextUsageFromTurns(turns)).toBeUndefined(); }); - it('skips multi-round turns because accumulated usage would overestimate context', () => { + it('does not scan past the latest terminal turn when its multi-round usage is accumulated', () => { const turns = [ makeTurn({ id: 'turn-1', status: 'completed', tokenUsage: usage(1000) }), makeTurn({ @@ -229,7 +254,7 @@ describe('deriveContextUsageFromTurns', () => { }), ]; - expect(deriveContextUsageFromTurns(turns)).toEqual(usage(1000)); + expect(deriveContextUsageFromTurns(turns)).toBeUndefined(); }); it('returns undefined for empty input or when no completed single-round turn has usage', () => { diff --git a/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.ts b/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.ts index 3f4b9961f..ae2708486 100644 --- a/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.ts +++ b/src/web-ui/src/flow_chat/utils/tokenUsageDisplay.ts @@ -36,14 +36,14 @@ function formatCompactNumber(value: number): string { } /** - * Derive the last completed single-round turn's token usage as a - * context-usage approximation. + * Derive the latest terminal turn-with-usage's token usage as a + * context-usage approximation when it represents exactly one model round. * * Used as a fallback to restore `session.currentTokenUsage` when a session is * hydrated from persisted history and no exact last-request usage was stored - * in session metadata. Only single-round turns are used: dialog turn usage - * accumulates across model rounds, so a multi-round turn's input total would - * badly overestimate the current context. + * in session metadata. Dialog turn usage accumulates across model rounds, so + * a multi-round latest turn is not usable. In that case we must not scan past + * it and mislabel an older request as the last request. */ export function deriveContextUsageFromTurns(turns: DialogTurn[] | undefined): TokenUsage | undefined { if (!turns) { @@ -53,22 +53,30 @@ export function deriveContextUsageFromTurns(turns: DialogTurn[] | undefined): To for (let i = turns.length - 1; i >= 0; i--) { const turn = turns[i]; const usage = turn.tokenUsage; - if (!usage) { + if ( + !usage + || ( + turn.status !== 'completed' + && turn.status !== 'error' + && turn.status !== 'cancelled' + ) + ) { continue; } + if ( - turn.status === 'completed' - || turn.status === 'error' - || turn.status === 'cancelled' + turn.modelRounds.length === 1 + && typeof usage.inputTokens === 'number' + && Number.isFinite(usage.inputTokens) + && usage.inputTokens > 0 ) { - if ( - turn.modelRounds.length === 1 - && typeof usage.inputTokens === 'number' - && usage.inputTokens > 0 - ) { - return usage; - } + return { + ...usage, + turnId: turn.id, + }; } + + return undefined; } return undefined; } diff --git a/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts b/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts index 921c5946a..ce019e666 100644 --- a/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts +++ b/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts @@ -5,6 +5,7 @@ import { createTauriCommandError } from '../errors/TauriCommandError'; import type { DialogTurnData, ModelRoundAttemptDiagnostic, + SessionContextUsage, SessionRelationship, SessionTurnCatalog, } from '@/shared/types/session-history'; @@ -241,6 +242,7 @@ export interface SessionViewRestoreTiming { export interface RestoreSessionViewResponse { session: SessionInfo; turns: DialogTurnData[]; + currentContextUsage?: SessionContextUsage | null; turnCatalog?: SessionTurnCatalog; contextRestoreState: 'ready' | 'pending'; isPartial?: boolean; diff --git a/src/web-ui/src/shared/types/session-history.ts b/src/web-ui/src/shared/types/session-history.ts index d865f6108..810492e95 100644 --- a/src/web-ui/src/shared/types/session-history.ts +++ b/src/web-ui/src/shared/types/session-history.ts @@ -43,6 +43,18 @@ export interface SessionCustomMetadata extends Record { titleParams?: Record | null; } +export type SessionContextUsageSource = 'model_request' | 'context_compression'; + +/** Exact session-level context usage persisted by the Agent Session runtime. */ +export interface SessionContextUsage { + turnId: string; + inputTokens: number; + outputTokens?: number; + totalTokens: number; + timestamp: number; + source: SessionContextUsageSource; +} + export interface SessionMetadata { sessionId: string; sessionName: string; @@ -75,6 +87,7 @@ export interface SessionMetadata { snapshotSessionId?: string; tags: string[]; customMetadata?: SessionCustomMetadata; + currentContextUsage?: SessionContextUsage; relationship?: SessionRelationship; todos?: any[]; workspacePath?: string;