From e3b423ea12d2cf4e3a3d611d7318f24295a8ea62 Mon Sep 17 00:00:00 2001 From: hanafish <1106510024@qq.com> Date: Wed, 22 Jul 2026 21:54:50 +0800 Subject: [PATCH 1/3] perf(cli): replace reconciliation polling with intent push Pre-commit hook ran. Total eslint: 0, total circular: 0 --- .../agent_sessions/cli/agent_core_bridge.rs | 3 + .../cli/commands/launch_profile.rs | 36 +-- .../cli/commands/resume_delete.rs | 2 +- .../src/agent_sessions/cli/commands/run.rs | 137 +++++++++-- .../src/agent_sessions/cli/commands/status.rs | 47 ++++ .../src/agent_sessions/cli/persistence/mod.rs | 71 ++++++ .../cli/persistence/session_crud.rs | 156 +++++++++++-- .../agent_sessions/cli/persistence/types.rs | 7 + .../cli/session_runner/cursor_usage.rs | 41 ++-- .../cli/session_runner/env_setup.rs | 106 +++++---- .../cli/session_runner/finalize.rs | 117 ++++++---- .../cli/session_runner/helpers.rs | 109 ++++++++- .../cli/session_runner/lifecycle.rs | 79 +++++-- .../cli/session_runner/oauth_setup.rs | 8 +- .../cli/session_runner/proxy_release.rs | 11 +- .../cli/session_runner/session.rs | 14 +- src-tauri/src/api/agent/test/cli.rs | 6 + src-tauri/src/commands/handler_list.inc | 1 + src/api/realtime/websocket/schemas.ts | 1 + src/api/tauri/rpc/procedures/cli.ts | 27 +++ src/api/tauri/rpc/procedures/index.ts | 1 + src/api/tauri/rpc/router.ts | 1 + src/api/tauri/rpc/schemas/cli.ts | 53 +++++ src/api/tauri/rpc/schemas/index.ts | 1 + .../control/turnIntentDispatchLifecycle.ts | 6 + .../SessionCore/control/turnLifecycle.ts | 20 +- .../hooks/session/useQueueDispatch.ts | 8 +- .../SessionCore/services/SessionService.ts | 6 +- .../cli/__tests__/cliTransport.test.ts | 70 ++++++ .../sync/adapters/cli/cliHistory.ts | 14 +- .../sync/adapters/cli/cliLifecycle.ts | 168 ++------------ .../sync/adapters/cli/cliTransport.ts | 94 ++------ .../adapters/cli/createCliEventHandler.ts | 77 +------ .../cliTurnLifecycleCoordinator.test.ts | 177 +++++++++++++++ .../cliSession/cliTurnLifecycleCoordinator.ts | 213 ++++++++++++++++++ .../cliSession/useBackgroundSessionMonitor.ts | 44 ++-- .../useWorkstationSidebarHandlers.ts | 2 + 37 files changed, 1410 insertions(+), 524 deletions(-) create mode 100644 src/api/tauri/rpc/procedures/cli.ts create mode 100644 src/api/tauri/rpc/schemas/cli.ts create mode 100644 src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts create mode 100644 src/hooks/cliSession/cliTurnLifecycleCoordinator.test.ts create mode 100644 src/hooks/cliSession/cliTurnLifecycleCoordinator.ts diff --git a/src-tauri/src/agent_sessions/cli/agent_core_bridge.rs b/src-tauri/src/agent_sessions/cli/agent_core_bridge.rs index a5baeff3fe..4dfc6a423a 100644 --- a/src-tauri/src/agent_sessions/cli/agent_core_bridge.rs +++ b/src-tauri/src/agent_sessions/cli/agent_core_bridge.rs @@ -201,8 +201,11 @@ fn respond_plan_approval( None, Some(AgentExecMode::Build.as_str().to_string()), None, + None, + None, ) .await + .map(|_| ()) }) } diff --git a/src-tauri/src/agent_sessions/cli/commands/launch_profile.rs b/src-tauri/src/agent_sessions/cli/commands/launch_profile.rs index 30ee63e677..30760e726c 100644 --- a/src-tauri/src/agent_sessions/cli/commands/launch_profile.rs +++ b/src-tauri/src/agent_sessions/cli/commands/launch_profile.rs @@ -7,31 +7,39 @@ use super::super::session_runner::launch_profiles::{ use std::collections::HashMap; #[tauri::command] -pub fn cli_launch_profile_get(agent_name: String) -> Result { - launch_profile_store::cli_launch_profile_get(agent_name) +pub async fn cli_launch_profile_get(agent_name: String) -> Result { + tokio::task::spawn_blocking(move || launch_profile_store::cli_launch_profile_get(agent_name)) + .await + .map_err(|err| format!("Task error: {err}"))? } #[tauri::command] -pub fn cli_launch_profile_update( +pub async fn cli_launch_profile_update( agent_name: String, permission_mode: CliPermissionMode, command_override: Option, args_override: Option>, env_override: Option>, ) -> Result { - launch_profile_store::cli_launch_profile_update(CliLaunchProfileUpdate { - agent_name, - permission_mode, - command_override, - args_override, - env_override, - // Experimental app-server transport opt-in is not exposed in the - // settings UI; `None` preserves whatever the store already holds. - transport: None, + tokio::task::spawn_blocking(move || { + launch_profile_store::cli_launch_profile_update(CliLaunchProfileUpdate { + agent_name, + permission_mode, + command_override, + args_override, + env_override, + // Experimental app-server transport opt-in is not exposed in the + // settings UI; `None` preserves whatever the store already holds. + transport: None, + }) }) + .await + .map_err(|err| format!("Task error: {err}"))? } #[tauri::command] -pub fn cli_launch_profile_reset(agent_name: String) -> Result { - launch_profile_store::cli_launch_profile_reset(agent_name) +pub async fn cli_launch_profile_reset(agent_name: String) -> Result { + tokio::task::spawn_blocking(move || launch_profile_store::cli_launch_profile_reset(agent_name)) + .await + .map_err(|err| format!("Task error: {err}"))? } diff --git a/src-tauri/src/agent_sessions/cli/commands/resume_delete.rs b/src-tauri/src/agent_sessions/cli/commands/resume_delete.rs index 439562becf..c8415d0d3c 100644 --- a/src-tauri/src/agent_sessions/cli/commands/resume_delete.rs +++ b/src-tauri/src/agent_sessions/cli/commands/resume_delete.rs @@ -82,7 +82,7 @@ pub async fn cli_agent_resume(session_id: String) -> Result<(), String> { let handle = tokio::spawn(async move { if let Err(e) = - session_runner::run_session(sid.clone(), input, cli_resume_id, None, None).await + session_runner::run_session(sid.clone(), input, cli_resume_id, None, None, None).await { tracing::error!("[CodeSession] Resume of {} failed: {}", sid, e); // Same fail-loud principle as the create path above: log the diff --git a/src-tauri/src/agent_sessions/cli/commands/run.rs b/src-tauri/src/agent_sessions/cli/commands/run.rs index 43e097a5fe..855edb308b 100644 --- a/src-tauri/src/agent_sessions/cli/commands/run.rs +++ b/src-tauri/src/agent_sessions/cli/commands/run.rs @@ -6,6 +6,15 @@ use super::super::persistence; use super::super::session_runner; use super::super::types::{KeySource, SessionStatus}; use agent_core::session::IdeContext; +use serde::Serialize; + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct CliRunReceipt { + pub session_id: String, + pub turn_intent_id: String, + pub status: SessionStatus, +} /// Prepend IDE context (open files, git status, etc.) to the user prompt /// so external CLI agents are aware of the user's IDE state. @@ -45,6 +54,30 @@ pub async fn cli_agent_run( ide_context: Option, mode: Option, images: Option>, +) -> Result<(), String> { + cli_agent_run_internal( + session_id, + user_input, + cli_resume_id, + ide_context, + mode, + images, + uuid::Uuid::new_v4().to_string(), + uuid::Uuid::new_v4().to_string(), + ) + .await +} + +#[allow(clippy::too_many_arguments)] +async fn cli_agent_run_internal( + session_id: String, + user_input: String, + cli_resume_id: Option, + ide_context: Option, + mode: Option, + images: Option>, + turn_intent_id: String, + client_message_id: String, ) -> Result<(), String> { tracing::info!( session_id = %session_id, @@ -65,9 +98,8 @@ pub async fn cli_agent_run( .map_err(|err| format!("Task error: {}", err))??; } - // Hold lock across check + spawn + insert to prevent duplicate agents from - // concurrent calls (e.g., double-click). tokio::spawn returns immediately so - // the lock is held only briefly. + // Hold the registry lock across acceptance persistence + spawn so two + // concurrent calls cannot both create a running intent for one session. let mut sessions = session_runner::RUNNING_SESSIONS.lock().await; // Guard: prevent duplicate parallel agents for the same session @@ -80,10 +112,32 @@ pub async fn cli_agent_run( } } + let persist_session_id = session_id.clone(); + let persist_turn_intent_id = turn_intent_id.clone(); + tokio::task::spawn_blocking(move || { + persistence::accept_cli_turn( + &persist_session_id, + &persist_turn_intent_id, + &client_message_id, + ) + .map_err(|err| format!("failed to accept CLI turn lifecycle: {err}")) + }) + .await + .map_err(|err| format!("Task error: {err}"))??; + + let mut running_msg = serde_json::json!({ + "type": "code_session.status_changed", + "session_id": session_id, + "status": "running", + }); + running_msg["turn_intent_id"] = serde_json::Value::String(turn_intent_id.clone()); + crate::api::websocket_handler::broadcast(running_msg.to_string()); + let sid = session_id.clone(); let cli_input = inject_ide_context_into_prompt(&user_input, ide_context.as_ref()); let resume_id = cli_resume_id.clone(); let agent_mode = mode.clone(); + let runner_turn_intent_id = turn_intent_id.clone(); tracing::info!(session_id = %session_id, "cli_agent_run: spawning background runner"); @@ -95,17 +149,34 @@ pub async fn cli_agent_run( resume_id, agent_mode.as_deref(), images, + Some(&runner_turn_intent_id), ) .await { tracing::error!("[CodeSession] Session {} failed: {}", sid, e); - session_runner::flush_cli_streams_for_session(&sid); + session_runner::flush_cli_streams_for_session(&sid).await; // Best-effort: if marking the row as Failed itself fails, log // it explicitly rather than silently dropping the persistence // error — the session row may be left in `Running` until the // health checker repairs it on next pass. - if let Err(persist_err) = - persistence::update_status_with_error(&sid, SessionStatus::Failed, &e) + let failed_sid = sid.clone(); + let failed_error = e.clone(); + let failed_intent = runner_turn_intent_id.clone(); + let persist_result = tokio::task::spawn_blocking(move || { + persistence::update_cli_turn_lifecycle( + &failed_sid, + SessionStatus::Failed, + Some(&failed_error), + Some(( + &failed_intent, + session_persistence::turn_intents::TurnIntentStatus::Failed, + )), + ) + }) + .await; + if let Err(persist_err) = persist_result + .map_err(|err| err.to_string()) + .and_then(|result| result) { tracing::error!( "[CodeSession] failed to mark session {} as Failed: {}", @@ -115,23 +186,24 @@ pub async fn cli_agent_run( } integrations::proxy::server::stop_session_proxy(&sid).await; session_runner::release_proxy_token_for_session_pub(&sid).await; + let mut failed_msg = serde_json::json!({ + "type": "code_session.status_changed", + "session_id": sid, + "status": "failed", + "error_message": e, + }); + failed_msg["turn_intent_id"] = + serde_json::Value::String(runner_turn_intent_id.clone()); + crate::api::websocket_handler::broadcast(failed_msg.to_string()); } // Remove finished entry from RUNNING_SESSIONS to prevent unbounded growth session_runner::RUNNING_SESSIONS.lock().await.remove(&sid); }); sessions.insert(session_id.clone(), handle); + drop(sessions); tracing::info!(session_id = %session_id, "cli_agent_run: background runner registered"); - persistence::update_status(&session_id, SessionStatus::Running) - .map_err(|err| format!("DB error updating status: {err}"))?; - let running_msg = serde_json::json!({ - "type": "code_session.status_changed", - "session_id": session_id, - "status": "running", - }); - crate::api::websocket_handler::broadcast(running_msg.to_string()); - Ok(()) } @@ -144,6 +216,7 @@ pub async fn cli_agent_run( /// If `model` or `account_id` is provided, updates the session config before /// re-running so the CLI uses the newly selected model/key. #[tauri::command] +#[allow(clippy::too_many_arguments)] pub async fn cli_agent_message( session_id: String, content: String, @@ -152,7 +225,11 @@ pub async fn cli_agent_message( ide_context: Option, mode: Option, images: Option>, -) -> Result<(), String> { + turn_intent_id: Option, + client_message_id: Option, +) -> Result { + let turn_intent_id = turn_intent_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); + let client_message_id = client_message_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); tracing::info!( session_id = %session_id, has_model_override = model.is_some(), @@ -231,9 +308,18 @@ pub async fn cli_agent_message( .map_err(|err| format!("DB error: {}", err))? .and_then(|s| s.cli_session_id) }; - let cli_resume_id = persistence::get_cli_session_id_for_account(&session_id, target_account_id) - .map_err(|err| format!("DB error: {}", err))? - .or_else(|| { + let resume_session_id = session_id.clone(); + let resume_account_id = target_account_id.map(str::to_string); + let account_scoped_resume_id = tokio::task::spawn_blocking(move || { + persistence::get_cli_session_id_for_account( + &resume_session_id, + resume_account_id.as_deref(), + ) + .map_err(|err| format!("DB error: {err}")) + }) + .await + .map_err(|err| format!("Task error: {err}"))??; + let cli_resume_id = account_scoped_resume_id.or_else(|| { if account_id .as_deref() .is_some_and(|new_account_id| session.account_id.as_deref() != Some(new_account_id)) @@ -294,15 +380,22 @@ pub async fn cli_agent_message( // Re-run the session with the new message tracing::info!(session_id = %session_id, "cli_agent_message: dispatching rerun"); - cli_agent_run( - session_id, + cli_agent_run_internal( + session_id.clone(), content, cli_resume_id, ide_context, mode, images, + turn_intent_id.clone(), + client_message_id, ) - .await + .await?; + Ok(CliRunReceipt { + session_id, + turn_intent_id, + status: SessionStatus::Running, + }) } /// Respond to a pending approval request from a CLI agent. diff --git a/src-tauri/src/agent_sessions/cli/commands/status.rs b/src-tauri/src/agent_sessions/cli/commands/status.rs index 29f201a254..97075751ea 100644 --- a/src-tauri/src/agent_sessions/cli/commands/status.rs +++ b/src-tauri/src/agent_sessions/cli/commands/status.rs @@ -4,6 +4,20 @@ use super::super::persistence::{self, CliHistoryMutation, CodeSession}; use super::super::session_runner; use agent_core::state::control_flow::CancelReason; +use serde::Serialize; +use std::collections::HashSet; + +const MAX_CLI_STATUS_BATCH_SESSIONS: usize = 256; + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct CliAgentStatusBatchItem { + pub session_id: String, + pub status: super::super::types::SessionStatus, + pub updated_at: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub turn_intent_id: Option, +} /// Get session status. #[tauri::command] @@ -15,6 +29,39 @@ pub async fn cli_agent_status(session_id: String) -> Result, .map_err(|e| format!("Task error: {}", e))? } +/// Minimal reconnect/focus recovery snapshot. Healthy sessions use pushed +/// lifecycle events; this bounded batch is only the repair path. +#[tauri::command] +pub async fn cli_agent_status_batch( + session_ids: Vec, +) -> Result, String> { + tokio::task::spawn_blocking(move || { + let mut seen = HashSet::new(); + let session_ids: Vec = session_ids + .into_iter() + .filter(|session_id| !session_id.is_empty() && seen.insert(session_id.clone())) + .take(MAX_CLI_STATUS_BATCH_SESSIONS) + .collect(); + let intents = session_persistence::turn_intents::latest_for_sessions(&session_ids) + .map_err(|err| format!("DB error loading turn intents: {err}"))?; + let rows = persistence::status_snapshots(&session_ids) + .map_err(|err| format!("DB error: {err}"))?; + Ok(rows + .into_iter() + .map(|session| CliAgentStatusBatchItem { + turn_intent_id: intents + .get(&session.session_id) + .map(|intent| intent.turn_intent_id.clone()), + session_id: session.session_id, + status: session.status, + updated_at: session.updated_at, + }) + .collect()) + }) + .await + .map_err(|err| format!("Task error: {err}"))? +} + /// Get the last ORGII-side history mutation that invalidated native CLI resume state. #[tauri::command] pub async fn cli_agent_history_mutation( diff --git a/src-tauri/src/agent_sessions/cli/persistence/mod.rs b/src-tauri/src/agent_sessions/cli/persistence/mod.rs index 6e31bd1029..ce02ac027e 100644 --- a/src-tauri/src/agent_sessions/cli/persistence/mod.rs +++ b/src-tauri/src/agent_sessions/cli/persistence/mod.rs @@ -13,6 +13,7 @@ pub use worktree_state::*; #[cfg(test)] mod resume_state_tests { use super::*; + use crate::agent_sessions::cli::types::SessionStatus; use crate::test_utils::test_env; use agent_core::foundation::session_bridge; @@ -52,6 +53,76 @@ mod resume_state_tests { .expect("create test CLI session"); } + #[test] + fn status_snapshots_return_only_requested_existing_sessions() { + let _sandbox = test_env::sandbox(); + create_test_session("cli-status-a", "account-a"); + create_test_session("cli-status-b", "account-b"); + update_status("cli-status-b", SessionStatus::Running).expect("mark running"); + + let rows = + status_snapshots(&["cli-status-b".to_string(), "cli-status-missing".to_string()]) + .expect("load status batch"); + + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].session_id, "cli-status-b"); + assert_eq!(rows[0].status, SessionStatus::Running); + assert!(!rows[0].updated_at.is_empty()); + } + + #[test] + fn cli_session_and_turn_intent_lifecycle_commit_atomically() { + let _sandbox = test_env::sandbox(); + let session_id = "cli-atomic-lifecycle"; + let turn_intent_id = "intent-atomic"; + create_test_session(session_id, "account-a"); + + accept_cli_turn(session_id, turn_intent_id, "message-atomic").expect("accept lifecycle"); + assert_eq!( + get_session(session_id) + .expect("load session") + .expect("session exists") + .status, + SessionStatus::Running + ); + assert_eq!( + session_persistence::turn_intents::list_for_session(session_id).expect("load intent") + [0] + .status, + session_persistence::turn_intents::TurnIntentStatus::Running + ); + + update_cli_turn_lifecycle( + session_id, + SessionStatus::Completed, + None, + Some(( + turn_intent_id, + session_persistence::turn_intents::TurnIntentStatus::Completed, + )), + ) + .expect("complete lifecycle"); + + let rejected = update_cli_turn_lifecycle( + session_id, + SessionStatus::Running, + None, + Some(( + turn_intent_id, + session_persistence::turn_intents::TurnIntentStatus::Running, + )), + ); + assert!(rejected.is_err()); + assert_eq!( + get_session(session_id) + .expect("load session") + .expect("session exists") + .status, + SessionStatus::Completed, + "failed intent transition must roll back the adjacent session status" + ); + } + #[test] fn cli_resume_state_is_scoped_by_account_and_restored_on_switch_back() { let _sandbox = test_env::sandbox(); diff --git a/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs b/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs index 54b456985b..aca1996321 100644 --- a/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs +++ b/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs @@ -1,5 +1,5 @@ use chrono::Utc; -use rusqlite::{params, OptionalExtension, Result as SqliteResult}; +use rusqlite::{params, Connection, OptionalExtension, Result as SqliteResult}; use agent_core::session::AgentExecMode; use database::db::get_connection; @@ -8,7 +8,9 @@ use super::super::types::{ session_defaults, KeySource, SessionRunner, SessionStatus, DEFAULT_CODE_SESSION_FLOW, PERSONAL_ORG_ID, }; -use super::types::{CliHistoryMutation, CodeSession, CreateCodeSessionParams}; +use super::types::{ + CliHistoryMutation, CliSessionStatusSnapshot, CodeSession, CreateCodeSessionParams, +}; pub(super) fn now_iso() -> String { Utc::now().to_rfc3339() @@ -166,6 +168,39 @@ pub fn list_sessions() -> SqliteResult> { rows.collect() } +/// Load only lifecycle fields for a bounded set of sessions. This is used by +/// reconnect/focus reconciliation and deliberately avoids hydrating complete +/// session rows or scanning unrelated sessions. +pub fn status_snapshots(session_ids: &[String]) -> SqliteResult> { + if session_ids.is_empty() { + return Ok(Vec::new()); + } + let conn = get_connection()?; + let placeholders = std::iter::repeat_n("?", session_ids.len()) + .collect::>() + .join(","); + let query = format!( + "SELECT session_id, status, updated_at FROM code_sessions WHERE session_id IN ({placeholders})" + ); + let mut stmt = conn.prepare(&query)?; + let rows = stmt.query_map(rusqlite::params_from_iter(session_ids), |row| { + let raw_status: String = row.get(1)?; + let status = SessionStatus::parse(&raw_status).ok_or_else(|| { + rusqlite::Error::FromSqlConversionFailure( + 1, + rusqlite::types::Type::Text, + format!("invalid code session status: {raw_status}").into(), + ) + })?; + Ok(CliSessionStatusSnapshot { + session_id: row.get(0)?, + status, + updated_at: row.get(2)?, + }) + })?; + rows.collect() +} + /// One page of sessions ordered by recent activity. Serves the sidebar's /// paginated category view without loading the whole table. pub fn list_sessions_page(limit: usize, offset: usize) -> SqliteResult> { @@ -184,21 +219,38 @@ pub fn list_sessions_page(limit: usize, offset: usize) -> SqliteResult SqliteResult { let conn = get_connection()?; + let affected = update_status_row(&conn, session_id, status, None)?; + if affected { + sync_orgtrack_mirror(session_id); + } + Ok(affected) +} + +fn update_status_row( + conn: &Connection, + session_id: &str, + status: SessionStatus, + error: Option<&str>, +) -> SqliteResult { let now = now_iso(); - let affected = if status.is_terminal() { - conn.execute( + let affected = match (status.is_terminal(), error) { + (true, Some(error)) => conn.execute( + "UPDATE code_sessions SET status = ?2, error_message = ?3, pid = NULL, updated_at = ?4 WHERE session_id = ?1", + params![session_id, status.as_ref(), error, now], + )?, + (false, Some(error)) => conn.execute( + "UPDATE code_sessions SET status = ?2, error_message = ?3, updated_at = ?4 WHERE session_id = ?1", + params![session_id, status.as_ref(), error, now], + )?, + (true, None) => conn.execute( "UPDATE code_sessions SET status = ?2, pid = NULL, updated_at = ?3 WHERE session_id = ?1", params![session_id, status.as_ref(), now], - )? - } else { - conn.execute( + )?, + (false, None) => conn.execute( "UPDATE code_sessions SET status = ?2, updated_at = ?3 WHERE session_id = ?1", params![session_id, status.as_ref(), now], - )? + )?, }; - if affected > 0 { - sync_orgtrack_mirror(session_id); - } Ok(affected > 0) } @@ -209,22 +261,76 @@ pub fn update_status_with_error( error: &str, ) -> SqliteResult { let conn = get_connection()?; - let now = now_iso(); - let affected = if status.is_terminal() { - conn.execute( - "UPDATE code_sessions SET status = ?2, error_message = ?3, pid = NULL, updated_at = ?4 WHERE session_id = ?1", - params![session_id, status.as_ref(), error, now], - )? - } else { - conn.execute( - "UPDATE code_sessions SET status = ?2, error_message = ?3, updated_at = ?4 WHERE session_id = ?1", - params![session_id, status.as_ref(), error, now], - )? - }; - if affected > 0 { + let affected = update_status_row(&conn, session_id, status, Some(error))?; + if affected { sync_orgtrack_mirror(session_id); } - Ok(affected > 0) + Ok(affected) +} + +/// Atomically accept a CLI turn: the session and its intent become running +/// together, so reconnect cannot observe a split-brain lifecycle snapshot. +pub fn accept_cli_turn( + session_id: &str, + turn_intent_id: &str, + client_message_id: &str, +) -> Result<(), String> { + let conn = get_connection().map_err(|err| err.to_string())?; + let tx = conn + .unchecked_transaction() + .map_err(|err| err.to_string())?; + if !update_status_row(&tx, session_id, SessionStatus::Running, None) + .map_err(|err| err.to_string())? + { + return Err(format!("session not found: {session_id}")); + } + session_persistence::turn_intents::upsert_initial_on( + &tx, + session_id, + turn_intent_id, + Some(client_message_id), + session_persistence::turn_intents::TurnIntentSource::UserSubmit, + session_persistence::turn_intents::TurnIntentStatus::Queued, + ) + .map_err(|err| err.to_string())?; + session_persistence::turn_intents::update_status_on( + &tx, + session_id, + turn_intent_id, + session_persistence::turn_intents::TurnIntentStatus::Running, + ) + .map_err(|err| err.to_string())?; + tx.commit().map_err(|err| err.to_string())?; + sync_orgtrack_mirror(session_id); + Ok(()) +} + +/// Atomically persist a CLI session status and the matching intent terminal. +pub fn update_cli_turn_lifecycle( + session_id: &str, + status: SessionStatus, + error: Option<&str>, + turn_intent: Option<(&str, session_persistence::turn_intents::TurnIntentStatus)>, +) -> Result<(), String> { + let conn = get_connection().map_err(|err| err.to_string())?; + let tx = conn + .unchecked_transaction() + .map_err(|err| err.to_string())?; + if !update_status_row(&tx, session_id, status, error).map_err(|err| err.to_string())? { + return Err(format!("session not found: {session_id}")); + } + if let Some((turn_intent_id, intent_status)) = turn_intent { + session_persistence::turn_intents::update_status_on( + &tx, + session_id, + turn_intent_id, + intent_status, + ) + .map_err(|err| err.to_string())?; + } + tx.commit().map_err(|err| err.to_string())?; + sync_orgtrack_mirror(session_id); + Ok(()) } /// Store the PID of the CLI subprocess. diff --git a/src-tauri/src/agent_sessions/cli/persistence/types.rs b/src-tauri/src/agent_sessions/cli/persistence/types.rs index aa6d7501dc..edd87831b0 100644 --- a/src-tauri/src/agent_sessions/cli/persistence/types.rs +++ b/src-tauri/src/agent_sessions/cli/persistence/types.rs @@ -3,6 +3,13 @@ use serde::{Deserialize, Serialize}; use super::super::types::KeySource; use super::super::types::SessionStatus; +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct CliSessionStatusSnapshot { + pub session_id: String, + pub status: SessionStatus, + pub updated_at: String, +} + /// A code generation session record. #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] diff --git a/src-tauri/src/agent_sessions/cli/session_runner/cursor_usage.rs b/src-tauri/src/agent_sessions/cli/session_runner/cursor_usage.rs index 484b594dfe..0b6368de34 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/cursor_usage.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/cursor_usage.rs @@ -78,19 +78,34 @@ pub(super) async fn fetch_cursor_usage_for_session( return; } - if let Err(err) = session_persistence::token_usage::insert_token_usage_record( - session_id, - "code", - summary.dominant_model.as_deref(), - account_id, - summary.input_tokens as i64, - summary.output_tokens as i64, - summary.cache_read_tokens as i64, - summary.cache_write_tokens as i64, - summary.total_tokens as i64, - 0, - None, - ) { + let persist_session_id = session_id.to_string(); + let persist_model = summary.dominant_model.clone(); + let persist_account_id = account_id.map(str::to_string); + let input_tokens = summary.input_tokens as i64; + let output_tokens = summary.output_tokens as i64; + let cache_read_tokens = summary.cache_read_tokens as i64; + let cache_write_tokens = summary.cache_write_tokens as i64; + let total_tokens = summary.total_tokens as i64; + let persist_result = tokio::task::spawn_blocking(move || { + session_persistence::token_usage::insert_token_usage_record( + &persist_session_id, + "code", + persist_model.as_deref(), + persist_account_id.as_deref(), + input_tokens, + output_tokens, + cache_read_tokens, + cache_write_tokens, + total_tokens, + 0, + None, + ) + }) + .await; + if let Err(err) = persist_result + .map_err(|join_err| join_err.to_string()) + .and_then(|db_result| db_result.map_err(|db_err| db_err.to_string())) + { tracing::warn!( "[CursorUsage] Failed to insert per-round token usage for session {}: {}", session_id, diff --git a/src-tauri/src/agent_sessions/cli/session_runner/env_setup.rs b/src-tauri/src/agent_sessions/cli/session_runner/env_setup.rs index ae90b0ca49..20e0d68df7 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/env_setup.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/env_setup.rs @@ -355,21 +355,8 @@ pub(super) fn apply_system_proxy_passthrough(env_vars: &mut HashMap, -) { - if !(matches!(agent, ModelType::Codex) && session.key_source == KeySource::HostedKey) { - return; - } - +fn ensure_codex_hosted_proxy_config(proxy_url_val: &str) { if let Some(home) = dirs::home_dir() { - let proxy_url_val = session.proxy_url.as_deref().unwrap_or(""); let codex_dir = home.join(".codex"); let config_file = codex_dir.join("config.toml"); @@ -419,6 +406,29 @@ pub(super) async fn setup_codex_hosted_proxy( } } } +} + +/// For a Codex hosted-key session: ensure `~/.codex/config.toml` has the proxy +/// `model_providers.proxy` section, then run `codex login --with-api-key` with +/// the proxy token. No-op for any other agent/key-source. Failures are logged +/// and the session continues. +pub(super) async fn setup_codex_hosted_proxy( + agent: &ModelType, + session: &CodeSession, + env_vars: &HashMap, +) { + if !(matches!(agent, ModelType::Codex) && session.key_source == KeySource::HostedKey) { + return; + } + + let proxy_url = session.proxy_url.clone().unwrap_or_default(); + if let Err(err) = tokio::task::spawn_blocking(move || { + ensure_codex_hosted_proxy_config(&proxy_url); + }) + .await + { + tracing::warn!("[CodeSession] Codex proxy config task failed: {err}"); + } let api_key_val = session.proxy_token.as_deref().unwrap_or(""); if !api_key_val.is_empty() { @@ -484,37 +494,43 @@ pub(super) async fn setup_opencode_sse_sanitizer( return; } - if let Ok(config_text) = std::fs::read_to_string( - dirs::config_dir() - .unwrap_or_default() - .join("opencode") - .join("opencode.json"), - ) { - if let Ok(config) = serde_json::from_str::(&config_text) { - let base_url = config - .get("provider") - .and_then(|p| p.get("anthropic")) - .and_then(|a| a.get("options")) - .and_then(|o| o.get("baseURL")) - .and_then(|v| v.as_str()); - if let Some(upstream) = base_url { - if !upstream.contains("127.0.0.1") && !upstream.contains("localhost") { - match integrations::proxy::sse_sanitizer::ensure_running(upstream).await { - Ok(local_url) => { - tracing::info!( - "[CodeSession] SSE sanitizer active: {} → {}", - local_url, - upstream - ); - env_vars.insert("ANTHROPIC_BASE_URL".to_string(), local_url); - } - Err(err) => { - tracing::warn!( - "[CodeSession] SSE sanitizer failed: {} — using direct connection", - err - ); - } - } + let upstream = tokio::task::spawn_blocking(|| { + let config_text = std::fs::read_to_string( + dirs::config_dir() + .unwrap_or_default() + .join("opencode") + .join("opencode.json"), + ) + .ok()?; + let config = serde_json::from_str::(&config_text).ok()?; + config + .get("provider")? + .get("anthropic")? + .get("options")? + .get("baseURL")? + .as_str() + .map(str::to_string) + }) + .await + .ok() + .flatten(); + + if let Some(upstream) = upstream { + if !upstream.contains("127.0.0.1") && !upstream.contains("localhost") { + match integrations::proxy::sse_sanitizer::ensure_running(&upstream).await { + Ok(local_url) => { + tracing::info!( + "[CodeSession] SSE sanitizer active: {} → {}", + local_url, + upstream + ); + env_vars.insert("ANTHROPIC_BASE_URL".to_string(), local_url); + } + Err(err) => { + tracing::warn!( + "[CodeSession] SSE sanitizer failed: {} — using direct connection", + err + ); } } } diff --git a/src-tauri/src/agent_sessions/cli/session_runner/finalize.rs b/src-tauri/src/agent_sessions/cli/session_runner/finalize.rs index e4abf4204a..9fa9b8ef5a 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/finalize.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/finalize.rs @@ -48,6 +48,7 @@ pub(super) async fn finalize_session_run( use_codex_app_server: bool, is_acp_agent: bool, synced_rule_files: &[std::path::PathBuf], + turn_intent_id: Option<&str>, outcome: SessionRunOutcome, ) { let SessionRunOutcome { @@ -62,29 +63,44 @@ pub(super) async fn finalize_session_run( let session_id = session.session_id.as_str(); let account_id = session.account_id.as_deref(); - if *agent == ModelType::Codex && session.key_source == KeySource::OwnKey { - let launched_access_token = env_vars.get("OPENAI_API_KEY").map(String::as_str); - if let Err(err) = sync_codex_cli_auth_to_key_vault(account_id, launched_access_token) { - tracing::warn!( - "[CodeSession] Failed to sync Codex CLI auth tokens: {}", - err - ); - } - if exit_code == 0 { - if let Some(account_id) = account_id { - if let Err(err) = KEY_SERVICE.reset_oauth_refresh_failures(account_id) { - tracing::warn!( - "[CodeSession] Failed to reset Codex OAuth refresh failures: {}", - err - ); + let setup_is_codex_own_key = + *agent == ModelType::Codex && session.key_source == KeySource::OwnKey; + let setup_access_token = env_vars.get("OPENAI_API_KEY").cloned(); + let setup_account_id = account_id.map(str::to_string); + let setup_session_id = session_id.to_string(); + let setup_cli_session_id = cli_session_id_out.clone(); + let _ = tokio::task::spawn_blocking(move || { + if setup_is_codex_own_key { + if let Err(err) = sync_codex_cli_auth_to_key_vault( + setup_account_id.as_deref(), + setup_access_token.as_deref(), + ) { + tracing::warn!( + "[CodeSession] Failed to sync Codex CLI auth tokens: {}", + err + ); + } + if exit_code == 0 { + if let Some(account_id) = setup_account_id.as_deref() { + if let Err(err) = KEY_SERVICE.reset_oauth_refresh_failures(account_id) { + tracing::warn!( + "[CodeSession] Failed to reset Codex OAuth refresh failures: {}", + err + ); + } } } } - } - - if let Some(ref cli_sid) = cli_session_id_out { - persistence::update_cli_session_id_for_account(session_id, account_id, cli_sid).ok(); - } + if let Some(cli_session_id) = setup_cli_session_id.as_deref() { + persistence::update_cli_session_id_for_account( + &setup_session_id, + setup_account_id.as_deref(), + cli_session_id, + ) + .ok(); + } + }) + .await; let raw_final_status = if cli_plan_approval_gate_reached { SessionStatus::Completed @@ -162,15 +178,25 @@ pub(super) async fn finalize_session_run( None }; - if *agent == ModelType::Codex + let should_record_oauth_failure = *agent == ModelType::Codex && session.key_source == KeySource::OwnKey && error_message .as_deref() - .is_some_and(is_cli_oauth_failure_message) - { - if let Some(account_id) = account_id { - if let Some(ref err_msg) = error_message { - if let Err(err) = KEY_SERVICE.record_oauth_refresh_failure(account_id, err_msg) { + .is_some_and(is_cli_oauth_failure_message); + + let persist_session_id = session_id.to_string(); + let persist_error_message = error_message.clone(); + let persist_turn_intent_id = turn_intent_id.map(str::to_string); + let persist_account_id = account_id.map(str::to_string); + let persist_result = tokio::task::spawn_blocking(move || { + if should_record_oauth_failure { + if let (Some(account_id), Some(error_message)) = ( + persist_account_id.as_deref(), + persist_error_message.as_deref(), + ) { + if let Err(err) = + KEY_SERVICE.record_oauth_refresh_failure(account_id, error_message) + { tracing::warn!( "[CodeSession] Failed to record Codex OAuth refresh failure: {}", err @@ -178,17 +204,29 @@ pub(super) async fn finalize_session_run( } } } - } - - if let Some(ref err_msg) = error_message { - if let Err(err) = persistence::update_status_with_error(session_id, final_status, err_msg) { - tracing::error!( - "[CodeSession] Failed to update final status with error: {}", - err - ); - } - } else if let Err(err) = persistence::update_status(session_id, final_status) { - tracing::error!("[CodeSession] Failed to update final status: {}", err); + let intent_status = persist_turn_intent_id.as_deref().map(|turn_intent_id| { + ( + turn_intent_id, + if raw_final_status == SessionStatus::Completed { + session_persistence::turn_intents::TurnIntentStatus::Completed + } else { + session_persistence::turn_intents::TurnIntentStatus::Failed + }, + ) + }); + persistence::update_cli_turn_lifecycle( + &persist_session_id, + final_status, + persist_error_message.as_deref(), + intent_status, + ) + }) + .await; + if let Err(err) = persist_result + .map_err(|join_err| join_err.to_string()) + .and_then(|result| result) + { + tracing::error!("[CodeSession] Failed to persist final lifecycle: {}", err); } if final_status.is_terminal() { @@ -213,7 +251,7 @@ pub(super) async fn finalize_session_run( } // Flush any pending streaming deltas before signaling session end - flush_and_broadcast(session_id); + flush_and_broadcast(session_id).await; let mut status_msg = serde_json::json!({ "type": "code_session.status_changed", @@ -226,6 +264,9 @@ pub(super) async fn finalize_session_run( if let Some(ref err_msg) = error_message { status_msg["error_message"] = serde_json::Value::String(err_msg.clone()); } + if let Some(turn_intent_id) = turn_intent_id { + status_msg["turn_intent_id"] = serde_json::Value::String(turn_intent_id.to_string()); + } websocket_handler::broadcast(status_msg.to_string()); // ── Worktree: commit changes on completion ── diff --git a/src-tauri/src/agent_sessions/cli/session_runner/helpers.rs b/src-tauri/src/agent_sessions/cli/session_runner/helpers.rs index 79d231f312..557c8594b7 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/helpers.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/helpers.rs @@ -55,7 +55,54 @@ pub(super) fn strip_ide_context(input: &str) -> String { /// /// Shared helper used by both the ACP flow (Copilot) and the standard /// CliAgentParser loop (all other agents). -pub(super) fn emit_chunk( +pub(super) async fn emit_chunk( + chunk: &core_types::activity::ActivityChunk, + session_id: &str, + sequence: &mut i64, +) { + let action_type = chunk.action_type.as_str(); + let is_delta = action_type.contains("delta") + && chunk + .result + .get("is_delta") + .and_then(|v| v.as_bool()) + .unwrap_or(false); + let delta_requires_flush = action_type == "tool_call_delta" + && (chunk + .result + .get("tool_call_id") + .and_then(|v| v.as_str()) + .is_some_and(|value| !value.is_empty()) + || chunk + .result + .get("tool_name") + .and_then(|v| v.as_str()) + .is_some_and(|value| !value.is_empty())); + + // Ordinary token deltas are memory-only and latency-sensitive. Keep them + // on the async runner; only chunks that may touch SQLite, the event cache, + // or filesystem side effects cross onto the blocking pool. + if is_delta && !delta_requires_flush { + emit_chunk_blocking(chunk, session_id, sequence); + return; + } + + let owned_chunk = chunk.clone(); + let owned_session_id = session_id.to_string(); + let initial_sequence = *sequence; + match tokio::task::spawn_blocking(move || { + let mut next_sequence = initial_sequence; + emit_chunk_blocking(&owned_chunk, &owned_session_id, &mut next_sequence); + next_sequence + }) + .await + { + Ok(next_sequence) => *sequence = next_sequence, + Err(err) => tracing::error!("[CodeSession] chunk persistence task failed: {err}"), + } +} + +fn emit_chunk_blocking( chunk: &core_types::activity::ActivityChunk, session_id: &str, sequence: &mut i64, @@ -91,7 +138,7 @@ pub(super) fn emit_chunk( .filter(|v| !v.is_empty()) .is_some(); if has_tool_identity { - flush_and_broadcast(session_id); + flush_and_broadcast_blocking(session_id); } } @@ -145,7 +192,7 @@ pub(super) fn emit_chunk( } else { // Non-streaming chunk (tool_call, user_message, etc.): flush any // pending streams before appending, same as UnifiedEventHandler. - flush_and_broadcast(session_id); + flush_and_broadcast_blocking(session_id); } // Persist non-delta chunks to DB (legacy mode). Native-transcript @@ -253,7 +300,7 @@ fn persist_streaming_complete_chunk( } /// Flush all pending CLI streams and broadcast completion events. -pub(super) fn flush_and_broadcast(session_id: &str) { +fn flush_and_broadcast_blocking(session_id: &str) { let mut sequence = next_chunk_sequence(session_id); for event in crate::agent_sessions::event_pipeline::streaming::cli_flush_session(session_id) { let stream_type = if event.action_type == "assistant" { @@ -270,8 +317,19 @@ pub(super) fn flush_and_broadcast(session_id: &str) { } } -pub fn flush_cli_streams_for_session(session_id: &str) { - flush_and_broadcast(session_id); +pub(super) async fn flush_and_broadcast(session_id: &str) { + let owned_session_id = session_id.to_string(); + if let Err(err) = tokio::task::spawn_blocking(move || { + flush_and_broadcast_blocking(&owned_session_id); + }) + .await + { + tracing::error!("[CodeSession] stream flush task failed: {err}"); + } +} + +pub async fn flush_cli_streams_for_session(session_id: &str) { + flush_and_broadcast(session_id).await; } /// Drop hook-derived live status for a finished managed session. The @@ -328,7 +386,7 @@ fn is_cli_file_edit_function(function_name: &str) -> bool { /// /// Non-fatal: snapshot failures are logged at `warn` level and never block the /// chunk from being persisted and broadcast. -pub(super) fn snapshot_cli_file_edit( +fn snapshot_cli_file_edit_blocking( session_id: &str, snapshot_id: &str, chunk: &core_types::activity::ActivityChunk, @@ -422,6 +480,30 @@ pub(super) fn snapshot_cli_file_edit( } } +pub(super) async fn snapshot_cli_file_edit( + session_id: &str, + snapshot_id: &str, + chunk: &core_types::activity::ActivityChunk, + repo_path: &str, +) { + let owned_session_id = session_id.to_string(); + let owned_snapshot_id = snapshot_id.to_string(); + let owned_chunk = chunk.clone(); + let owned_repo_path = repo_path.to_string(); + if let Err(err) = tokio::task::spawn_blocking(move || { + snapshot_cli_file_edit_blocking( + &owned_session_id, + &owned_snapshot_id, + &owned_chunk, + &owned_repo_path, + ); + }) + .await + { + tracing::warn!("[cli_snapshot] snapshot task failed: {err}"); + } +} + /// Save base64 data-URL images to `~/.orgii/session-images/` and return file paths. /// /// Delegates to `agent_core::images::persist_images` which uses content-hash @@ -436,7 +518,18 @@ pub(super) async fn persist_attached_images( return vec![]; } - let paths = agent_core::persistence::images::persist_images(imgs); + let owned_images = imgs.to_vec(); + let paths = match tokio::task::spawn_blocking(move || { + agent_core::persistence::images::persist_images(&owned_images) + }) + .await + { + Ok(paths) => paths, + Err(err) => { + tracing::warn!("[CodeSession] image persistence task failed: {err}"); + Vec::new() + } + }; if !paths.is_empty() { tracing::info!( diff --git a/src-tauri/src/agent_sessions/cli/session_runner/lifecycle.rs b/src-tauri/src/agent_sessions/cli/session_runner/lifecycle.rs index cf064a895b..699800c5d6 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/lifecycle.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/lifecycle.rs @@ -55,18 +55,23 @@ pub async fn terminate_process_tree(pid: i64, _label: &str) { /// in the database — callers are responsible for setting the appropriate final status /// (e.g., Cancelled for user cancel, or nothing before a re-run). pub async fn kill_running_agent(session_id: &str) -> bool { - let had_running_task = { + let running_task = { let mut sessions = RUNNING_SESSIONS.lock().await; - if let Some(handle) = sessions.remove(session_id) { - flush_cli_streams_for_session(session_id); - handle.abort(); - true - } else { - false - } + sessions.remove(session_id) }; + let had_running_task = running_task.is_some(); + if let Some(handle) = running_task { + // The flush can touch SQLite. Never hold the global runner registry + // lock across that blocking-pool round trip or unrelated sessions' + // start/stop operations would serialize behind it. + flush_cli_streams_for_session(session_id).await; + handle.abort(); + } - if let Ok(Some(session)) = persistence::get_session(session_id) { + let process_session_id = session_id.to_string(); + let persisted_session = + tokio::task::spawn_blocking(move || persistence::get_session(&process_session_id)).await; + if let Ok(Ok(Some(session))) = persisted_session { if let Some(pid) = session.pid { terminate_process_tree(pid, session_id).await; } @@ -92,15 +97,39 @@ pub async fn cancel_session(session_id: &str, reason: CancelReason) -> Result s, - Err(err) => { + let lookup_session_id = session_id.to_string(); + let (session, active_turn_intent_id) = match tokio::task::spawn_blocking(move || { + let session = + persistence::get_session(&lookup_session_id).map_err(|err| err.to_string())?; + let latest = session_persistence::turn_intents::latest_for_sessions(std::slice::from_ref( + &lookup_session_id, + )) + .map_err(|err| err.to_string())? + .remove(&lookup_session_id) + .filter(|intent| { + intent.status == session_persistence::turn_intents::TurnIntentStatus::Running + }) + .map(|intent| intent.turn_intent_id); + Ok::<_, String>((session, latest)) + }) + .await + { + Ok(Ok(result)) => result, + Ok(Err(err)) => { tracing::warn!( session_id = %session_id, error = %err, "cli::cancel_session: get_session DB error; broadcast will lack session metadata" ); - None + (None, None) + } + Err(err) => { + tracing::warn!( + session_id = %session_id, + error = %err, + "cli::cancel_session: status lookup task failed" + ); + (None, None) } }; @@ -112,8 +141,23 @@ pub async fn cancel_session(session_id: &str, reason: CancelReason) -> Result Result Result s, - _ => return, - }; + let load_session_id = session_id.to_string(); + let session = + match tokio::task::spawn_blocking(move || persistence::get_session(&load_session_id)).await + { + Ok(Ok(Some(s))) => s, + _ => return, + }; if session.key_source != KeySource::HostedKey { return; diff --git a/src-tauri/src/agent_sessions/cli/session_runner/session.rs b/src-tauri/src/agent_sessions/cli/session_runner/session.rs index fac067d77d..6fea7beea7 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/session.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/session.rs @@ -48,6 +48,7 @@ pub async fn run_session( cli_resume_id: Option, mode: Option<&str>, images: Option>, + turn_intent_id: Option<&str>, ) -> Result<(), String> { let session = persistence::get_session(&session_id) .map_err(|e| format!("DB error: {}", e))? @@ -332,18 +333,6 @@ pub async fn run_session( .map_err(|e| format!("DB: failed to store user_input: {}", e))?; } - if let Err(err) = persistence::update_status(&session_id, SessionStatus::Running) { - tracing::error!("[CodeSession] Failed to update status to running: {}", err); - return Err(format!("DB error updating status: {}", err)); - } - - let running_msg = serde_json::json!({ - "type": "code_session.status_changed", - "session_id": session_id, - "status": "running", - }); - websocket_handler::broadcast(running_msg.to_string()); - // Start per-session MITM proxy if needed let needs_mitm = session.key_source == KeySource::HostedKey && agent.needs_mitm_proxy(); @@ -694,6 +683,7 @@ pub async fn run_session( use_codex_app_server, is_acp_agent, &synced_rule_files, + turn_intent_id, super::finalize::SessionRunOutcome { exit_code, cli_session_id_out, diff --git a/src-tauri/src/api/agent/test/cli.rs b/src-tauri/src/api/agent/test/cli.rs index 1c9d441ea7..6bd8689849 100644 --- a/src-tauri/src/api/agent/test/cli.rs +++ b/src-tauri/src/api/agent/test/cli.rs @@ -328,6 +328,8 @@ pub async fn test_cursor_cli_account_switch( None, None, None, + None, + None, ) .await { @@ -450,6 +452,8 @@ pub async fn test_claude_code_cli_account_switch( None, None, None, + None, + None, ) .await { @@ -609,6 +613,8 @@ pub async fn test_codex_cli_account_switch( None, None, None, + None, + None, ) .await { diff --git a/src-tauri/src/commands/handler_list.inc b/src-tauri/src/commands/handler_list.inc index 7663db9c86..b06774cee9 100644 --- a/src-tauri/src/commands/handler_list.inc +++ b/src-tauri/src/commands/handler_list.inc @@ -532,6 +532,7 @@ agent_sessions::cli::commands::cli_agent_run, agent_sessions::cli::commands::cli_agent_message, agent_sessions::cli::commands::cli_agent_approval_response, agent_sessions::cli::commands::cli_agent_status, +agent_sessions::cli::commands::cli_agent_status_batch, agent_sessions::cli::commands::cli_agent_history_mutation, agent_sessions::cli::commands::cli_agent_cancel, agent_sessions::cli::commands::cli_agent_tui_release, diff --git a/src/api/realtime/websocket/schemas.ts b/src/api/realtime/websocket/schemas.ts index 2502df9ba4..4066dde8b0 100644 --- a/src/api/realtime/websocket/schemas.ts +++ b/src/api/realtime/websocket/schemas.ts @@ -345,6 +345,7 @@ export const CodeEditorWebSocketMessageSchema = z data: z.unknown().optional(), payload: z.unknown().optional(), status: z.unknown().optional(), + turn_intent_id: z.string().optional(), files: z.array(z.unknown()).optional(), timestamp: z.number().optional(), }) diff --git a/src/api/tauri/rpc/procedures/cli.ts b/src/api/tauri/rpc/procedures/cli.ts new file mode 100644 index 0000000000..2135df3dec --- /dev/null +++ b/src/api/tauri/rpc/procedures/cli.ts @@ -0,0 +1,27 @@ +import { z } from "zod/v4"; + +import { defineProcedure } from "../invoke"; +import * as schemas from "../schemas"; + +export const cli = { + message: defineProcedure("cli_agent_message") + .input(schemas.cli.CliMessageInputSchema) + .output(schemas.cli.CliRunReceiptSchema) + .build(), + status: defineProcedure("cli_agent_status") + .input(schemas.cli.CliSessionIdInputSchema) + .output(schemas.cli.CliStatusSchema.nullable()) + .build(), + statusBatch: defineProcedure("cli_agent_status_batch") + .input(schemas.cli.CliStatusBatchInputSchema) + .output(z.array(schemas.cli.CliStatusBatchItemSchema)) + .build(), + chunks: defineProcedure("cli_agent_chunks") + .input(schemas.cli.CliSessionIdInputSchema) + .output(schemas.cli.CliChunksSchema) + .build(), + cancel: defineProcedure("cli_agent_cancel") + .input(schemas.cli.CliCancelInputSchema) + .output(z.boolean()) + .build(), +} as const; diff --git a/src/api/tauri/rpc/procedures/index.ts b/src/api/tauri/rpc/procedures/index.ts index c52fbf1857..7cfe6e4a21 100644 --- a/src/api/tauri/rpc/procedures/index.ts +++ b/src/api/tauri/rpc/procedures/index.ts @@ -18,3 +18,4 @@ export { terminal } from "./terminal"; export { tools } from "./tools"; export { validation } from "./validation"; export { workspaceMemory } from "./workspaceMemory"; +export { cli } from "./cli"; diff --git a/src/api/tauri/rpc/router.ts b/src/api/tauri/rpc/router.ts index f524ef47b9..374e3b5b84 100644 --- a/src/api/tauri/rpc/router.ts +++ b/src/api/tauri/rpc/router.ts @@ -47,6 +47,7 @@ export const procedures = { mcp: p.mcp, flow: p.flow, humanSession: p.humanSession, + cli: p.cli, } as const; // ============================================================================ diff --git a/src/api/tauri/rpc/schemas/cli.ts b/src/api/tauri/rpc/schemas/cli.ts new file mode 100644 index 0000000000..96e91b70bc --- /dev/null +++ b/src/api/tauri/rpc/schemas/cli.ts @@ -0,0 +1,53 @@ +import { z } from "zod/v4"; + +import { ActivityChunkSchema } from "@src/api/realtime/websocket/schemas"; + +export const CliMessageInputSchema = z.object({ + sessionId: z.string().min(1), + content: z.string(), + turnIntentId: z.string().min(1), + clientMessageId: z.string().min(1), + model: z.string().optional(), + accountId: z.string().optional(), + ideContext: z.unknown().optional(), + mode: z.string().optional(), + images: z.array(z.string()).optional(), +}); + +export const CliRunReceiptSchema = z.object({ + sessionId: z.string(), + turnIntentId: z.string(), + status: z.string(), +}); + +export const CliSessionIdInputSchema = z.object({ + sessionId: z.string().min(1), +}); + +export const CliCancelInputSchema = CliSessionIdInputSchema.extend({ + reason: z.string().optional(), +}); + +export const CliStatusSchema = z + .object({ + sessionId: z.string(), + status: z.string(), + updatedAt: z.string(), + errorMessage: z.string().nullable().optional(), + totalTokens: z.number().optional(), + transcriptSource: z.string().optional(), + }) + .passthrough(); + +export const CliStatusBatchInputSchema = z.object({ + sessionIds: z.array(z.string().min(1)).max(256), +}); + +export const CliStatusBatchItemSchema = z.object({ + sessionId: z.string(), + status: z.string(), + updatedAt: z.string(), + turnIntentId: z.string().optional(), +}); + +export const CliChunksSchema = z.array(ActivityChunkSchema); diff --git a/src/api/tauri/rpc/schemas/index.ts b/src/api/tauri/rpc/schemas/index.ts index bc53d4a11e..ee4d5097f3 100644 --- a/src/api/tauri/rpc/schemas/index.ts +++ b/src/api/tauri/rpc/schemas/index.ts @@ -25,3 +25,4 @@ export * as mcp from "./mcp"; export * as flow from "./flow"; export * as humanSession from "./humanSession"; export * as sessionCore from "./sessionCore"; +export * as cli from "./cli"; diff --git a/src/engines/SessionCore/control/turnIntentDispatchLifecycle.ts b/src/engines/SessionCore/control/turnIntentDispatchLifecycle.ts index 22267a7f46..44f914bb20 100644 --- a/src/engines/SessionCore/control/turnIntentDispatchLifecycle.ts +++ b/src/engines/SessionCore/control/turnIntentDispatchLifecycle.ts @@ -60,6 +60,12 @@ export function waitForTurnIntentDispatch( }); } +export function getTurnIntentDispatch( + turnIntentId: string +): TurnIntentDispatch | undefined { + return recentDispatches.get(turnIntentId); +} + export function resetTurnIntentDispatchLifecycleForTests(): void { recentDispatches.clear(); waiters.clear(); diff --git a/src/engines/SessionCore/control/turnLifecycle.ts b/src/engines/SessionCore/control/turnLifecycle.ts index f9657ca383..c9bb30c909 100644 --- a/src/engines/SessionCore/control/turnLifecycle.ts +++ b/src/engines/SessionCore/control/turnLifecycle.ts @@ -184,8 +184,17 @@ export function beginTurnDispatch(sessionId: string): number { * turns, org-coordinator dispatches) and confirms a pending dispatch. * Never downgrades "stopping" — a late running ack must not cancel a Stop. */ -export function markTurnRunning(sessionId: string): void { +export function markTurnRunning( + sessionId: string, + options: { generation?: number } = {} +): void { const state = getState(sessionId); + if ( + options.generation !== undefined && + options.generation !== state.generation + ) { + return; + } if (state.phase === "working" || state.phase === "stopping") return; if (state.phase === "idle") { state.generation += 1; @@ -288,6 +297,15 @@ export function getLastTurnTerminal( return stateBySession.get(sessionId)?.lastTerminal ?? null; } +/** Release all retained lifecycle state when a session is permanently removed. */ +export function clearTurnLifecycleSession(sessionId: string): void { + const state = stateBySession.get(sessionId); + if (!state) return; + clearDeadman(state); + stateBySession.delete(sessionId); + bumpSignal(); +} + export function resetTurnLifecycleForTests(): void { for (const state of stateBySession.values()) { clearDeadman(state); diff --git a/src/engines/SessionCore/hooks/session/useQueueDispatch.ts b/src/engines/SessionCore/hooks/session/useQueueDispatch.ts index c1610d3557..3ad7f54668 100644 --- a/src/engines/SessionCore/hooks/session/useQueueDispatch.ts +++ b/src/engines/SessionCore/hooks/session/useQueueDispatch.ts @@ -66,7 +66,6 @@ import { queueEditingAtom, queueFlushRequestAtom, } from "@src/store/ui/messageQueueAtom"; -import { invokeTauri } from "@src/util/platform/tauri/init"; import { resolveModelForMessage } from "@src/util/session/resolveModelForMessage"; import { selectionFromSession } from "@src/util/session/selectionFromSession"; import { @@ -117,10 +116,9 @@ async function getBackendDispatchVerdict( ): Promise { try { if (isCliSession(sessionId)) { - const status = (await invokeTauri("cli_agent_status", { sessionId })) as { - status?: string; - } | null; - return classifyBackendSessionStatus(status?.status); + // CLI finality is push-owned by CliTurnLifecycleCoordinator. Re-reading + // status here would reintroduce one polling RPC per queued turn. + return "ready"; } if (isAgentSession(sessionId)) { const meta = await getSession(sessionId); diff --git a/src/engines/SessionCore/services/SessionService.ts b/src/engines/SessionCore/services/SessionService.ts index cab9bdf942..b05c09a701 100644 --- a/src/engines/SessionCore/services/SessionService.ts +++ b/src/engines/SessionCore/services/SessionService.ts @@ -24,6 +24,7 @@ import { respondQuestion, sessionLaunch, } from "@src/api/tauri/agent"; +import { rpc } from "@src/api/tauri/rpc"; import { ROUTES } from "@src/config/routes"; import { getAdapterForSession } from "@src/engines/SessionCore/sync/types"; import { @@ -257,10 +258,7 @@ export const SessionService = { }; } - const session = await invokeTauri<{ - sessionId: string; - status: string; - } | null>("cli_agent_status", { sessionId }); + const session = await rpc.cli.status({ sessionId }); if (!session) { throw new Error(`CLI session not found: ${sessionId}`); diff --git a/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts b/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts new file mode 100644 index 0000000000..2aece764bd --- /dev/null +++ b/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts @@ -0,0 +1,70 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +import { sendCliMessage } from "../cliTransport"; + +const mocks = vi.hoisted(() => ({ + enterIntervention: vi.fn(), + message: vi.fn(), + registerReceipt: vi.fn(), +})); + +vi.mock("@src/api/tauri/agent", () => ({ + enterAgentOrgSessionIntervention: mocks.enterIntervention, +})); +vi.mock("@src/api/tauri/rpc", () => ({ + rpc: { cli: { message: mocks.message, cancel: vi.fn() } }, +})); +vi.mock("@src/hooks/cliSession/cliTurnLifecycleCoordinator", () => ({ + cliTurnLifecycleCoordinator: { registerReceipt: mocks.registerReceipt }, +})); + +describe("sendCliMessage acceptance boundary", () => { + beforeEach(() => { + vi.clearAllMocks(); + mocks.message.mockResolvedValue({ + sessionId: "cliagent-worker", + turnIntentId: "intent-1", + status: "running", + }); + mocks.enterIntervention.mockResolvedValue(undefined); + }); + + it("resolves from the receipt without status or history reconciliation", async () => { + await expect( + sendCliMessage({ + sessionId: "cliagent-worker", + content: "continue", + turnIntentId: "intent-1", + clientMessageId: "message-1", + }) + ).resolves.toBeUndefined(); + + expect(mocks.message).toHaveBeenCalledWith({ + sessionId: "cliagent-worker", + content: "continue", + turnIntentId: "intent-1", + clientMessageId: "message-1", + }); + expect(mocks.registerReceipt).toHaveBeenCalledWith({ + sessionId: "cliagent-worker", + turnIntentId: "intent-1", + status: "running", + }); + }); + + it("rejects when the backend command rejects", async () => { + mocks.message.mockRejectedValue(new Error("ipc unavailable")); + + await expect( + sendCliMessage({ + sessionId: "cliagent-worker", + content: "retry", + isResume: true, + turnIntentId: "intent-2", + clientMessageId: "message-2", + }) + ).rejects.toThrow("ipc unavailable"); + + expect(mocks.registerReceipt).not.toHaveBeenCalled(); + }); +}); diff --git a/src/engines/SessionCore/sync/adapters/cli/cliHistory.ts b/src/engines/SessionCore/sync/adapters/cli/cliHistory.ts index 88b69c5d9e..82b8fc51da 100644 --- a/src/engines/SessionCore/sync/adapters/cli/cliHistory.ts +++ b/src/engines/SessionCore/sync/adapters/cli/cliHistory.ts @@ -1,5 +1,6 @@ -import { convertFileSrc, invoke as tauriInvoke } from "@tauri-apps/api/core"; +import { convertFileSrc } from "@tauri-apps/api/core"; +import { rpc } from "@src/api/tauri/rpc"; import type { SessionEvent } from "@src/engines/SessionCore/core/types"; import { processChunksRust } from "@src/engines/SessionCore/ingestion/rustBridge"; import { createLogger } from "@src/hooks/logger"; @@ -34,9 +35,7 @@ export async function loadCliHistory( sessionId: string, signal: AbortSignal ): Promise { - const chunks = await tauriInvoke("cli_agent_chunks", { - sessionId, - }); + const chunks = (await rpc.cli.chunks({ sessionId })) as ActivityChunk[]; if (signal.aborted || !Array.isArray(chunks)) return []; const events = await processChunksRust(chunks, sessionId); if (signal.aborted) return []; @@ -49,10 +48,9 @@ export async function postLoadCliSession( ): Promise { const result: PostLoadResult = {}; try { - const storedSession = await tauriInvoke( - "cli_agent_status", - { sessionId } - ); + const storedSession = (await rpc.cli.status({ + sessionId, + })) as StoredSession | null; if (signal.aborted || !storedSession) return result; registerSessionTranscriptSource(sessionId, storedSession.transcriptSource); diff --git a/src/engines/SessionCore/sync/adapters/cli/cliLifecycle.ts b/src/engines/SessionCore/sync/adapters/cli/cliLifecycle.ts index 0119336a7c..f2f5bd377e 100644 --- a/src/engines/SessionCore/sync/adapters/cli/cliLifecycle.ts +++ b/src/engines/SessionCore/sync/adapters/cli/cliLifecycle.ts @@ -1,10 +1,9 @@ -import { invoke as tauriInvoke } from "@tauri-apps/api/core"; - -import { confirmTurnRunning } from "@src/engines/SessionCore/control/turnLifecycle"; -import { loadSessionAtom } from "@src/engines/SessionCore/core/atoms"; +import { + type TurnTerminalStatus, + markTurnRunning, +} from "@src/engines/SessionCore/control/turnLifecycle"; import { isTurnBlockingRuntimeEvent } from "@src/engines/SessionCore/core/runningEventGate"; import { eventStoreProxy } from "@src/engines/SessionCore/core/store/EventStoreProxy"; -import type { SessionEvent } from "@src/engines/SessionCore/core/types"; import { createLogger } from "@src/hooks/logger"; import { setSessionRuntimeStatusAtom } from "@src/store/session/cliSessionStatusAtom"; import type { CliSessionStatus } from "@src/types/session/session"; @@ -13,16 +12,8 @@ import { isStoreInitialized, } from "@src/util/core/state/instrumentedStore"; -import { isNativeTranscriptSession } from "../../nativeTranscriptReconcile"; -import { loadCliHistory } from "./cliHistory"; - const log = createLogger("CliAdapter"); -export type CliStatusResponse = { - status?: CliSessionStatus; - updatedAt?: string; -}; - const CLI_TERMINAL_STATUSES = new Set([ "completed", "failed", @@ -33,56 +24,12 @@ const CLI_TERMINAL_STATUSES = new Set([ "archived", ]); -export const protectedRunningTurnBySession = new Map< - string, - { content: string; startedAt: number } ->(); - export function isCliTerminalStatus( status: CliSessionStatus | undefined ): status is CliSessionStatus { return status !== undefined && CLI_TERMINAL_STATUSES.has(status); } -export async function readCliStatus( - sessionId: string -): Promise { - return (await tauriInvoke("cli_agent_status", { - sessionId, - })) as CliStatusResponse | null; -} - -export async function waitForCliRunBoundary( - sessionId: string, - previousStatus: CliStatusResponse | null -): Promise { - const deadline = Date.now() + 15_000; - const previousUpdatedAt = previousStatus?.updatedAt; - const previousWasTerminal = isCliTerminalStatus(previousStatus?.status); - let lastStatus: CliStatusResponse | null = null; - while (Date.now() < deadline) { - lastStatus = await readCliStatus(sessionId); - const hasNewStatus = - !previousUpdatedAt || lastStatus?.updatedAt !== previousUpdatedAt; - const hasDurableBoundary = - Boolean(previousUpdatedAt) && lastStatus?.updatedAt !== previousUpdatedAt; - if (lastStatus?.status === "running" && hasNewStatus) { - return lastStatus; - } - if ( - isCliTerminalStatus(lastStatus?.status) && - (hasDurableBoundary || !previousWasTerminal) - ) { - return lastStatus; - } - await new Promise((resolve) => setTimeout(resolve, 100)); - } - - throw new Error( - `CLI run boundary was not observed for ${sessionId}; lastStatus=${JSON.stringify(lastStatus)}` - ); -} - async function closeObservedCliTerminalEvents( sessionId: string, status: CliSessionStatus @@ -111,111 +58,40 @@ async function closeObservedCliTerminalEvents( ); } -export function markCliRuntimeRunning(sessionId: string): void { - // FSM running ack is visibility-independent: the dispatch reserved the - // turn, so promote it to "working" even for background sessions. - confirmTurnRunning(sessionId); +export function markCliRuntimeRunning( + sessionId: string, + generation?: number +): void { + markTurnRunning(sessionId, { generation }); if (!isStoreInitialized()) return; - const store = getInstrumentedStore(); - store.set(setSessionRuntimeStatusAtom, { + getInstrumentedStore().set(setSessionRuntimeStatusAtom, { sessionId, status: "running", source: "sync", }); } -export function isProtectedCliTurnTerminal( - sessionId: string, - status: CliSessionStatus | undefined -): boolean { - return ( - isCliTerminalStatus(status) && protectedRunningTurnBySession.has(sessionId) - ); -} - export function markObservedCliTerminalStatus( sessionId: string, status: CliSessionStatus | undefined ): void { if (!isCliTerminalStatus(status) || !isStoreInitialized()) return; - if (isProtectedCliTurnTerminal(sessionId, status)) return; - const store = getInstrumentedStore(); - store.set(setSessionRuntimeStatusAtom, { sessionId, status, source: "sync" }); + getInstrumentedStore().set(setSessionRuntimeStatusAtom, { + sessionId, + status, + source: "sync", + }); void closeObservedCliTerminalEvents(sessionId, status).catch((error) => { log.warn("[cliAdapter] failed to close terminal CLI events:", error); }); } -export async function waitForCliTerminalBoundary( - sessionId: string, - previousUpdatedAt: string | null | undefined, - timeoutMs = 90_000 -): Promise { - const deadline = Date.now() + timeoutMs; - let lastStatus: CliStatusResponse | null = null; - while (Date.now() < deadline) { - lastStatus = await readCliStatus(sessionId); - const hasNewStatus = - !previousUpdatedAt || lastStatus?.updatedAt !== previousUpdatedAt; - if (hasNewStatus && isCliTerminalStatus(lastStatus?.status)) { - return lastStatus; - } - await new Promise((resolve) => setTimeout(resolve, 250)); - } - return lastStatus; -} - -async function refreshLoadedCliHistory( - sessionId: string -): Promise { - if (!isStoreInitialized()) return []; - const events = await loadCliHistory(sessionId, new AbortController().signal); - if (events.length === 0) return events; - // Native-transcript sessions render the live turn from in-memory events - // only. The replay is read here purely to observe persistence for the send - // handshake; terminal reconcile remains the single on-screen replacement. - if (!isNativeTranscriptSession(sessionId)) { - await eventStoreProxy.mergeEvents(events, sessionId); - getInstrumentedStore().set(loadSessionAtom, { sessionId, events }); - } - return events; -} - -function eventContainsText(event: SessionEvent, text: string): boolean { - return JSON.stringify(event).includes(text); -} - -export async function waitForPersistedCliUserEvent( - sessionId: string, - content: string -): Promise { - const deadline = Date.now() + 15_000; - let lastEventCount = 0; - while (Date.now() < deadline) { - const events = await refreshLoadedCliHistory(sessionId); - lastEventCount = events.length; - if (events.some((event) => eventContainsText(event, content))) { - return events; - } - await new Promise((resolve) => setTimeout(resolve, 250)); - } - throw new Error( - `CLI user event was not persisted for ${sessionId}; eventCount=${lastEventCount}` - ); -} - -export function hasRuntimeOutputAfterUserEvent( - events: SessionEvent[], - content: string -): boolean { - let userIndex = -1; - for (let index = events.length - 1; index >= 0; index -= 1) { - const event = events[index]; - if (event.source === "user" && eventContainsText(event, content)) { - userIndex = index; - break; - } +export function cliTerminalStatus( + status: CliSessionStatus +): TurnTerminalStatus { + if (status === "failed" || status === "error" || status === "timeout") { + return "failed"; } - if (userIndex < 0) return false; - return events.slice(userIndex + 1).some((event) => event.source !== "user"); + if (status === "cancelled" || status === "abandoned") return "cancelled"; + return "completed"; } diff --git a/src/engines/SessionCore/sync/adapters/cli/cliTransport.ts b/src/engines/SessionCore/sync/adapters/cli/cliTransport.ts index 52d486fb17..b19d70e490 100644 --- a/src/engines/SessionCore/sync/adapters/cli/cliTransport.ts +++ b/src/engines/SessionCore/sync/adapters/cli/cliTransport.ts @@ -1,25 +1,13 @@ -import { invoke as tauriInvoke } from "@tauri-apps/api/core"; - import { enterAgentOrgSessionIntervention } from "@src/api/tauri/agent"; import type { CancelReason } from "@src/api/tauri/agent/session"; -import { - getTurnGeneration, - markTurnTerminal, - toTurnTerminalStatus, -} from "@src/engines/SessionCore/control/turnLifecycle"; +import { rpc } from "@src/api/tauri/rpc"; +import { cliTurnLifecycleCoordinator } from "@src/hooks/cliSession/cliTurnLifecycleCoordinator"; import type { AdapterSendInput } from "../../types"; -import { - hasRuntimeOutputAfterUserEvent, - isCliTerminalStatus, - markCliRuntimeRunning, - markObservedCliTerminalStatus, - protectedRunningTurnBySession, - readCliStatus, - waitForCliRunBoundary, - waitForCliTerminalBoundary, - waitForPersistedCliUserEvent, -} from "./cliLifecycle"; + +function newMessageId(): string { + return crypto.randomUUID(); +} export async function sendCliMessage(input: AdapterSendInput): Promise { const { @@ -35,67 +23,27 @@ export async function sendCliMessage(input: AdapterSendInput): Promise { if (!isResume && content.trim()) { await enterAgentOrgSessionIntervention(sessionId); } - const previousStatus = await readCliStatus(sessionId); - protectedRunningTurnBySession.set(sessionId, { - content, - startedAt: Date.now(), - }); - markCliRuntimeRunning(sessionId); - try { - await tauriInvoke("cli_agent_message", { - sessionId, - content, - ...(model ? { model } : {}), - ...(accountId ? { accountId } : {}), - ...(mode ? { mode } : {}), - ...(imageDataUrls && imageDataUrls.length > 0 - ? { images: imageDataUrls } - : {}), - ...(adeContext ? { ideContext: adeContext } : {}), - }); - } catch (error) { - protectedRunningTurnBySession.delete(sessionId); - throw error; - } - const acceptedStatus = await waitForCliRunBoundary(sessionId, previousStatus); - markCliRuntimeRunning(sessionId); - // Capture this dispatch's generation so a late terminal can never close a - // newer turn. - const dispatchGeneration = getTurnGeneration(sessionId); - const persistedEvents = await waitForPersistedCliUserEvent( + const turnIntentId = input.turnIntentId ?? newMessageId(); + const clientMessageId = input.clientMessageId ?? newMessageId(); + const receipt = await rpc.cli.message({ sessionId, - content - ); - const acceptedTerminalIsCurrentTurn = - isCliTerminalStatus(acceptedStatus?.status) && - hasRuntimeOutputAfterUserEvent(persistedEvents, content); - if (acceptedTerminalIsCurrentTurn) { - protectedRunningTurnBySession.delete(sessionId); - markObservedCliTerminalStatus(sessionId, acceptedStatus.status); - markTurnTerminal( - sessionId, - toTurnTerminalStatus(acceptedStatus?.status ?? "completed"), - { generation: dispatchGeneration } - ); - return; - } - - void waitForCliTerminalBoundary( - sessionId, - acceptedStatus?.updatedAt ?? previousStatus?.updatedAt - ).then((terminalStatus) => { - if (!isCliTerminalStatus(terminalStatus?.status)) return; - protectedRunningTurnBySession.delete(sessionId); - markObservedCliTerminalStatus(sessionId, terminalStatus.status); - markTurnTerminal(sessionId, toTurnTerminalStatus(terminalStatus.status), { - generation: dispatchGeneration, - }); + content, + turnIntentId, + clientMessageId, + ...(model ? { model } : {}), + ...(accountId ? { accountId } : {}), + ...(mode ? { mode } : {}), + ...(imageDataUrls && imageDataUrls.length > 0 + ? { images: imageDataUrls } + : {}), + ...(adeContext ? { ideContext: adeContext } : {}), }); + cliTurnLifecycleCoordinator.registerReceipt(receipt); } export async function stopCliSession( sessionId: string, reason: CancelReason ): Promise { - await tauriInvoke("cli_agent_cancel", { sessionId, reason }); + await rpc.cli.cancel({ sessionId, reason }); } diff --git a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts index d15b1d63fb..902fda5ff3 100644 --- a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts +++ b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts @@ -1,8 +1,4 @@ import type { MergeStatus } from "@src/api/tauri/rpc/schemas/validation"; -import { - markTurnTerminal, - toTurnTerminalStatus, -} from "@src/engines/SessionCore/control/turnLifecycle"; import { eventStoreProxy } from "@src/engines/SessionCore/core/store/EventStoreProxy"; import type { SessionEvent } from "@src/engines/SessionCore/core/types"; import { normalizeChunkRust } from "@src/engines/SessionCore/ingestion/rustBridge"; @@ -47,11 +43,7 @@ import { capStreamContent } from "../shared/subagentTracking"; import type { AgentWSEvent, PermissionRequestEvent } from "../shared/types"; import { isCliTerminalStatus, - isProtectedCliTurnTerminal, - markCliRuntimeRunning, markObservedCliTerminalStatus, - protectedRunningTurnBySession, - readCliStatus, } from "./cliLifecycle"; import { buildCliStreamingEvent } from "./streamingEvent"; @@ -75,7 +67,6 @@ export function createCliEventHandler( let thinkStreamId = ""; let thinkStartedAt = ""; let observedTerminalStatus: CliSessionStatus | undefined; - let finalAssistantSettleTimer: ReturnType | undefined; const finalizedStreamEventIds = new Set(); const toolCallDeltaBuffers = new Map< number, @@ -123,49 +114,9 @@ export function createCliEventHandler( function reconcileTerminalEventsIfNeeded(): void { if (!observedTerminalStatus) return; - clearFinalAssistantSettleTimer(); markObservedCliTerminalStatus(sessionId, observedTerminalStatus); } - function scheduleFinalAssistantSettleFallback(): void { - if (observedTerminalStatus) return; - const protectedTurn = protectedRunningTurnBySession.get(sessionId); - if (!protectedTurn) return; - clearFinalAssistantSettleTimer(); - finalAssistantSettleTimer = setTimeout(() => { - if (observedTerminalStatus) return; - if (protectedRunningTurnBySession.get(sessionId) !== protectedTurn) { - return; - } - void readCliStatus(sessionId) - .then((statusResponse) => { - if (observedTerminalStatus) return; - if (protectedRunningTurnBySession.get(sessionId) !== protectedTurn) { - return; - } - const terminalStatus = isCliTerminalStatus(statusResponse?.status) - ? statusResponse.status - : "completed"; - observedTerminalStatus = terminalStatus; - protectedRunningTurnBySession.delete(sessionId); - callbacks.onStatusChange?.(terminalStatus); - markObservedCliTerminalStatus(sessionId, terminalStatus); - markTurnTerminal(sessionId, toTurnTerminalStatus(terminalStatus)); - clearMessageStream(); - clearThinkingStream(); - clearToolCallDeltaBuffers(); - setStreamingMode(false); - callbacks.onAgentComplete?.(); - }) - .catch((error) => { - log.warn( - "[CliAdapter] final assistant settle fallback failed:", - error - ); - }); - }, 1_500); - } - function asString(value: unknown): string | undefined { return typeof value === "string" && value.length > 0 ? value : undefined; } @@ -303,10 +254,6 @@ export function createCliEventHandler( } function handleActivity(chunk: ActivityChunk): void { - if (!observedTerminalStatus) { - callbacks.onStatusChange?.("running"); - } - if ( chunk.function === "user_message" && (chunk.action_type === "raw" || chunk.action_type === "raw_event") @@ -388,11 +335,8 @@ export function createCliEventHandler( // Final message/thinking chunks replace any TS typewriter placeholder. if (isMessageType || isThinkingType) { const tempId = isMessageType ? msgStreamId : thinkStreamId; - const isFinalAssistantMessage = - isMessageType && chunk.result?.is_full_content === true; const reconcileAfterFinalEvent = () => { reconcileTerminalEventsIfNeeded(); - if (isFinalAssistantMessage) scheduleFinalAssistantSettleFallback(); }; normalizeChunkRust(chunk, sessionId) .then((event) => { @@ -455,7 +399,6 @@ export function createCliEventHandler( clearMessageStream(); const reconcileAfterCompleteMessage = () => { reconcileTerminalEventsIfNeeded(); - scheduleFinalAssistantSettleFallback(); }; if (tsTempId && tsTempId !== completeEvent.id) { eventStoreProxy @@ -485,20 +428,12 @@ export function createCliEventHandler( } } - function handleStatusChange(status: string, errorMessage?: string): void { + function handleStatusChange(status: string): void { const terminalStatus = isCliTerminalStatus(status as CliSessionStatus) ? (status as CliSessionStatus) : undefined; - if (isProtectedCliTurnTerminal(sessionId, terminalStatus)) { - markCliRuntimeRunning(sessionId); - return; - } - - callbacks.onStatusChange?.(status, errorMessage); - if (terminalStatus) { observedTerminalStatus = terminalStatus; - clearFinalAssistantSettleTimer(); clearMessageStream(); clearThinkingStream(); clearToolCallDeltaBuffers(); @@ -510,7 +445,6 @@ export function createCliEventHandler( if (isSessionRuntimeExecuting(status)) { observedTerminalStatus = undefined; - protectedRunningTurnBySession.delete(sessionId); cancelled = false; } } @@ -557,10 +491,7 @@ export function createCliEventHandler( } else if (raw.type === "agent:streaming_complete") { handleStreamingComplete(raw); } else if (raw.type === "code_session.status_changed") { - handleStatusChange( - raw.status as string, - raw.error_message as string | undefined - ); + handleStatusChange(raw.status as string); } else if (raw.type === "code_session.token_usage_updated") { const total = raw.total_tokens; if (typeof total === "number") callbacks.onTokenUpdate?.(total); @@ -590,14 +521,12 @@ export function createCliEventHandler( }, reset(): void { - clearFinalAssistantSettleTimer(); clearMessageStream(); clearThinkingStream(); clearToolCallDeltaBuffers(); observedTerminalStatus = undefined; cancelled = false; - streaming = false; - eventStoreProxy.setStreaming(false, sessionId); + setStreamingMode(false); }, get isStreaming(): boolean { diff --git a/src/hooks/cliSession/cliTurnLifecycleCoordinator.test.ts b/src/hooks/cliSession/cliTurnLifecycleCoordinator.test.ts new file mode 100644 index 0000000000..58dfd488ea --- /dev/null +++ b/src/hooks/cliSession/cliTurnLifecycleCoordinator.test.ts @@ -0,0 +1,177 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +import { + publishTurnIntentDispatch, + resetTurnIntentDispatchLifecycleForTests, +} from "@src/engines/SessionCore/control/turnIntentDispatchLifecycle"; +import { + beginTurnDispatch, + getTurnPhase, + resetTurnLifecycleForTests, +} from "@src/engines/SessionCore/control/turnLifecycle"; + +import { CliTurnLifecycleCoordinator } from "./cliTurnLifecycleCoordinator"; + +describe("CliTurnLifecycleCoordinator", () => { + beforeEach(() => { + resetTurnLifecycleForTests(); + resetTurnIntentDispatchLifecycleForTests(); + }); + + it("binds current terminals to the dispatched generation and drops stale intents", () => { + const coordinator = new CliTurnLifecycleCoordinator(vi.fn()); + const sessionId = "cliagent-current"; + const generation = beginTurnDispatch(sessionId); + publishTurnIntentDispatch("intent-current", { sessionId, generation }); + + expect( + coordinator.handleStatus({ + sessionId, + status: "running", + turnIntentId: "intent-current", + }) + ).toBe(true); + expect(getTurnPhase(sessionId)).toBe("working"); + + expect( + coordinator.handleStatus({ + sessionId, + status: "completed", + turnIntentId: "intent-old", + }) + ).toBe(false); + expect(getTurnPhase(sessionId)).toBe("working"); + + expect( + coordinator.handleStatus({ + sessionId, + status: "completed", + turnIntentId: "intent-current", + }) + ).toBe(true); + expect(getTurnPhase(sessionId)).toBe("idle"); + expect(coordinator.activeSessionCount).toBe(0); + }); + + it("creates and closes a local generation for an unknown cross-window intent", () => { + const coordinator = new CliTurnLifecycleCoordinator(vi.fn()); + const sessionId = "cliagent-cross-window"; + + coordinator.handleStatus({ + sessionId, + status: "running", + turnIntentId: "remote-intent", + }); + expect(getTurnPhase(sessionId)).toBe("working"); + + coordinator.handleStatus({ + sessionId, + status: "failed", + turnIntentId: "remote-intent", + }); + expect(getTurnPhase(sessionId)).toBe("idle"); + + expect( + coordinator.handleStatus({ + sessionId, + status: "running", + turnIntentId: "remote-intent", + }) + ).toBe(false); + expect(getTurnPhase(sessionId)).toBe("idle"); + }); + + it("does not let an unattributed terminal close a tracked intent", () => { + const coordinator = new CliTurnLifecycleCoordinator(vi.fn()); + const sessionId = "cliagent-attributed"; + coordinator.handleStatus({ + sessionId, + status: "running", + turnIntentId: "intent-attributed", + }); + + expect(coordinator.handleStatus({ sessionId, status: "completed" })).toBe( + false + ); + expect(getTurnPhase(sessionId)).toBe("working"); + + coordinator.clearSession(sessionId); + expect(getTurnPhase(sessionId)).toBe("idle"); + expect(coordinator.activeSessionCount).toBe(0); + }); + + it("rejects excess unknown intents without evicting active generations", () => { + const coordinator = new CliTurnLifecycleCoordinator(vi.fn()); + for (let index = 0; index < 256; index += 1) { + expect( + coordinator.handleStatus({ + sessionId: `cliagent-cap-${index}`, + status: "running", + turnIntentId: `intent-cap-${index}`, + }) + ).toBe(true); + } + + expect( + coordinator.handleStatus({ + sessionId: "cliagent-cap-overflow", + status: "running", + turnIntentId: "intent-cap-overflow", + }) + ).toBe(false); + expect(coordinator.activeSessionCount).toBe(256); + + expect( + coordinator.handleStatus({ + sessionId: "cliagent-cap-0", + status: "completed", + turnIntentId: "intent-cap-0", + }) + ).toBe(true); + expect(coordinator.activeSessionCount).toBe(255); + }); + + it("shares one batch request across concurrent reconnect and focus triggers", async () => { + let resolveBatch!: (value: never[]) => void; + const loadBatch = vi.fn( + () => + new Promise((resolve) => { + resolveBatch = resolve; + }) + ); + const coordinator = new CliTurnLifecycleCoordinator(loadBatch); + coordinator.handleStatus({ + sessionId: "cliagent-reconcile", + status: "running", + turnIntentId: "intent-reconcile", + }); + + const first = coordinator.reconcile(); + const second = coordinator.reconcile(); + + expect(first).toBe(second); + expect(loadBatch).toHaveBeenCalledOnce(); + expect(loadBatch).toHaveBeenCalledWith({ + sessionIds: ["cliagent-reconcile"], + }); + resolveBatch([]); + await first; + }); + + it("does not reconcile while the document is hidden", async () => { + const loadBatch = vi.fn(async () => []); + const coordinator = new CliTurnLifecycleCoordinator(loadBatch); + coordinator.handleStatus({ + sessionId: "cliagent-hidden", + status: "running", + turnIntentId: "intent-hidden", + }); + const originalDocument = globalThis.document; + vi.stubGlobal("document", { visibilityState: "hidden" }); + + await coordinator.reconcile(); + + expect(loadBatch).not.toHaveBeenCalled(); + vi.stubGlobal("document", originalDocument); + }); +}); diff --git a/src/hooks/cliSession/cliTurnLifecycleCoordinator.ts b/src/hooks/cliSession/cliTurnLifecycleCoordinator.ts new file mode 100644 index 0000000000..6616b72387 --- /dev/null +++ b/src/hooks/cliSession/cliTurnLifecycleCoordinator.ts @@ -0,0 +1,213 @@ +import { rpc } from "@src/api/tauri/rpc"; +import { getTurnIntentDispatch } from "@src/engines/SessionCore/control/turnIntentDispatchLifecycle"; +import { + beginTurnDispatch, + clearTurnLifecycleSession, + getTurnGeneration, + getTurnPhase, + markTurnTerminal, +} from "@src/engines/SessionCore/control/turnLifecycle"; +import { + cliTerminalStatus, + isCliTerminalStatus, + markCliRuntimeRunning, + markObservedCliTerminalStatus, +} from "@src/engines/SessionCore/sync/adapters/cli/cliLifecycle"; +import { + type SessionStatus, + sessionsAtom, + updateSessionStatus, +} from "@src/store/session"; +import type { CliSessionStatus } from "@src/types/session/session"; +import { + getInstrumentedStore, + isStoreInitialized, +} from "@src/util/core/state/instrumentedStore"; +import { isCliSession } from "@src/util/session/sessionDispatch"; + +export interface CliRunReceipt { + sessionId: string; + turnIntentId: string; + status: string; +} + +export interface CliLifecycleStatus { + sessionId: string; + status: string; + updatedAt?: string; + turnIntentId?: string; +} + +interface ActiveCliTurn { + turnIntentId: string; + generation: number; +} + +const MAX_ACTIVE_SESSIONS = 256; +const MAX_RECENT_TERMINALS = 256; +const RECONCILE_STATUSES = new Set([ + "pending", + "running", + "waiting_for_user", + "waiting_for_funds", +]); + +type BatchLoader = (input: { + sessionIds: string[]; +}) => Promise; + +export class CliTurnLifecycleCoordinator { + private readonly activeBySession = new Map(); + private readonly recentTerminalIntents = new Set(); + private reconcilePromise: Promise | null = null; + + constructor(private readonly loadStatusBatch: BatchLoader) {} + + get activeSessionCount(): number { + return this.activeBySession.size; + } + + registerReceipt(receipt: CliRunReceipt): void { + this.handleStatus({ + sessionId: receipt.sessionId, + turnIntentId: receipt.turnIntentId, + status: receipt.status, + }); + } + + handleStatus(event: CliLifecycleStatus): boolean { + if (!isCliSession(event.sessionId)) return false; + const status = event.status as CliSessionStatus; + const turnIntentId = event.turnIntentId; + const existing = this.activeBySession.get(event.sessionId); + + if (status === "running") { + if (!turnIntentId || this.recentTerminalIntents.has(turnIntentId)) + return false; + if (existing?.turnIntentId === turnIntentId) { + markCliRuntimeRunning(event.sessionId, existing.generation); + return false; + } + + const dispatch = getTurnIntentDispatch(turnIntentId); + if (dispatch && dispatch.sessionId !== event.sessionId) return false; + if (dispatch && existing && dispatch.generation < existing.generation) { + return false; + } + // Unknown cross-window intents need a retained generation so their + // terminal cannot close a newer turn. Never evict another active + // session merely to admit one beyond the bounded coordinator capacity. + if ( + !dispatch && + !existing && + this.activeBySession.size >= MAX_ACTIVE_SESSIONS + ) { + return false; + } + + const generation = + dispatch?.generation ?? beginTurnDispatch(event.sessionId); + markCliRuntimeRunning(event.sessionId, generation); + if (getTurnGeneration(event.sessionId) !== generation) return false; + this.setActive(event.sessionId, { turnIntentId, generation }); + return true; + } + + if (!isCliTerminalStatus(status)) return false; + if (!turnIntentId && existing) return false; + if (turnIntentId && this.recentTerminalIntents.has(turnIntentId)) + return false; + if (turnIntentId && existing && existing.turnIntentId !== turnIntentId) { + return false; + } + + const dispatch = turnIntentId + ? getTurnIntentDispatch(turnIntentId) + : undefined; + if (dispatch && dispatch.sessionId !== event.sessionId) return false; + const generation = existing?.generation ?? dispatch?.generation; + markTurnTerminal(event.sessionId, cliTerminalStatus(status), { + generation, + }); + markObservedCliTerminalStatus(event.sessionId, status); + if (isStoreInitialized()) { + updateSessionStatus(event.sessionId, status as SessionStatus); + } + this.activeBySession.delete(event.sessionId); + if (turnIntentId) this.rememberTerminal(turnIntentId); + return true; + } + + reconcile(): Promise { + if ( + typeof document !== "undefined" && + document.visibilityState === "hidden" + ) { + return Promise.resolve(); + } + if (this.reconcilePromise) return this.reconcilePromise; + + const sessionIds = this.collectReconcileSessionIds(); + if (sessionIds.length === 0) return Promise.resolve(); + this.reconcilePromise = this.loadStatusBatch({ sessionIds }) + .then((statuses) => { + for (const status of statuses) this.handleStatus(status); + }) + .finally(() => { + this.reconcilePromise = null; + }); + return this.reconcilePromise; + } + + clearSession(sessionId: string): void { + this.activeBySession.delete(sessionId); + clearTurnLifecycleSession(sessionId); + } + + resetForTests(): void { + this.activeBySession.clear(); + this.recentTerminalIntents.clear(); + this.reconcilePromise = null; + } + + private collectReconcileSessionIds(): string[] { + const ids = new Set(this.activeBySession.keys()); + if (isStoreInitialized()) { + for (const session of getInstrumentedStore().get(sessionsAtom)) { + if ( + isCliSession(session.session_id) && + (RECONCILE_STATUSES.has(session.status) || + getTurnPhase(session.session_id) !== "idle") + ) { + ids.add(session.session_id); + } + } + } + return [...ids].slice(0, MAX_ACTIVE_SESSIONS); + } + + private setActive(sessionId: string, active: ActiveCliTurn): void { + this.activeBySession.delete(sessionId); + this.activeBySession.set(sessionId, active); + } + + private rememberTerminal(turnIntentId: string): void { + this.recentTerminalIntents.delete(turnIntentId); + this.recentTerminalIntents.add(turnIntentId); + while (this.recentTerminalIntents.size > MAX_RECENT_TERMINALS) { + const oldest = this.recentTerminalIntents.values().next().value as + | string + | undefined; + if (!oldest) break; + this.recentTerminalIntents.delete(oldest); + } + } +} + +export const cliTurnLifecycleCoordinator = new CliTurnLifecycleCoordinator( + rpc.cli.statusBatch +); + +export function clearCliTurnLifecycleSession(sessionId: string): void { + cliTurnLifecycleCoordinator.clearSession(sessionId); +} diff --git a/src/hooks/cliSession/useBackgroundSessionMonitor.ts b/src/hooks/cliSession/useBackgroundSessionMonitor.ts index 6293476af3..abb1b8a463 100644 --- a/src/hooks/cliSession/useBackgroundSessionMonitor.ts +++ b/src/hooks/cliSession/useBackgroundSessionMonitor.ts @@ -1,15 +1,15 @@ /** * useBackgroundSessionMonitor Hook * - * Listens for WebSocket status changes on background ("fire and forget") - * CLI sessions and delivers system notifications + in-app toasts when - * they complete or fail. + * Owns the single window-level CLI lifecycle status subscription. It routes + * every CLI status through the global coordinator and additionally delivers + * notifications for background ("fire and forget") sessions. * * This hook runs at the app root level (via GlobalSessionSync) so it is * always active, regardless of which view the user is on. * - * It complements the cliAdapter sync, which only tracks the *active* session. - * This hook watches ALL background sessions globally. + * Active adapters remain responsible for transcript/UI mirroring only; turn + * finality for active and background sessions is owned here. */ import { useAtomValue } from "jotai"; import { useEffect, useRef } from "react"; @@ -20,14 +20,11 @@ import { notifyTaskCompletion, } from "@src/api/services/notification"; import Message from "@src/components/Message"; -import { - markTurnTerminal, - toTurnTerminalStatus, -} from "@src/engines/SessionCore/control/turnLifecycle"; -import { type SessionStatus, updateSessionStatus } from "@src/store/session"; import { notificationSettingsAtom } from "@src/store/ui/notificationAtom"; import { isTerminalStatus } from "@src/types/session/session"; +import { cliTurnLifecycleCoordinator } from "./cliTurnLifecycleCoordinator"; + interface BackgroundStatusMessage { type: "code_session.status_changed"; session_id: string; @@ -36,6 +33,7 @@ interface BackgroundStatusMessage { session_name?: string; error_message?: string; exit_code?: number; + turn_intent_id?: string; } export function useBackgroundSessionMonitor(): void { @@ -52,15 +50,18 @@ export function useBackgroundSessionMonitor(): void { const unsubscribe = wsClient.on("code_session.status_changed", (raw) => { const msg = raw as unknown as BackgroundStatusMessage; + const applied = cliTurnLifecycleCoordinator.handleStatus({ + sessionId: msg.session_id, + status: msg.status, + turnIntentId: msg.turn_intent_id, + }); if (!msg.background) return; if (!isTerminalStatus(msg.status)) return; + if (!applied) return; const sessionName = msg.session_name || "Background session"; - markTurnTerminal(msg.session_id, toTurnTerminalStatus(msg.status)); - updateSessionStatus(msg.session_id, msg.status as SessionStatus); - if (msg.status === "completed") { notifyTaskCompletion( `"${sessionName}" completed — ready for review`, @@ -95,6 +96,21 @@ export function useBackgroundSessionMonitor(): void { } }); - return unsubscribe; + const reconcile = () => { + void cliTurnLifecycleCoordinator.reconcile(); + }; + const unsubscribeConnected = wsClient.on("connected", reconcile); + const handleVisibilityChange = () => { + if (document.visibilityState === "visible") reconcile(); + }; + document.addEventListener("visibilitychange", handleVisibilityChange); + window.addEventListener("focus", reconcile); + + return () => { + unsubscribe(); + unsubscribeConnected(); + document.removeEventListener("visibilitychange", handleVisibilityChange); + window.removeEventListener("focus", reconcile); + }; }, []); } diff --git a/src/scaffold/NavigationSidebar/connectors/useWorkstationSidebarHandlers.ts b/src/scaffold/NavigationSidebar/connectors/useWorkstationSidebarHandlers.ts index 2cf7100d62..070213b672 100644 --- a/src/scaffold/NavigationSidebar/connectors/useWorkstationSidebarHandlers.ts +++ b/src/scaffold/NavigationSidebar/connectors/useWorkstationSidebarHandlers.ts @@ -30,6 +30,7 @@ import { isSessionTaggedToCloudOrg, sessionOrgTagsAtom, } from "@src/features/TeamCollaboration/sessionOrgTagsAtom"; +import { clearCliTurnLifecycleSession } from "@src/hooks/cliSession/cliTurnLifecycleCoordinator"; import { createLogger } from "@src/hooks/logger"; import type { GoToNewSessionOptions } from "@src/hooks/navigation/useAppNavigation"; import type { NavigationMenuItem } from "@src/scaffold/NavigationSidebar/components/NavigationMenu/config"; @@ -198,6 +199,7 @@ export function useWorkstationSidebarHandlers({ } if (isCliSession(sessionId)) { await invokeTauri("cli_agent_delete", { sessionId }); + clearCliTurnLifecycleSession(sessionId); } else if (isHumanSession(sessionId)) { await deleteHumanSession(sessionId); } else { From 9d0f0b1a0b756b3c8e84644048fa981648ea7c2a Mon Sep 17 00:00:00 2001 From: hanafish <1106510024@qq.com> Date: Wed, 29 Jul 2026 10:56:22 +0800 Subject: [PATCH 2/3] fix(rpc): preserve current lifecycle contracts --- .../sync/adapters/cli/__tests__/cliTransport.test.ts | 2 ++ .../SessionCore/sync/adapters/cli/createCliEventHandler.ts | 1 + 2 files changed, 3 insertions(+) diff --git a/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts b/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts index 2aece764bd..5b6438444c 100644 --- a/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts +++ b/src/engines/SessionCore/sync/adapters/cli/__tests__/cliTransport.test.ts @@ -36,6 +36,7 @@ describe("sendCliMessage acceptance boundary", () => { content: "continue", turnIntentId: "intent-1", clientMessageId: "message-1", + turnIntentSource: "user_submit", }) ).resolves.toBeUndefined(); @@ -62,6 +63,7 @@ describe("sendCliMessage acceptance boundary", () => { isResume: true, turnIntentId: "intent-2", clientMessageId: "message-2", + turnIntentSource: "user_submit", }) ).rejects.toThrow("ipc unavailable"); diff --git a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts index 902fda5ff3..d3f279ec23 100644 --- a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts +++ b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts @@ -67,6 +67,7 @@ export function createCliEventHandler( let thinkStreamId = ""; let thinkStartedAt = ""; let observedTerminalStatus: CliSessionStatus | undefined; + let finalAssistantSettleTimer: ReturnType | undefined; const finalizedStreamEventIds = new Set(); const toolCallDeltaBuffers = new Map< number, From e746dd4c6f3f6aa9e2a3356ea3eb44e08c2fea62 Mon Sep 17 00:00:00 2001 From: hanafish <1106510024@qq.com> Date: Wed, 29 Jul 2026 13:27:13 +0800 Subject: [PATCH 3/3] fix(cli): preserve runner lifecycle contracts --- .../src/agent_sessions/cli/commands/run.rs | 29 ++++++------- .../cli/persistence/session_crud.rs | 1 + .../cli/session_runner/session.rs | 2 +- .../session_runner/session/transport_acp.rs | 4 +- .../session/transport_app_server.rs | 42 +++++++++---------- .../session/transport_standard.rs | 27 ++++++------ .../adapters/cli/createCliEventHandler.ts | 7 ---- 7 files changed, 51 insertions(+), 61 deletions(-) diff --git a/src-tauri/src/agent_sessions/cli/commands/run.rs b/src-tauri/src/agent_sessions/cli/commands/run.rs index 855edb308b..5b51135611 100644 --- a/src-tauri/src/agent_sessions/cli/commands/run.rs +++ b/src-tauri/src/agent_sessions/cli/commands/run.rs @@ -38,11 +38,9 @@ fn inject_ide_context_into_prompt(user_input: &str, ide_context: Option<&IdeCont /// tab close). Non-TUI sessions and already-terminal rows are left alone. #[tauri::command] pub async fn cli_agent_tui_release(session_id: String) -> Result { - tokio::task::spawn_blocking(move || { - super::super::tui_bridge::release_tui_session(&session_id) - }) - .await - .map_err(|e| format!("Task error: {}", e))? + tokio::task::spawn_blocking(move || super::super::tui_bridge::release_tui_session(&session_id)) + .await + .map_err(|e| format!("Task error: {}", e))? } /// Run a code session (spawn CLI agent in background). @@ -192,8 +190,7 @@ async fn cli_agent_run_internal( "status": "failed", "error_message": e, }); - failed_msg["turn_intent_id"] = - serde_json::Value::String(runner_turn_intent_id.clone()); + failed_msg["turn_intent_id"] = serde_json::Value::String(runner_turn_intent_id.clone()); crate::api::websocket_handler::broadcast(failed_msg.to_string()); } // Remove finished entry from RUNNING_SESSIONS to prevent unbounded growth @@ -320,15 +317,15 @@ pub async fn cli_agent_message( .await .map_err(|err| format!("Task error: {err}"))??; let cli_resume_id = account_scoped_resume_id.or_else(|| { - if account_id - .as_deref() - .is_some_and(|new_account_id| session.account_id.as_deref() != Some(new_account_id)) - { - None - } else { - fresh_cli_session_id - } - }); + if account_id + .as_deref() + .is_some_and(|new_account_id| session.account_id.as_deref() != Some(new_account_id)) + { + None + } else { + fresh_cli_session_id + } + }); tracing::info!( session_id = %session_id, diff --git a/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs b/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs index aca1996321..391b1e9201 100644 --- a/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs +++ b/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs @@ -289,6 +289,7 @@ pub fn accept_cli_turn( session_id, turn_intent_id, Some(client_message_id), + None, session_persistence::turn_intents::TurnIntentSource::UserSubmit, session_persistence::turn_intents::TurnIntentStatus::Queued, ) diff --git a/src-tauri/src/agent_sessions/cli/session_runner/session.rs b/src-tauri/src/agent_sessions/cli/session_runner/session.rs index 6fea7beea7..05553b6740 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/session.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/session.rs @@ -21,7 +21,7 @@ use key_vault::key_store::{KeyService, ModelType, KEY_SERVICE}; use super::super::launch_profile_store::resolve_cli_launch_profile; use super::super::persistence; -use super::super::types::{KeySource, SessionStatus}; +use super::super::types::KeySource; use super::command::{ build_command_with_launch_profile, launch_profile_env, CliCommandBuildRequest, }; diff --git a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_acp.rs b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_acp.rs index 9993c3a60e..d858bf8f25 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_acp.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_acp.rs @@ -93,9 +93,9 @@ pub(super) async fn run_acp_branch( let timeout_result = tokio::time::timeout(session_timeout, async { while let Some(chunk) = chunk_rx.recv().await { if let Some(snap_id) = &pre_message_snapshot_id { - snapshot_cli_file_edit(&session_id, snap_id, &chunk, &snapshot_working_dir); + snapshot_cli_file_edit(&session_id, snap_id, &chunk, &snapshot_working_dir).await; } - emit_chunk(&chunk, &session_id, sequence); + emit_chunk(&chunk, &session_id, sequence).await; } }) .await; diff --git a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_app_server.rs b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_app_server.rs index 84db0fb82c..604a75d3e7 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_app_server.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_app_server.rs @@ -7,9 +7,9 @@ use tokio::process::Child; use crate::api::websocket_handler; +use super::super::super::persistence; use super::super::helpers::{emit_chunk, snapshot_cli_file_edit}; use super::super::launch_profiles::ResolvedCliLaunchProfile; -use super::super::super::persistence; pub(super) struct AppServerOutcome { pub(super) exit_code: i32, @@ -71,11 +71,9 @@ pub(super) async fn run_codex_app_server_branch( if cli_session_id_out.is_none() { if let Some(ref tid) = chunk.thread_id { cli_session_id_out = Some(tid.clone()); - if let Err(err) = persistence::update_cli_session_id_for_account( - &session_id, - account_id, - tid, - ) { + if let Err(err) = + persistence::update_cli_session_id_for_account(&session_id, account_id, tid) + { tracing::warn!( "[CodeSession] Failed to bind early cli_session_id: {}", err @@ -92,9 +90,9 @@ pub(super) async fn run_codex_app_server_branch( } } if let Some(snap_id) = &pre_message_snapshot_id { - snapshot_cli_file_edit(&session_id, snap_id, &chunk, &snapshot_working_dir); + snapshot_cli_file_edit(&session_id, snap_id, &chunk, &snapshot_working_dir).await; } - emit_chunk(&chunk, &session_id, sequence); + emit_chunk(&chunk, &session_id, sequence).await; } }) .await; @@ -106,21 +104,19 @@ pub(super) async fn run_codex_app_server_branch( codex_app_server_turn_ok = result.turn_status != "failed"; if let Some(ref usage) = result.usage { let round_model = usage.model.as_deref().or(model); - if let Err(err) = - session_persistence::token_usage::insert_token_usage_record( - &session_id, - "code", - round_model, - account_id, - usage.input_tokens as i64, - usage.output_tokens as i64, - usage.cache_read_tokens as i64, - usage.cache_write_tokens as i64, - usage.total_tokens as i64, - 0, - None, - ) - { + if let Err(err) = session_persistence::token_usage::insert_token_usage_record( + &session_id, + "code", + round_model, + account_id, + usage.input_tokens as i64, + usage.output_tokens as i64, + usage.cache_read_tokens as i64, + usage.cache_write_tokens as i64, + usage.total_tokens as i64, + 0, + None, + ) { tracing::warn!( "[CodeSession] Failed to insert per-round token usage: {}", err diff --git a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_standard.rs b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_standard.rs index 72146a9682..e3b412887d 100644 --- a/src-tauri/src/agent_sessions/cli/session_runner/session/transport_standard.rs +++ b/src-tauri/src/agent_sessions/cli/session_runner/session/transport_standard.rs @@ -13,8 +13,13 @@ use tokio::sync::Mutex; use crate::api::websocket_handler; use key_vault::key_store::ModelType; +use super::super::super::persistence; +use super::super::super::persistence::CodeSession; +use super::super::super::types::SessionStatus; use super::super::command::create_parser; -use super::super::helpers::{clear_live_status, emit_chunk, flush_and_broadcast, snapshot_cli_file_edit}; +use super::super::helpers::{ + clear_live_status, emit_chunk, flush_and_broadcast, snapshot_cli_file_edit, +}; use super::super::oauth_setup::{ is_cli_chunk_replay_unsafe, is_cli_oauth_failure_message, is_cli_oauth_stderr_retry_candidate, is_retryable_cli_oauth_failure_chunk, is_retryable_overloaded_chunk, @@ -23,9 +28,6 @@ use super::super::plan_approval::{ create_plan_content_from_chunk, is_successful_mode_tool, plan_candidate_path_from_chunk, register_cli_plan_approval, register_synthetic_cli_plan_approval, }; -use super::super::super::persistence; -use super::super::super::persistence::CodeSession; -use super::super::super::types::SessionStatus; const CLI_PLAN_GATE_NATURAL_EXIT_GRACE_SECS: u64 = 45; @@ -164,7 +166,8 @@ pub(super) async fn run_standard_branch( snap_id, &chunk, &snapshot_working_dir, - ); + ) + .await; } if is_successful_mode_tool(&chunk, "enter_plan_mode") { cli_plan_active = true; @@ -188,7 +191,7 @@ pub(super) async fn run_standard_branch( .await { Ok(plan_chunk) => { - emit_chunk(&plan_chunk, &session_id, sequence); + emit_chunk(&plan_chunk, &session_id, sequence).await; cli_plan_registered_this_turn = true; cli_plan_approval_gate_triggered = true; } @@ -217,7 +220,7 @@ pub(super) async fn run_standard_branch( .await { Ok(plan_chunk) => { - emit_chunk(&plan_chunk, &session_id, sequence); + emit_chunk(&plan_chunk, &session_id, sequence).await; cli_plan_registered_this_turn = true; cli_plan_approval_gate_triggered = true; } @@ -242,7 +245,7 @@ pub(super) async fn run_standard_branch( .await { Ok(plan_chunk) => { - emit_chunk(&plan_chunk, &session_id, sequence); + emit_chunk(&plan_chunk, &session_id, sequence).await; cli_plan_registered_this_turn = true; cli_plan_approval_gate_triggered = true; } @@ -263,7 +266,7 @@ pub(super) async fn run_standard_branch( } cli_plan_active = false; } - emit_chunk(&chunk, &session_id, sequence); + emit_chunk(&chunk, &session_id, sequence).await; if cli_plan_approval_gate_triggered && !cli_plan_gate_announced { cli_plan_gate_announced = true; tracing::info!( @@ -275,7 +278,7 @@ pub(super) async fn run_standard_branch( // instead of holding Stop for up to the 45s drain window // while the child process winds down. The final // status_changed after child exit is idempotent. - flush_and_broadcast(&session_id); + flush_and_broadcast(&session_id).await; // The plan card supersedes any hook-derived // waiting/working entry for this turn. clear_live_status( @@ -416,9 +419,9 @@ pub(super) async fn run_standard_branch( break; } if let Some(snap_id) = &pre_message_snapshot_id { - snapshot_cli_file_edit(&session_id, snap_id, chunk, &snapshot_working_dir); + snapshot_cli_file_edit(&session_id, snap_id, chunk, &snapshot_working_dir).await; } - emit_chunk(chunk, &session_id, sequence); + emit_chunk(chunk, &session_id, sequence).await; } } diff --git a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts index d3f279ec23..f6495942ca 100644 --- a/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts +++ b/src/engines/SessionCore/sync/adapters/cli/createCliEventHandler.ts @@ -67,7 +67,6 @@ export function createCliEventHandler( let thinkStreamId = ""; let thinkStartedAt = ""; let observedTerminalStatus: CliSessionStatus | undefined; - let finalAssistantSettleTimer: ReturnType | undefined; const finalizedStreamEventIds = new Set(); const toolCallDeltaBuffers = new Map< number, @@ -97,12 +96,6 @@ export function createCliEventHandler( toolCallDeltaBuffers.clear(); } - function clearFinalAssistantSettleTimer(): void { - if (!finalAssistantSettleTimer) return; - clearTimeout(finalAssistantSettleTimer); - finalAssistantSettleTimer = undefined; - } - function rememberFinalizedStreamEvent(eventId: string): void { if (finalizedStreamEventIds.has(eventId)) return; while (finalizedStreamEventIds.size >= MAX_FINALIZED_STREAM_IDS) {