Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions src/apps/desktop/src/api/agentic_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -502,6 +502,7 @@ pub struct RestoreSessionWithTurnsResponse {
pub struct RestoreSessionViewResponse {
pub session: SessionResponse,
pub turns: Vec<DialogTurnData>,
pub current_context_usage: Option<SessionContextUsage>,
pub turn_catalog: SessionTurnCatalog,
pub context_restore_state: String,
pub is_partial: bool,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
8 changes: 8 additions & 0 deletions src/apps/desktop/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
51 changes: 42 additions & 9 deletions src/apps/desktop/src/runtime/session_application.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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")]
Expand Down Expand Up @@ -129,6 +125,7 @@ fn local_command_turn_record_request(
pub(crate) struct DesktopSessionViewRestore {
pub session: Session,
pub turns: Vec<DialogTurnData>,
pub current_context_usage: Option<SessionContextUsage>,
pub total_turn_count: usize,
pub turn_catalog: SessionTurnCatalog,
pub timings: SessionViewRestoreTiming,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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();
Expand All @@ -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",
Expand Down Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
157 changes: 157 additions & 0 deletions src/crates/assembly/core/src/agentic/session/context_usage.rs
Original file line number Diff line number Diff line change
@@ -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<SessionManager>,
}

impl SessionContextUsageSubscriber {
pub fn new(session_manager: Arc<SessionManager>) -> 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(&not_applied).is_none());
}
}
2 changes: 2 additions & 0 deletions src/crates/assembly/core/src/agentic/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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::*;
Expand Down
Loading