From 8e84f507c3a7f1934379ee355f58e97560548829 Mon Sep 17 00:00:00 2001 From: Neonforge <48338160+Neonforge98@users.noreply.github.com> Date: Tue, 25 Aug 2026 07:34:59 +0800 Subject: [PATCH] fix(runtime): make turn acceptance retry-safe Atomically claim native and CLI turn intents before runtime, transcript, configuration, process, or scheduler side effects. Return durable lifecycle receipts, preserve Project composer-to-WorkItemRun identity, and expose status recovery after IPC response loss. --- .../agent-core/src/core/session/scheduler.rs | 163 +++++-- .../src/foundation/session_bridge.rs | 50 +++ .../session/message/project_bootstrap.rs | 17 +- .../state/commands/session/message/send.rs | 302 ++++++++++--- .../src/state/commands/session/mod.rs | 43 ++ .../src/work_run_service/enqueue.rs | 37 +- .../src/work_run_service/mod.rs | 25 +- .../src/work_run_service/read.rs | 51 ++- .../src/work_run_service/tests.rs | 58 ++- .../src/agent_core_bridge.rs | 25 ++ .../session-persistence/src/turn_intents.rs | 108 ++++- .../src/agent_sessions/cli/commands/run.rs | 405 +++++++++++++++--- .../cli/persistence/session_crud.rs | 91 +++- src-tauri/src/commands/handler_list.inc | 1 + .../schemas/__tests__/cliRunReceipt.test.ts | 35 ++ src/api/tauri/rpc/schemas/cli.ts | 14 +- 16 files changed, 1253 insertions(+), 172 deletions(-) create mode 100644 src/api/tauri/rpc/schemas/__tests__/cliRunReceipt.test.ts diff --git a/src-tauri/crates/agent-core/src/core/session/scheduler.rs b/src-tauri/crates/agent-core/src/core/session/scheduler.rs index 29e6fed79a..36ee16416a 100644 --- a/src-tauri/crates/agent-core/src/core/session/scheduler.rs +++ b/src-tauri/crates/agent-core/src/core/session/scheduler.rs @@ -38,7 +38,7 @@ use futures::FutureExt; use std::any::Any; -use std::collections::HashSet; +use std::collections::HashMap; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::Arc; use tokio::sync::{mpsc, Mutex as TokioMutex}; @@ -163,6 +163,40 @@ pub struct EnqueueResult { pub duplicate: bool, } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ClientMessageClaim { + Claimed, + SameIntentDuplicate, + DifferentIntentDuplicate, +} + +fn claim_client_message( + owners: &mut HashMap, + client_message_id: &str, + turn_intent_id: &str, +) -> ClientMessageClaim { + match owners.get(client_message_id) { + Some(owner_turn_intent_id) if owner_turn_intent_id == turn_intent_id => { + ClientMessageClaim::SameIntentDuplicate + } + Some(_) => ClientMessageClaim::DifferentIntentDuplicate, + None => { + owners.insert(client_message_id.to_string(), turn_intent_id.to_string()); + ClientMessageClaim::Claimed + } + } +} + +fn release_client_message( + owners: &mut HashMap, + client_message_id: &str, + turn_intent_id: &str, +) { + if owners.get(client_message_id).map(String::as_str) == Some(turn_intent_id) { + owners.remove(client_message_id); + } +} + // ============================================ // DialogScheduler // ============================================ @@ -198,7 +232,7 @@ pub struct DialogScheduler { processing: Arc, /// Whether the job the worker is currently executing is a [`ScheduledKind::Turn`]. processing_turn: Arc, - client_message_ids: Arc>>, + client_message_owners: Arc>>, } impl DialogScheduler { @@ -217,7 +251,7 @@ impl DialogScheduler { generation: Arc::new(AtomicU64::new(0)), processing: Arc::new(std::sync::atomic::AtomicBool::new(false)), processing_turn: Arc::new(std::sync::atomic::AtomicBool::new(false)), - client_message_ids: Arc::new(TokioMutex::new(HashSet::new())), + client_message_owners: Arc::new(TokioMutex::new(HashMap::new())), } } /// Ensure the worker is spawned and return a reference to the sender. @@ -237,7 +271,7 @@ impl DialogScheduler { generation: Arc::clone(&self.generation), processing: Arc::clone(&self.processing), processing_turn: Arc::clone(&self.processing_turn), - client_message_ids: Arc::clone(&self.client_message_ids), + client_message_owners: Arc::clone(&self.client_message_owners), }; tokio::spawn(worker.run()); @@ -262,22 +296,37 @@ impl DialogScheduler { let message_id = msg.message_id.clone(); if let Some(client_message_id) = msg.client_message_id.as_ref() { - let mut ids = self.client_message_ids.lock().await; - if !ids.insert(client_message_id.clone()) { - // This request minted its own durable intent before enqueue, - // but an equivalent client message is already queued/running. - // The scheduler is the single authority that knows the request - // was coalesced, so it also closes that new intent here. - crate::foundation::session_bridge::update_turn_intent_status( - &self.session_id, - &msg.turn_intent_id, - crate::foundation::session_bridge::TurnIntentBridgeStatus::Coalesced, - ); - return Ok(EnqueueResult { - message_id, - queue_position: 0, - duplicate: true, - }); + let claim = { + let mut owners = self.client_message_owners.lock().await; + claim_client_message(&mut owners, client_message_id, &msg.turn_intent_id) + }; + match claim { + ClientMessageClaim::Claimed => {} + ClientMessageClaim::SameIntentDuplicate => { + // A response-loss retry owns the same durable intent as the + // queued/running message. Keep that row in its original + // state so the retry can reconcile its exact receipt. + return Ok(EnqueueResult { + message_id, + queue_position: 0, + duplicate: true, + }); + } + ClientMessageClaim::DifferentIntentDuplicate => { + // A distinct logical intent reused an in-flight idempotency + // key. It will never create a round, so close only the new + // intent as coalesced and preserve the original owner. + crate::foundation::session_bridge::update_turn_intent_status( + &self.session_id, + &msg.turn_intent_id, + crate::foundation::session_bridge::TurnIntentBridgeStatus::Coalesced, + ); + return Ok(EnqueueResult { + message_id, + queue_position: 0, + duplicate: true, + }); + } } } @@ -297,10 +346,12 @@ impl DialogScheduler { Err(mpsc::error::TrySendError::Full(rejected)) => { self.pending.fetch_sub(1, Ordering::Relaxed); if let Some(client_message_id) = rejected.client_message_id.as_ref() { - self.client_message_ids - .lock() - .await - .remove(client_message_id); + let mut owners = self.client_message_owners.lock().await; + release_client_message( + &mut owners, + client_message_id, + &rejected.turn_intent_id, + ); } crate::foundation::session_bridge::update_turn_intent_status( &self.session_id, @@ -315,10 +366,12 @@ impl DialogScheduler { Err(mpsc::error::TrySendError::Closed(rejected)) => { self.pending.fetch_sub(1, Ordering::Relaxed); if let Some(client_message_id) = rejected.client_message_id.as_ref() { - self.client_message_ids - .lock() - .await - .remove(client_message_id); + let mut owners = self.client_message_owners.lock().await; + release_client_message( + &mut owners, + client_message_id, + &rejected.turn_intent_id, + ); } crate::foundation::session_bridge::update_turn_intent_status( &self.session_id, @@ -339,8 +392,8 @@ impl DialogScheduler { pub fn invalidate_pending(&self) { self.generation.fetch_add(1, Ordering::AcqRel); self.pending.store(0, Ordering::Release); - if let Ok(mut ids) = self.client_message_ids.try_lock() { - ids.clear(); + if let Ok(mut owners) = self.client_message_owners.try_lock() { + owners.clear(); } // Lifecycle: every still-queued / optimistic intent for this // session walks to `stale`. The worker drops queued-but-stale @@ -393,7 +446,7 @@ struct WorkerTask { generation: Arc, processing: Arc, processing_turn: Arc, - client_message_ids: Arc>>, + client_message_owners: Arc>>, } impl WorkerTask { @@ -418,10 +471,8 @@ impl WorkerTask { crate::foundation::session_bridge::TurnIntentBridgeStatus::Stale, ); if let Some(client_message_id) = msg.client_message_id.as_ref() { - self.client_message_ids - .lock() - .await - .remove(client_message_id); + let mut owners = self.client_message_owners.lock().await; + release_client_message(&mut owners, client_message_id, &msg.turn_intent_id); } self.broadcast_idle_status(); continue; @@ -553,10 +604,8 @@ impl WorkerTask { } if let Some(client_message_id) = client_message_id.as_ref() { - self.client_message_ids - .lock() - .await - .remove(client_message_id); + let mut owners = self.client_message_owners.lock().await; + release_client_message(&mut owners, client_message_id, &turn_intent_id); } self.processing_turn.store(false, Ordering::Relaxed); self.processing.store(false, Ordering::Relaxed); @@ -583,6 +632,42 @@ mod tests { use super::*; use std::sync::atomic::{AtomicUsize, Ordering}; + #[test] + fn client_message_claim_preserves_the_exact_intent_owner() { + let mut owners = HashMap::new(); + assert_eq!( + claim_client_message(&mut owners, "client-1", "intent-owner"), + ClientMessageClaim::Claimed + ); + assert_eq!( + claim_client_message(&mut owners, "client-1", "intent-owner"), + ClientMessageClaim::SameIntentDuplicate + ); + assert_eq!( + claim_client_message(&mut owners, "client-1", "intent-other"), + ClientMessageClaim::DifferentIntentDuplicate + ); + assert_eq!( + owners.get("client-1").map(String::as_str), + Some("intent-owner") + ); + + // A stale worker must not release a newer owner that claimed the same + // idempotency key after invalidation. + owners.clear(); + assert_eq!( + claim_client_message(&mut owners, "client-1", "intent-new"), + ClientMessageClaim::Claimed + ); + release_client_message(&mut owners, "client-1", "intent-owner"); + assert_eq!( + owners.get("client-1").map(String::as_str), + Some("intent-new") + ); + release_client_message(&mut owners, "client-1", "intent-new"); + assert!(!owners.contains_key("client-1")); + } + #[tokio::test] async fn invalidated_pending_message_is_skipped() { let scheduler = DialogScheduler::new("session-a", 8); diff --git a/src-tauri/crates/agent-core/src/foundation/session_bridge.rs b/src-tauri/crates/agent-core/src/foundation/session_bridge.rs index 116293e068..8c3b24f4fc 100644 --- a/src-tauri/crates/agent-core/src/foundation/session_bridge.rs +++ b/src-tauri/crates/agent-core/src/foundation/session_bridge.rs @@ -499,6 +499,22 @@ pub type UpsertTurnIntentFn = fn( status: TurnIntentBridgeStatus, ); +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TurnIntentBridgeClaim { + pub duplicate: bool, + pub status: TurnIntentBridgeStatus, + pub client_message_id: Option, +} + +pub type ClaimTurnIntentFn = fn( + session_id: &str, + turn_intent_id: &str, + client_message_id: Option<&str>, + org_run_id: Option<&str>, + source: TurnIntentBridgeSource, + status: TurnIntentBridgeStatus, +) -> Result; + pub type UpdateTurnIntentStatusFn = fn(session_id: &str, turn_intent_id: &str, new_status: TurnIntentBridgeStatus); @@ -508,6 +524,7 @@ pub type GetTurnIntentStatusFn = pub type MarkPendingTurnIntentsStaleFn = fn(session_id: &str); static UPSERT_TURN_INTENT: OnceLock = OnceLock::new(); +static CLAIM_TURN_INTENT: OnceLock = OnceLock::new(); static UPDATE_TURN_INTENT_STATUS: OnceLock = OnceLock::new(); static GET_TURN_INTENT_STATUS: OnceLock = OnceLock::new(); static MARK_PENDING_TURN_INTENTS_STALE: OnceLock = OnceLock::new(); @@ -516,6 +533,10 @@ pub fn register_upsert_turn_intent(implementation: UpsertTurnIntentFn) { let _ = UPSERT_TURN_INTENT.set(implementation); } +pub fn register_claim_turn_intent(implementation: ClaimTurnIntentFn) { + let _ = CLAIM_TURN_INTENT.set(implementation); +} + pub fn register_update_turn_intent_status(implementation: UpdateTurnIntentStatusFn) { let _ = UPDATE_TURN_INTENT_STATUS.set(implementation); } @@ -553,6 +574,35 @@ pub fn upsert_turn_intent( } } +/// Atomically reserve an exact logical intent before any turn side effects. +/// +/// Unlike the best-effort projection helper above, this is an acceptance +/// boundary: missing registration or persistence failure is returned to the +/// caller so it cannot continue with an unowned execution. +pub fn claim_turn_intent( + session_id: &str, + turn_intent_id: &str, + client_message_id: Option<&str>, + org_run_id: Option<&str>, + source: TurnIntentBridgeSource, + status: TurnIntentBridgeStatus, +) -> Result { + if session_id.is_empty() || turn_intent_id.is_empty() { + return Err("turn intent claim requires session_id and turn_intent_id".to_string()); + } + let implementation = CLAIM_TURN_INTENT + .get() + .ok_or_else(|| "turn intent claim persistence is not registered".to_string())?; + implementation( + session_id, + turn_intent_id, + client_message_id, + org_run_id, + source, + status, + ) +} + /// Patch the status of an existing lifecycle row. Illegal transitions are /// silently rejected by the implementation — callers do not need to handle /// the error case. diff --git a/src-tauri/crates/agent-core/src/state/commands/session/message/project_bootstrap.rs b/src-tauri/crates/agent-core/src/state/commands/session/message/project_bootstrap.rs index 2ba60d350b..f1c63a63bb 100644 --- a/src-tauri/crates/agent-core/src/state/commands/session/message/project_bootstrap.rs +++ b/src-tauri/crates/agent-core/src/state/commands/session/message/project_bootstrap.rs @@ -13,8 +13,8 @@ use crate::foundation::session_bridge::TurnIntentBridgeSource; use project_management::projects::types::{ - EnqueueWorkItemRunRequest, WorkItemRun, WorkItemRunTarget, WorkItemRunTargetSnapshot, - WorkItemRunTrigger, PERSONAL_ORG_ID, + EnqueueWorkItemRunRequest, WorkItemRunTarget, WorkItemRunTargetSnapshot, WorkItemRunTrigger, + PERSONAL_ORG_ID, }; /// Bootstrap called from the message-accept path. Project mode is an explicit @@ -55,7 +55,7 @@ pub(super) async fn enqueue_project_turn_if_needed( turn_intent_id: &str, client_message_id: Option<&str>, source: TurnIntentBridgeSource, -) -> Result, String> { +) -> Result, String> { if content.trim().is_empty() || turn_intent_id.starts_with("wir_") { return Ok(None); } @@ -129,15 +129,18 @@ pub(super) async fn enqueue_project_turn_if_needed( "content": content, "displayText": display_text, "clientMessageId": client_message_id, + "originTurnIntentId": turn_intent_id, }), idempotency_key: format!("project-session-turn:{session_id}:{turn_intent_id}"), max_attempts: 3, parent_run_id: None, }; - tokio::task::spawn_blocking(move || project_management::work_run_service::enqueue(request)) - .await - .map_err(|err| format!("Project WorkItemRun enqueue worker failed: {err}"))? - .map(Some) + tokio::task::spawn_blocking(move || { + project_management::work_run_service::enqueue_with_receipt(request) + }) + .await + .map_err(|err| format!("Project WorkItemRun enqueue worker failed: {err}"))? + .map(Some) } /// Blocking core, also driven directly by the `Track this` command — diff --git a/src-tauri/crates/agent-core/src/state/commands/session/message/send.rs b/src-tauri/crates/agent-core/src/state/commands/session/message/send.rs index 183b8153fc..762e745e57 100644 --- a/src-tauri/crates/agent-core/src/state/commands/session/message/send.rs +++ b/src-tauri/crates/agent-core/src/state/commands/session/message/send.rs @@ -65,6 +65,112 @@ pub(super) fn terminal_intent_status_override( } } +fn duplicate_turn_intent_response( + session_id: String, + model: String, + turn_intent_id: &str, + client_message_id: Option<&str>, + status: crate::foundation::session_bridge::TurnIntentBridgeStatus, +) -> AgentResponse { + AgentResponse { + content: serde_json::json!({ + "queued": status.is_in_flight(), + "messageId": client_message_id.unwrap_or(turn_intent_id), + "queuePosition": 0, + "duplicate": true, + "turnIntentStatus": status.as_str(), + "effectiveTurnIntentId": turn_intent_id, + }) + .to_string(), + session_id, + model, + } +} + +fn project_turn_intent_response( + session_id: String, + model: String, + origin_turn_intent_id: &str, + client_message_id: Option<&str>, + effective_turn_intent_id: String, + status: &str, + duplicate: bool, +) -> AgentResponse { + AgentResponse { + content: serde_json::json!({ + "queued": matches!(status, "queued" | "running"), + "durableRunId": effective_turn_intent_id, + "effectiveTurnIntentId": effective_turn_intent_id, + "turnIntentStatus": status, + "messageId": client_message_id.unwrap_or(origin_turn_intent_id), + "queuePosition": 0, + "duplicate": duplicate, + }) + .to_string(), + session_id, + model, + } +} + +fn accepted_turn_intent_response( + session_id: String, + model: String, + turn_intent_id: &str, + message_id: String, + queue_position: usize, + duplicate: bool, + status: crate::foundation::session_bridge::TurnIntentBridgeStatus, +) -> AgentResponse { + AgentResponse { + content: serde_json::json!({ + "queued": status.is_in_flight(), + "messageId": message_id, + "queuePosition": queue_position, + "duplicate": duplicate, + "turnIntentStatus": status.as_str(), + "effectiveTurnIntentId": turn_intent_id, + }) + .to_string(), + session_id, + model, + } +} + +/// Reject a newly claimed ordinary turn when preparation exits before a +/// scheduler or steering boundary accepts ownership. +struct TurnIntentClaimGuard { + session_id: String, + turn_intent_id: String, + armed: bool, +} + +impl TurnIntentClaimGuard { + fn new(session_id: &str, turn_intent_id: &str) -> Self { + Self { + session_id: session_id.to_string(), + turn_intent_id: turn_intent_id.to_string(), + armed: true, + } + } + + fn disarm(&mut self) { + self.armed = false; + } +} + +impl Drop for TurnIntentClaimGuard { + fn drop(&mut self) { + if !self.armed { + return; + } + crate::foundation::session_bridge::update_turn_intent_status( + &self.session_id, + &self.turn_intent_id, + crate::foundation::session_bridge::TurnIntentBridgeStatus::Rejected, + ); + } +} + /// Implementation of agent_send_message. #[allow(clippy::too_many_arguments)] pub(crate) async fn send_message_impl( @@ -108,6 +214,63 @@ pub(crate) async fn send_message_impl( // ── 1. Resolve session identity (unified — single code path) ───────── let identity = resolve_session_identity(state, &session_id, overrides).await?; + // Project mode already reserves its durable execution through the + // WorkItemRun idempotency key. Route it before the ordinary turn-intent + // claim so one logical submission never acquires two execution owners. + if !is_resume { + super::project_bootstrap::ensure_project_root_work_item(&session_id, &content).await?; + if let Some(receipt) = super::project_bootstrap::enqueue_project_turn_if_needed( + &session_id, + &content, + display_text.as_deref(), + &effective_turn_intent_id, + client_message_id.as_deref(), + source, + ) + .await? + { + tracing::info!( + session_id = %session_id, + run_id = %receipt.run.id, + duplicate = receipt.duplicate, + "queued Project turn through durable WorkItem dispatcher" + ); + let turn_status = + project_management::work_run_service::turn_intent_status(receipt.run.status); + return Ok(project_turn_intent_response( + session_id, + identity.model, + &effective_turn_intent_id, + client_message_id.as_deref(), + receipt.run.id, + turn_status, + receipt.duplicate, + )); + } + } + + // Ordinary native turns claim their exact lifecycle row before goal, + // runtime, transcript, steering, or scheduler side effects. Exact IPC + // retries return the original durable state and cannot execute twice. + let intent_claim = crate::foundation::session_bridge::claim_turn_intent( + &session_id, + &effective_turn_intent_id, + client_message_id.as_deref(), + intent_org_run_id.as_deref(), + source, + crate::foundation::session_bridge::TurnIntentBridgeStatus::Queued, + )?; + if intent_claim.duplicate { + return Ok(duplicate_turn_intent_response( + session_id, + identity.model, + &effective_turn_intent_id, + intent_claim.client_message_id.as_deref(), + intent_claim.status, + )); + } + let mut intent_claim_guard = TurnIntentClaimGuard::new(&session_id, &effective_turn_intent_id); + // Goal loop: a real user submission becomes (or replaces) the // session's standing goal and resets the continuation counter. // `Queue`-sourced messages (goal continuations, queued flushes) and @@ -261,14 +424,6 @@ pub(crate) async fn send_message_impl( // part of accepting the control action. If the durable takeover row // cannot be written, do not inject a message that Wake may race. persist_direct_user_intervention(direct_user_intervention.clone()).await?; - crate::foundation::session_bridge::upsert_turn_intent( - &session_id, - &effective_turn_intent_id, - client_message_id.as_deref(), - effective_intent_org_run_id.as_deref(), - source, - crate::foundation::session_bridge::TurnIntentBridgeStatus::Queued, - ); session_handle .steering_queue .lock() @@ -291,6 +446,7 @@ pub(crate) async fn send_message_impl( }; if !reclaimed { + intent_claim_guard.disarm(); tracing::info!( "[agent_send_message] Steering message into active turn for session {} (intent={})", session_id, @@ -303,6 +459,8 @@ pub(crate) async fn send_message_impl( "messageId": effective_turn_intent_id, "queuePosition": 0, "duplicate": false, + "turnIntentStatus": "queued", + "effectiveTurnIntentId": effective_turn_intent_id, }) .to_string(), session_id, @@ -356,45 +514,6 @@ pub(crate) async fn send_message_impl( } } - // ── 4b. Project root WorkItem bootstrap (orgtrack/v1 §7.2) ────────── - // - // The first accepted non-empty submission of a Project session with - // no active WorkItem creates and links its root. Resumes replay an - // already-accepted submission, so they never bootstrap. - if !is_resume { - super::project_bootstrap::ensure_project_root_work_item(&session_id, &content).await?; - if let Some(run) = super::project_bootstrap::enqueue_project_turn_if_needed( - &session_id, - &content, - display_text.as_deref(), - &effective_turn_intent_id, - client_message_id.as_deref(), - source, - ) - .await? - { - tracing::info!( - session_id = %session_id, - run_id = %run.id, - "queued Project turn through durable WorkItem dispatcher" - ); - return Ok(AgentResponse { - content: serde_json::json!({ - "queued": true, - "durableRunId": run.id, - "messageId": client_message_id - .as_deref() - .unwrap_or(&effective_turn_intent_id), - "queuePosition": 0, - "duplicate": false, - }) - .to_string(), - session_id, - model: effective_model, - }); - } - } - // ── 5. Build the processing closure ────────────────────────────────── let sid_for_closure = session_id.clone(); let content_for_closure = content.clone(); @@ -707,6 +826,7 @@ pub(crate) async fn send_message_impl( .enqueue(msg) .await .map_err(|err| format!("Failed to enqueue message: {err}"))?; + intent_claim_guard.disarm(); tracing::info!( "[agent_send_message] Enqueued message {} at position {} for session {}", @@ -715,15 +835,83 @@ pub(crate) async fn send_message_impl( session_id ); - Ok(AgentResponse { - content: serde_json::json!({ - "queued": true, - "messageId": enqueue_result.message_id, - "queuePosition": enqueue_result.queue_position, - "duplicate": enqueue_result.duplicate, - }) - .to_string(), + let acknowledged_status = crate::foundation::session_bridge::get_turn_intent_status( + &session_id, + &effective_turn_intent_id, + ) + .unwrap_or(crate::foundation::session_bridge::TurnIntentBridgeStatus::Queued); + Ok(accepted_turn_intent_response( session_id, - model: effective_model, - }) + effective_model, + &effective_turn_intent_id, + enqueue_result.message_id, + enqueue_result.queue_position, + enqueue_result.duplicate, + acknowledged_status, + )) +} + +#[cfg(test)] +mod turn_intent_receipt_tests { + use super::{ + accepted_turn_intent_response, duplicate_turn_intent_response, project_turn_intent_response, + }; + use crate::foundation::session_bridge::TurnIntentBridgeStatus; + + #[test] + fn duplicate_receipt_preserves_the_exact_durable_status() { + let response = duplicate_turn_intent_response( + "session-1".to_string(), + "model-1".to_string(), + "intent-1", + Some("message-original"), + TurnIntentBridgeStatus::Completed, + ); + let payload: serde_json::Value = + serde_json::from_str(&response.content).expect("typed duplicate receipt"); + assert_eq!(payload["duplicate"], true); + assert_eq!(payload["queued"], false); + assert_eq!(payload["messageId"], "message-original"); + assert_eq!(payload["turnIntentStatus"], "completed"); + assert_eq!(payload["effectiveTurnIntentId"], "intent-1"); + } + + #[test] + fn accepted_receipt_exposes_the_effective_intent() { + let response = accepted_turn_intent_response( + "session-1".to_string(), + "model-1".to_string(), + "intent-1", + "message-1".to_string(), + 2, + false, + TurnIntentBridgeStatus::Queued, + ); + let payload: serde_json::Value = + serde_json::from_str(&response.content).expect("typed accepted receipt"); + assert_eq!(payload["duplicate"], false); + assert_eq!(payload["queued"], true); + assert_eq!(payload["queuePosition"], 2); + assert_eq!(payload["turnIntentStatus"], "queued"); + assert_eq!(payload["effectiveTurnIntentId"], "intent-1"); + } + + #[test] + fn project_receipt_keeps_origin_and_effective_intents_distinct() { + let response = project_turn_intent_response( + "session-1".to_string(), + "model-1".to_string(), + "intent-x", + Some("message-x"), + "wir-y".to_string(), + "queued", + true, + ); + let payload: serde_json::Value = + serde_json::from_str(&response.content).expect("Project receipt"); + assert_eq!(payload["messageId"], "message-x"); + assert_eq!(payload["effectiveTurnIntentId"], "wir-y"); + assert_eq!(payload["turnIntentStatus"], "queued"); + assert_eq!(payload["duplicate"], true); + } } diff --git a/src-tauri/crates/agent-core/src/state/commands/session/mod.rs b/src-tauri/crates/agent-core/src/state/commands/session/mod.rs index dee3d8cf7d..04bdc214cc 100644 --- a/src-tauri/crates/agent-core/src/state/commands/session/mod.rs +++ b/src-tauri/crates/agent-core/src/state/commands/session/mod.rs @@ -87,6 +87,49 @@ pub async fn agent_session_info( })) } +#[derive(Debug, Clone, serde::Serialize)] +#[serde(rename_all = "camelCase")] +pub struct TurnIntentStatusReceipt { + pub status: String, + pub effective_turn_intent_id: String, +} + +/// Read the durable owner and lifecycle for one exact logical turn. +/// +/// Project mode maps the composer intent to its WorkItemRun identity; ordinary +/// native and CLI turns retain their original intent id. +#[tauri::command] +pub async fn agent_turn_intent_status( + session_id: String, + turn_intent_id: String, +) -> Result, String> { + let lookup_session_id = session_id.clone(); + let lookup_turn_intent_id = turn_intent_id.clone(); + let project_run = tokio::task::spawn_blocking(move || { + project_management::work_run_service::find_project_session_turn( + &lookup_session_id, + &lookup_turn_intent_id, + ) + }) + .await + .map_err(|err| format!("Project turn receipt lookup worker failed: {err}"))??; + if let Some(run) = project_run { + return Ok(Some(TurnIntentStatusReceipt { + status: project_management::work_run_service::turn_intent_status(run.status) + .to_string(), + effective_turn_intent_id: run.id, + })); + } + + Ok( + crate::foundation::session_bridge::get_turn_intent_status(&session_id, &turn_intent_id) + .map(|status| TurnIntentStatusReceipt { + status: status.as_str().to_string(), + effective_turn_intent_id: turn_intent_id, + }), + ) +} + /// Remove a session (cleanup). #[tauri::command] pub async fn agent_session_remove( diff --git a/src-tauri/crates/project-management/src/work_run_service/enqueue.rs b/src-tauri/crates/project-management/src/work_run_service/enqueue.rs index 6341a855f3..14c45174ab 100644 --- a/src-tauri/crates/project-management/src/work_run_service/enqueue.rs +++ b/src-tauri/crates/project-management/src/work_run_service/enqueue.rs @@ -5,7 +5,7 @@ use crate::projects::io::helpers::{conn, now_ms}; use crate::projects::types::{EnqueueWorkItemRunRequest, WorkItemRun, WorkItemRunUsage}; use super::store::{append_audit, canonical_standalone_org_id, db, require_run, scope_key}; -use super::{error, DEFAULT_LEASE_MS, MAX_RUN_ATTEMPTS}; +use super::{error, EnqueueWorkItemRunReceipt, DEFAULT_LEASE_MS, MAX_RUN_ATTEMPTS}; #[derive(Debug)] struct WorkItemExecutionContext { @@ -203,6 +203,13 @@ fn hydrate_target_snapshot( /// returns the existing Run. Reusing the key with different content is a /// typed conflict. pub fn enqueue(request: EnqueueWorkItemRunRequest) -> Result { + Ok(enqueue_with_initial_delay(request, 0)?.run) +} + +/// Enqueue with an explicit idempotency receipt for response-loss recovery. +pub fn enqueue_with_receipt( + request: EnqueueWorkItemRunRequest, +) -> Result { enqueue_with_initial_delay(request, 0) } @@ -215,19 +222,19 @@ pub fn enqueue(request: EnqueueWorkItemRunRequest) -> Result Result { - enqueue_with_initial_delay(request, DEFAULT_LEASE_MS) + Ok(enqueue_with_initial_delay(request, DEFAULT_LEASE_MS)?.run) } fn enqueue_with_initial_delay( request: EnqueueWorkItemRunRequest, initial_delay_ms: i64, -) -> Result { +) -> Result { let mut connection = conn()?; let tx = db(connection.transaction_with_behavior(TransactionBehavior::Immediate))?; - let run = enqueue_in_transaction(&tx, request, initial_delay_ms)?; + let receipt = enqueue_in_transaction_with_receipt(&tx, request, initial_delay_ms)?; db(tx.commit())?; crate::projects::events::notify_work_item_dispatch_ready(); - Ok(run) + Ok(receipt) } /// Internal composition point for producers that must commit domain state and @@ -235,9 +242,17 @@ fn enqueue_with_initial_delay( /// The caller owns the surrounding `IMMEDIATE` transaction. pub(crate) fn enqueue_in_transaction( tx: &Transaction<'_>, - mut request: EnqueueWorkItemRunRequest, + request: EnqueueWorkItemRunRequest, initial_delay_ms: i64, ) -> Result { + Ok(enqueue_in_transaction_with_receipt(tx, request, initial_delay_ms)?.run) +} + +fn enqueue_in_transaction_with_receipt( + tx: &Transaction<'_>, + mut request: EnqueueWorkItemRunRequest, + initial_delay_ms: i64, +) -> Result { if request.work_item_id.trim().is_empty() || request.idempotency_key.trim().is_empty() { return Err(format!( "{}:work_item_id and idempotency_key are required", @@ -273,7 +288,10 @@ pub(crate) fn enqueue_in_transaction( request.idempotency_key )); } - return require_run(tx, &run_id); + return Ok(EnqueueWorkItemRunReceipt { + run: require_run(tx, &run_id)?, + duplicate: true, + }); } let attempt = if let Some(parent_run_id) = request.parent_run_id.as_deref() { @@ -364,5 +382,8 @@ pub(crate) fn enqueue_in_transaction( "attempt": attempt, }), )?; - require_run(tx, &run_id) + Ok(EnqueueWorkItemRunReceipt { + run: require_run(tx, &run_id)?, + duplicate: false, + }) } diff --git a/src-tauri/crates/project-management/src/work_run_service/mod.rs b/src-tauri/crates/project-management/src/work_run_service/mod.rs index 52e022f035..d111727cf0 100644 --- a/src-tauri/crates/project-management/src/work_run_service/mod.rs +++ b/src-tauri/crates/project-management/src/work_run_service/mod.rs @@ -24,10 +24,11 @@ pub use dispatch::{ has_claimable_dispatch, next_dispatch_due_at_ms, }; pub(crate) use enqueue::enqueue_in_transaction; -pub use enqueue::{enqueue, enqueue_for_inline_dispatch}; +pub use enqueue::{enqueue, enqueue_for_inline_dispatch, enqueue_with_receipt}; pub(crate) use read::read_in_transaction; pub use read::{ - latest_for_session, list_active_session_runs, list_for_work_item, read, routine_origin, + find_project_session_turn, latest_for_session, list_active_session_runs, list_for_work_item, + read, routine_origin, }; pub use terminal::{ classify_failure, mark_waiting, record_dispatch_failure, record_run_terminal, @@ -65,3 +66,23 @@ pub enum WorkItemRunTerminalOutcome { Failed, Cancelled, } + +#[derive(Debug, Clone)] +pub struct EnqueueWorkItemRunReceipt { + pub run: crate::projects::types::WorkItemRun, + pub duplicate: bool, +} + +/// Canonical projection from a WorkItemRun lifecycle to a turn receipt. +pub fn turn_intent_status(status: crate::projects::types::WorkItemRunStatus) -> &'static str { + use crate::projects::types::WorkItemRunStatus; + match status { + WorkItemRunStatus::Queued + | WorkItemRunStatus::Deferred + | WorkItemRunStatus::Dispatching => "queued", + WorkItemRunStatus::Running | WorkItemRunStatus::Waiting => "running", + WorkItemRunStatus::Succeeded => "completed", + WorkItemRunStatus::Failed => "failed", + WorkItemRunStatus::Cancelled => "cancelled", + } +} diff --git a/src-tauri/crates/project-management/src/work_run_service/read.rs b/src-tauri/crates/project-management/src/work_run_service/read.rs index 246c00c0a2..f8405bb7fb 100644 --- a/src-tauri/crates/project-management/src/work_run_service/read.rs +++ b/src-tauri/crates/project-management/src/work_run_service/read.rs @@ -1,7 +1,7 @@ use rusqlite::{params, OptionalExtension, Transaction}; use crate::projects::io::helpers::conn; -use crate::projects::types::WorkItemRun; +use crate::projects::types::{WorkItemRun, WorkItemRunTarget}; use super::store::{canonical_standalone_org_id, db, require_run, scope_key}; use super::{error, MAX_RUN_ATTEMPTS}; @@ -11,6 +11,55 @@ pub fn read(run_id: &str) -> Result { require_run(&connection, run_id) } +/// Recover the effective WorkItemRun identity for one Project composer turn. +pub fn find_project_session_turn( + session_id: &str, + origin_turn_intent_id: &str, +) -> Result, String> { + if session_id.trim().is_empty() || origin_turn_intent_id.trim().is_empty() { + return Err(format!( + "{}:session_id and origin_turn_intent_id are required", + error::INVALID_REQUEST + )); + } + let connection = conn()?; + let idempotency_key = format!("project-session-turn:{session_id}:{origin_turn_intent_id}"); + let ids = { + let mut statement = db(connection.prepare( + "SELECT id FROM pm_work_item_runs + WHERE idempotency_key = ?1 + ORDER BY created_at DESC, id DESC", + ))?; + let rows = db(statement.query_map([&idempotency_key], |row| row.get::<_, String>(0)))?; + db(rows.collect::>>())? + }; + let mut matched = None; + for run_id in ids { + let run = require_run(&connection, &run_id)?; + let target_matches = matches!( + &run.target_snapshot.target, + WorkItemRunTarget::ResumeSession { session_id: target_session_id } + if target_session_id == session_id + ); + let origin_matches = run + .input + .get("originTurnIntentId") + .and_then(serde_json::Value::as_str) + == Some(origin_turn_intent_id); + if !target_matches || !origin_matches { + continue; + } + if matched.is_some() { + return Err(format!( + "{}:multiple Work Item Runs map session {session_id} intent {origin_turn_intent_id}", + error::IDEMPOTENCY_CONFLICT + )); + } + matched = Some(run); + } + Ok(matched) +} + pub(crate) fn read_in_transaction( tx: &Transaction<'_>, run_id: &str, diff --git a/src-tauri/crates/project-management/src/work_run_service/tests.rs b/src-tauri/crates/project-management/src/work_run_service/tests.rs index eb35c09235..2e010cc693 100644 --- a/src-tauri/crates/project-management/src/work_run_service/tests.rs +++ b/src-tauri/crates/project-management/src/work_run_service/tests.rs @@ -209,11 +209,13 @@ fn enqueue_is_atomic_and_idempotent() { let _sandbox = test_env::sandbox(); seed(); - let first = enqueue(request("manual:1")).expect("enqueue"); - let replay = enqueue(request("manual:1")).expect("idempotent replay"); - assert_eq!(first.id, replay.id); - assert_eq!(first.status, WorkItemRunStatus::Queued); - assert_eq!(first.target_snapshot.work_item_revision, 0); + let first = enqueue_with_receipt(request("manual:1")).expect("enqueue"); + let replay = enqueue_with_receipt(request("manual:1")).expect("idempotent replay"); + assert!(!first.duplicate); + assert!(replay.duplicate); + assert_eq!(first.run.id, replay.run.id); + assert_eq!(first.run.status, WorkItemRunStatus::Queued); + assert_eq!(first.run.target_snapshot.work_item_revision, 0); let connection = conn().expect("connection"); let run_count: i64 = connection @@ -230,6 +232,52 @@ fn enqueue_is_atomic_and_idempotent() { assert_eq!(dispatch_count, 1); } +#[test] +fn work_run_status_projects_to_the_turn_receipt_exhaustively() { + for status in [ + WorkItemRunStatus::Queued, + WorkItemRunStatus::Deferred, + WorkItemRunStatus::Dispatching, + ] { + assert_eq!(turn_intent_status(status), "queued"); + } + for status in [WorkItemRunStatus::Running, WorkItemRunStatus::Waiting] { + assert_eq!(turn_intent_status(status), "running"); + } + assert_eq!( + turn_intent_status(WorkItemRunStatus::Succeeded), + "completed" + ); + assert_eq!(turn_intent_status(WorkItemRunStatus::Failed), "failed"); + assert_eq!( + turn_intent_status(WorkItemRunStatus::Cancelled), + "cancelled" + ); +} + +#[test] +fn project_session_turn_recovers_its_effective_run_identity() { + let _sandbox = test_env::sandbox(); + seed(); + let mut request = request("project-session-turn:session-1:intent-x"); + request.target_snapshot = WorkItemRunTargetSnapshot::new(WorkItemRunTarget::ResumeSession { + session_id: "session-1".to_string(), + }); + request.input = serde_json::json!({ + "content": "continue", + "originTurnIntentId": "intent-x", + }); + let receipt = enqueue_with_receipt(request).expect("enqueue Project turn"); + + let recovered = find_project_session_turn("session-1", "intent-x") + .expect("lookup Project turn") + .expect("mapped WorkItemRun"); + assert_eq!(recovered.id, receipt.run.id); + assert!(find_project_session_turn("session-1", "intent-other") + .expect("lookup missing intent") + .is_none()); +} + #[test] fn idempotency_key_rejects_different_request() { let _sandbox = test_env::sandbox(); diff --git a/src-tauri/crates/session-persistence/src/agent_core_bridge.rs b/src-tauri/crates/session-persistence/src/agent_core_bridge.rs index 3c296f918d..c8e4292951 100644 --- a/src-tauri/crates/session-persistence/src/agent_core_bridge.rs +++ b/src-tauri/crates/session-persistence/src/agent_core_bridge.rs @@ -154,6 +154,30 @@ fn upsert_turn_intent_adapter( } } +fn claim_turn_intent_adapter( + session_id: &str, + turn_intent_id: &str, + client_message_id: Option<&str>, + org_run_id: Option<&str>, + source: session_bridge::TurnIntentBridgeSource, + status: session_bridge::TurnIntentBridgeStatus, +) -> Result { + let claim = turn_intents::claim_initial( + session_id, + turn_intent_id, + client_message_id, + org_run_id, + map_bridge_source(source), + map_bridge_status(status), + ) + .map_err(|err| err.to_string())?; + Ok(session_bridge::TurnIntentBridgeClaim { + duplicate: claim.duplicate, + status: map_persisted_status(claim.row.status), + client_message_id: claim.row.client_message_id, + }) +} + fn update_turn_intent_status_adapter( session_id: &str, turn_intent_id: &str, @@ -209,6 +233,7 @@ pub fn register() { session_bridge::register_record_token_usage(record_token_usage_adapter); session_bridge::register_record_usage_telemetry_batch(record_usage_telemetry_batch_adapter); session_bridge::register_upsert_turn_intent(upsert_turn_intent_adapter); + session_bridge::register_claim_turn_intent(claim_turn_intent_adapter); session_bridge::register_update_turn_intent_status(update_turn_intent_status_adapter); session_bridge::register_get_turn_intent_status(get_turn_intent_status_adapter); session_bridge::register_mark_pending_turn_intents_stale( diff --git a/src-tauri/crates/session-persistence/src/turn_intents.rs b/src-tauri/crates/session-persistence/src/turn_intents.rs index 1f8aeeb6f2..6e1be69176 100644 --- a/src-tauri/crates/session-persistence/src/turn_intents.rs +++ b/src-tauri/crates/session-persistence/src/turn_intents.rs @@ -31,7 +31,7 @@ use rusqlite::{params, Connection, OptionalExtension, Result as SqliteResult}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use super::connection::get_connection; +use super::connection::{begin_immediate, get_connection, with_sessions_writer}; // ============================================ // Source / status enums (wire-stable strings) @@ -172,7 +172,7 @@ impl TurnIntentStatus { // Domain types // ============================================ -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct TurnIntentRow { pub session_id: String, @@ -185,6 +185,18 @@ pub struct TurnIntentRow { pub updated_at: String, } +/// Result of atomically reserving one logical turn identity. +/// +/// `duplicate = false` means this caller inserted the canonical lifecycle +/// row and may proceed with externally visible turn side effects. A duplicate +/// returns the original row unchanged, including its original +/// `client_message_id`, so a retry cannot replace the first caller's identity. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TurnIntentClaim { + pub row: TurnIntentRow, + pub duplicate: bool, +} + #[derive(Debug, thiserror::Error)] pub enum IntentError { #[error("sqlite error: {0}")] @@ -241,6 +253,60 @@ fn transition_allowed(from: TurnIntentStatus, to: TurnIntentStatus) -> bool { // CRUD // ============================================ +/// Atomically reserve a turn intent before any runtime/transcript/scheduler +/// mutation. +/// +/// A read-then-upsert sequence is not a claim: two concurrent IPC calls can +/// both observe absence and both execute side effects before one loses the +/// eventual `INSERT OR IGNORE`. This transaction serializes the decision and +/// returns the existing durable owner to every loser. +pub fn claim_initial( + session_id: &str, + turn_intent_id: &str, + client_message_id: Option<&str>, + org_run_id: Option<&str>, + source: TurnIntentSource, + status: TurnIntentStatus, +) -> Result { + with_sessions_writer(|| { + let conn = get_connection()?; + let tx = begin_immediate(&conn)?; + if get_intent(&tx, session_id, turn_intent_id)?.is_some() { + // Reuse the canonical upsert's NULL-only org-run backfill and + // mismatch fence. A duplicate identity from another Agent Org + // run must fail, not inherit the first run's execution. + let row = upsert_initial_on( + &tx, + session_id, + turn_intent_id, + client_message_id, + org_run_id, + source, + status, + )?; + tx.commit()?; + return Ok(TurnIntentClaim { + row, + duplicate: true, + }); + } + let row = upsert_initial_on( + &tx, + session_id, + turn_intent_id, + client_message_id, + org_run_id, + source, + status, + )?; + tx.commit()?; + Ok(TurnIntentClaim { + row, + duplicate: false, + }) + }) +} + fn row_from_sql(row: &rusqlite::Row<'_>) -> SqliteResult { let source_str: String = row.get(4)?; let status_str: String = row.get(5)?; @@ -821,6 +887,44 @@ mod tests { }); } + #[test] + fn concurrent_exact_claim_has_one_durable_owner() { + with_temp_orgii_home(|| { + let barrier = std::sync::Arc::new(std::sync::Barrier::new(3)); + let handles = ["client-a", "client-b"].map(|client_message_id| { + let barrier = std::sync::Arc::clone(&barrier); + std::thread::spawn(move || { + barrier.wait(); + claim_initial( + "test-session-concurrent-claim", + "intent-shared", + Some(client_message_id), + None, + TurnIntentSource::UserSubmit, + TurnIntentStatus::Queued, + ) + .expect("atomic claim") + }) + }); + barrier.wait(); + let claims = handles.map(|handle| handle.join().expect("claim thread")); + + assert_eq!(claims.iter().filter(|claim| !claim.duplicate).count(), 1); + assert_eq!(claims.iter().filter(|claim| claim.duplicate).count(), 1); + assert_eq!(claims[0].row, claims[1].row); + assert!(matches!( + claims[0].row.client_message_id.as_deref(), + Some("client-a" | "client-b") + )); + assert_eq!( + list_for_session("test-session-concurrent-claim") + .expect("list claimed lifecycle") + .len(), + 1 + ); + }); + } + #[test] fn restart_reconciliation_closes_every_in_flight_intent() { with_temp_orgii_home(|| { diff --git a/src-tauri/src/agent_sessions/cli/commands/run.rs b/src-tauri/src/agent_sessions/cli/commands/run.rs index fc30b7e84c..7b1a118fe8 100644 --- a/src-tauri/src/agent_sessions/cli/commands/run.rs +++ b/src-tauri/src/agent_sessions/cli/commands/run.rs @@ -7,13 +7,16 @@ use super::super::session_runner; use super::super::types::{KeySource, SessionStatus}; use agent_core::session::IdeContext; use serde::{Deserialize, Serialize}; +use session_persistence::turn_intents::TurnIntentStatus; #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase")] pub struct CliRunReceipt { pub session_id: String, pub turn_intent_id: String, - pub status: SessionStatus, + pub effective_turn_intent_id: String, + pub status: String, + pub duplicate: bool, } /// Start one CLI turn on an existing session row. @@ -108,6 +111,41 @@ pub async fn cli_agent_tui_release(session_id: String) -> Result { .map_err(|e| format!("Task error: {}", e))? } +async fn claim_cli_turn( + session_id: &str, + turn: &TurnIdentity, +) -> Result { + let sid = session_id.to_string(); + let intent_id = turn.turn_intent_id.clone(); + let message_id = turn.client_message_id.clone(); + tokio::task::spawn_blocking(move || persistence::claim_cli_turn(&sid, &intent_id, &message_id)) + .await + .map_err(|err| format!("Task error: {err}"))? +} + +async fn reject_claimed_cli_turn(session_id: &str, turn_intent_id: &str) { + let sid = session_id.to_string(); + let intent_id = turn_intent_id.to_string(); + let reject_result = + tokio::task::spawn_blocking(move || persistence::reject_claimed_cli_turn(&sid, &intent_id)) + .await; + match reject_result { + Ok(Ok(())) => {} + Ok(Err(error)) => tracing::warn!( + session_id, + turn_intent_id, + %error, + "failed to reject claimed CLI turn after preparation error" + ), + Err(error) => tracing::warn!( + session_id, + turn_intent_id, + %error, + "CLI turn rejection worker failed" + ), + } +} + /// Run a code session (spawn CLI agent in background). #[tauri::command] pub async fn cli_agent_run(mut request: CliRunRequest) -> Result<(), String> { @@ -117,7 +155,30 @@ pub async fn cli_agent_run(mut request: CliRunRequest) -> Result<(), String> { ); let control_lock = session_runner::session_control_lock(&request.session_id).await; let _control_guard = control_lock.lock().await; - run_turn(request, turn).await + if route_project_turn_if_needed( + &request.session_id, + &request.user_input, + &turn.turn_intent_id, + &turn.client_message_id, + ) + .await? + .is_some() + { + return Ok(()); + } + let claim = claim_cli_turn(&request.session_id, &turn).await?; + if claim.duplicate { + return Ok(()); + } + let session_id = request.session_id.clone(); + let turn_intent_id = turn.turn_intent_id.clone(); + match run_turn(request, turn).await { + Ok(_) => Ok(()), + Err(error) => { + reject_claimed_cli_turn(&session_id, &turn_intent_id).await; + Err(error) + } + } } /// Create the root Work Item on the first non-empty Project-mode turn. @@ -179,7 +240,7 @@ async fn enqueue_project_turn_if_needed( user_input: &str, turn_intent_id: &str, client_message_id: &str, -) -> Result, String> { +) -> Result, String> { if user_input.trim().is_empty() || turn_intent_id.starts_with("wir_") { return Ok(None); } @@ -228,21 +289,34 @@ async fn enqueue_project_turn_if_needed( "content": user_input, "displayText": user_input, "clientMessageId": client_message_id, + "originTurnIntentId": turn_intent_id, }), idempotency_key: format!("project-session-turn:{session_id}:{turn_intent_id}"), max_attempts: 3, parent_run_id: None, }; - tokio::task::spawn_blocking(move || project_management::work_run_service::enqueue(request)) - .await - .map_err(|err| format!("Project CLI WorkItemRun enqueue worker failed: {err}"))? - .map(Some) + tokio::task::spawn_blocking(move || { + project_management::work_run_service::enqueue_with_receipt(request) + }) + .await + .map_err(|err| format!("Project CLI WorkItemRun enqueue worker failed: {err}"))? + .map(Some) +} + +async fn route_project_turn_if_needed( + session_id: &str, + user_input: &str, + turn_intent_id: &str, + client_message_id: &str, +) -> Result, String> { + bootstrap_project_root_if_needed(session_id, user_input).await?; + enqueue_project_turn_if_needed(session_id, user_input, turn_intent_id, client_message_id).await } /// Shared turn body behind both `cli_agent_run` and `cli_agent_message`: /// persist acceptance under the registry lock, broadcast `running`, then spawn /// the background runner. -async fn run_turn(request: CliRunRequest, turn: TurnIdentity) -> Result<(), String> { +async fn run_turn(request: CliRunRequest, turn: TurnIdentity) -> Result { let CliRunRequest { session_id, user_input, @@ -286,23 +360,6 @@ async fn run_turn(request: CliRunRequest, turn: TurnIdentity) -> Result<(), Stri .map_err(|err| format!("Task error: {}", err))??; } - bootstrap_project_root_if_needed(&session_id, &user_input).await?; - if let Some(run) = enqueue_project_turn_if_needed( - &session_id, - &user_input, - &turn_intent_id, - &client_message_id, - ) - .await? - { - tracing::info!( - session_id = %session_id, - run_id = %run.id, - "queued CLI Project turn through durable WorkItem dispatcher" - ); - return Ok(()); - } - // 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; @@ -437,19 +494,13 @@ async fn run_turn(request: CliRunRequest, turn: TurnIdentity) -> Result<(), Stri drop(sessions); tracing::info!(session_id = %session_id, "cli_agent_run: background runner registered"); - Ok(()) + Ok(TurnIntentStatus::Running) } -/// Send a follow-up message to a running or completed session. -/// -/// Kills any existing running agent (OS process + proxy), re-allocates a fresh -/// proxy token (the previous one was released on completion), loads the CLI -/// session ID for resume, then re-runs with the new input. -/// -/// 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] -pub async fn cli_agent_message(request: CliMessageRequest) -> Result { +async fn execute_claimed_cli_message( + request: CliMessageRequest, + turn: TurnIdentity, +) -> Result { let CliMessageRequest { session_id, content, @@ -458,10 +509,9 @@ pub async fn cli_agent_message(request: CliMessageRequest) -> Result Result Result Result Result { + let turn = TurnIdentity::from_client( + request.turn_intent_id.take(), + request.client_message_id.take(), + ); + let receipt_turn_intent_id = turn.turn_intent_id.clone(); + let session_id = request.session_id.clone(); + + // Serialize the durable claim with kill + spawn. An exact response-loss + // retry returns here before config mutation or process termination. + let control_lock = session_runner::session_control_lock(&session_id).await; + let _control_guard = control_lock.lock().await; + if let Some(project_receipt) = route_project_turn_if_needed( + &session_id, + &request.content, + &turn.turn_intent_id, + &turn.client_message_id, + ) + .await? + { + let status = + project_management::work_run_service::turn_intent_status(project_receipt.run.status); + return Ok(CliRunReceipt { + session_id, + turn_intent_id: receipt_turn_intent_id, + effective_turn_intent_id: project_receipt.run.id, + status: status.to_string(), + duplicate: project_receipt.duplicate, + }); + } + + let claim = claim_cli_turn(&session_id, &turn).await?; + if claim.duplicate { + return Ok(CliRunReceipt { + session_id, + turn_intent_id: receipt_turn_intent_id, + effective_turn_intent_id: turn.turn_intent_id, + status: claim.status.as_str().to_string(), + duplicate: true, + }); + } + + let execution_result = execute_claimed_cli_message(request, turn).await; + + let accepted_status = match execution_result { + Ok(status) => status, + Err(error) => { + reject_claimed_cli_turn(&session_id, &receipt_turn_intent_id).await; + return Err(error); + } + }; + Ok(CliRunReceipt { session_id, - turn_intent_id, - status: SessionStatus::Running, + effective_turn_intent_id: receipt_turn_intent_id.clone(), + turn_intent_id: receipt_turn_intent_id, + status: accepted_status.as_str().to_string(), + duplicate: false, }) } @@ -681,3 +789,202 @@ pub async fn cli_agent_approval_response( ) .await } + +#[cfg(test)] +mod tests { + use super::*; + use crate::agent_sessions::cli::persistence::CreateCodeSessionParams; + use crate::test_utils::test_env; + + fn create_test_session(session_id: &str) { + persistence::create_session( + session_id, + &CreateCodeSessionParams { + name: Some("CLI exact intent test".to_string()), + flow: None, + runner: None, + cli_agent_type: "claude_code".to_string(), + model: Some("original-model".to_string()), + tier: None, + account_id: Some("account-a".to_string()), + repo_path: Some("/tmp".to_string()), + branch: None, + worktree_path: None, + worktree_base_ref: None, + proxy_token: None, + proxy_url: None, + hosted_token: None, + proxy_session_id: None, + isolate: None, + background: Some(false), + key_source: Some("own_key".to_string()), + additional_directories: None, + parent_session_id: None, + org_member_id: None, + org_id: None, + project_id: None, + project_name: None, + project_slug: None, + work_item_id: None, + agent_role: None, + product_mode: None, + }, + ) + .expect("create CLI test session"); + } + + fn seed_running_intent(session_id: &str, turn_intent_id: &str) { + let claim = persistence::claim_cli_turn(session_id, turn_intent_id, "message-original") + .expect("claim original CLI intent"); + assert!(!claim.duplicate); + persistence::accept_cli_turn(session_id, turn_intent_id, "message-original") + .expect("accept original CLI intent"); + } + + #[test] + fn cli_receipt_keeps_origin_and_effective_project_intents_distinct() { + let value = serde_json::to_value(CliRunReceipt { + session_id: "cliagent-project".to_string(), + turn_intent_id: "intent-composer-x".to_string(), + effective_turn_intent_id: "wir_effective-y".to_string(), + status: "queued".to_string(), + duplicate: false, + }) + .expect("serialize CLI receipt"); + + assert_eq!(value["turnIntentId"], "intent-composer-x"); + assert_eq!(value["effectiveTurnIntentId"], "wir_effective-y"); + assert_eq!(value["status"], "queued"); + } + + #[tokio::test] + async fn initial_exact_replay_does_not_mutate_or_spawn() { + let _sandbox = test_env::sandbox(); + let session_id = "cli-exact-replay-initial"; + let turn_intent_id = "intent-initial"; + create_test_session(session_id); + seed_running_intent(session_id, turn_intent_id); + + cli_agent_run(CliRunRequest { + session_id: session_id.to_string(), + user_input: "response-loss retry".to_string(), + mode: Some("plan".to_string()), + turn_intent_id: Some(turn_intent_id.to_string()), + client_message_id: Some("message-retry".to_string()), + ..Default::default() + }) + .await + .expect("initial replay receipt"); + + let session = persistence::get_session(session_id) + .expect("load CLI session") + .expect("CLI session exists"); + assert_eq!(session.agent_exec_mode.as_deref(), Some("build")); + assert!( + !session_runner::RUNNING_SESSIONS + .lock() + .await + .contains_key(session_id), + "initial duplicate spawned a second runner" + ); + } + + #[tokio::test] + async fn running_exact_replay_does_not_mutate_kill_or_spawn() { + let _sandbox = test_env::sandbox(); + let session_id = "cli-exact-replay-running"; + let turn_intent_id = "intent-running"; + create_test_session(session_id); + seed_running_intent(session_id, turn_intent_id); + + let sentinel = tokio::spawn(std::future::pending::<()>()); + assert!(session_runner::RUNNING_SESSIONS + .lock() + .await + .insert(session_id.to_string(), sentinel) + .is_none()); + + let receipt = cli_agent_message(CliMessageRequest { + session_id: session_id.to_string(), + content: "response-loss retry".to_string(), + model: Some("must-not-be-written".to_string()), + account_id: Some("account-b".to_string()), + mode: Some("plan".to_string()), + turn_intent_id: Some(turn_intent_id.to_string()), + client_message_id: Some("message-retry".to_string()), + ..Default::default() + }) + .await + .expect("running replay receipt"); + + let handle = session_runner::RUNNING_SESSIONS + .lock() + .await + .remove(session_id) + .expect("original runner remains registered"); + let original_runner_survived = !handle.is_finished(); + handle.abort(); + let session = persistence::get_session(session_id) + .expect("load CLI session") + .expect("CLI session exists"); + + assert!(receipt.duplicate); + assert_eq!(receipt.status, "running"); + assert_eq!(receipt.turn_intent_id, turn_intent_id); + assert_eq!(receipt.effective_turn_intent_id, turn_intent_id); + assert!( + original_runner_survived, + "duplicate retry killed the runner" + ); + assert_eq!(session.model.as_deref(), Some("original-model")); + assert_eq!(session.account_id.as_deref(), Some("account-a")); + assert_eq!(session.agent_exec_mode.as_deref(), Some("build")); + } + + #[tokio::test] + async fn completed_exact_replay_does_not_mutate_or_spawn() { + let _sandbox = test_env::sandbox(); + let session_id = "cli-exact-replay-completed"; + let turn_intent_id = "intent-completed"; + create_test_session(session_id); + seed_running_intent(session_id, turn_intent_id); + persistence::update_cli_turn_lifecycle( + session_id, + SessionStatus::Completed, + None, + Some((turn_intent_id, TurnIntentStatus::Completed)), + ) + .expect("complete original CLI intent"); + + let receipt = cli_agent_message(CliMessageRequest { + session_id: session_id.to_string(), + content: "late response-loss retry".to_string(), + model: Some("must-not-be-written".to_string()), + account_id: Some("account-b".to_string()), + mode: Some("plan".to_string()), + turn_intent_id: Some(turn_intent_id.to_string()), + client_message_id: Some("message-retry".to_string()), + ..Default::default() + }) + .await + .expect("completed replay receipt"); + + let session = persistence::get_session(session_id) + .expect("load CLI session") + .expect("CLI session exists"); + assert!(receipt.duplicate); + assert_eq!(receipt.status, "completed"); + assert_eq!(receipt.turn_intent_id, turn_intent_id); + assert_eq!(receipt.effective_turn_intent_id, turn_intent_id); + assert_eq!(session.model.as_deref(), Some("original-model")); + assert_eq!(session.account_id.as_deref(), Some("account-a")); + assert_eq!(session.agent_exec_mode.as_deref(), Some("build")); + assert!( + !session_runner::RUNNING_SESSIONS + .lock() + .await + .contains_key(session_id), + "completed duplicate spawned a second runner" + ); + } +} 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 4a3136abd1..240b1e48aa 100644 --- a/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs +++ b/src-tauri/src/agent_sessions/cli/persistence/session_crud.rs @@ -2,7 +2,7 @@ use chrono::Utc; use rusqlite::{params, Connection, OptionalExtension, Result as SqliteResult}; use agent_core::session::AgentExecMode; -use database::db::get_connection; +use database::db::{begin_immediate, get_connection, with_sessions_writer}; use super::super::types::{ session_defaults, KeySource, SessionRunner, SessionStatus, DEFAULT_CODE_SESSION_FLOW, @@ -337,6 +337,95 @@ pub fn accept_cli_turn( ) } +/// Result of reserving one exact CLI follow-up identity. +/// +/// The command holds the per-session control lock while calling this helper, +/// and this transaction makes the durable row visible before any destructive +/// follow-up side effect (account/model mutation, killing the old process, or +/// spawning its replacement). A response-loss retry therefore observes the +/// original lifecycle instead of executing the turn twice. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct CliTurnClaim { + pub duplicate: bool, + pub status: session_persistence::turn_intents::TurnIntentStatus, +} + +pub fn claim_cli_turn( + session_id: &str, + turn_intent_id: &str, + client_message_id: &str, +) -> Result { + with_sessions_writer(|| { + let conn = get_connection().map_err(|err| err.to_string())?; + let tx = begin_immediate(&conn).map_err(|err| err.to_string())?; + let session_exists = tx + .query_row( + "SELECT 1 FROM code_sessions WHERE session_id = ?1", + [session_id], + |_| Ok(()), + ) + .optional() + .map_err(|err| err.to_string())? + .is_some(); + if !session_exists { + return Err(format!("session not found: {session_id}")); + } + + if let Some(existing) = + session_persistence::turn_intents::get_intent(&tx, session_id, turn_intent_id) + .map_err(|err| err.to_string())? + { + tx.commit().map_err(|err| err.to_string())?; + return Ok(CliTurnClaim { + duplicate: true, + status: existing.status, + }); + } + + let claimed = session_persistence::turn_intents::upsert_initial_on( + &tx, + session_id, + turn_intent_id, + Some(client_message_id), + None, + session_persistence::turn_intents::TurnIntentSource::UserSubmit, + session_persistence::turn_intents::TurnIntentStatus::Queued, + ) + .map_err(|err| err.to_string())?; + tx.commit().map_err(|err| err.to_string())?; + Ok(CliTurnClaim { + duplicate: false, + status: claimed.status, + }) + }) +} + +/// Close a newly claimed intent when preparation fails before `run_turn` +/// accepts it. The guarded canonical transition makes this a no-op for an +/// intent that has already advanced beyond `queued`. +pub fn reject_claimed_cli_turn(session_id: &str, turn_intent_id: &str) -> Result<(), String> { + with_sessions_writer(|| { + let conn = get_connection().map_err(|err| err.to_string())?; + let tx = begin_immediate(&conn).map_err(|err| err.to_string())?; + let existing = + session_persistence::turn_intents::get_intent(&tx, session_id, turn_intent_id) + .map_err(|err| err.to_string())?; + if existing.is_some_and(|row| { + row.status == session_persistence::turn_intents::TurnIntentStatus::Queued + }) { + session_persistence::turn_intents::update_status_on( + &tx, + session_id, + turn_intent_id, + session_persistence::turn_intents::TurnIntentStatus::Rejected, + ) + .map_err(|err| err.to_string())?; + } + tx.commit().map_err(|err| err.to_string())?; + Ok(()) + }) +} + /// `accept_cli_turn` for a resumed session: same atomic acceptance, but the /// intent is sourced as `Resume` and has no client message behind it — resume /// replays the session's stored `user_input` instead of a fresh submit. diff --git a/src-tauri/src/commands/handler_list.inc b/src-tauri/src/commands/handler_list.inc index f7206a98e1..d20b22472b 100644 --- a/src-tauri/src/commands/handler_list.inc +++ b/src-tauri/src/commands/handler_list.inc @@ -1222,6 +1222,7 @@ agent_core::flow_awareness::commands::flow_clear_session, agent_core::state::commands::agent_session_list, agent_core::state::commands::agent_session_cancel, agent_core::state::commands::agent_session_info, +agent_core::state::commands::agent_turn_intent_status, agent_core::state::commands::agent_session_manual_compact, agent_core::state::commands::housekeeper_context_compaction_status, agent_core::state::commands::housekeeper_context_compaction_set_enabled, diff --git a/src/api/tauri/rpc/schemas/__tests__/cliRunReceipt.test.ts b/src/api/tauri/rpc/schemas/__tests__/cliRunReceipt.test.ts new file mode 100644 index 0000000000..83782de4aa --- /dev/null +++ b/src/api/tauri/rpc/schemas/__tests__/cliRunReceipt.test.ts @@ -0,0 +1,35 @@ +import { describe, expect, it } from "vitest"; + +import { CliRunReceiptSchema } from "../cli"; + +describe("CliRunReceiptSchema", () => { + it("preserves the origin and effective turn identities", () => { + expect( + CliRunReceiptSchema.parse({ + sessionId: "cliagent-1", + turnIntentId: "intent-x", + effectiveTurnIntentId: "wir-y", + status: "queued", + duplicate: true, + }) + ).toEqual({ + sessionId: "cliagent-1", + turnIntentId: "intent-x", + effectiveTurnIntentId: "wir-y", + status: "queued", + duplicate: true, + }); + }); + + it("rejects an unrecognized lifecycle status", () => { + expect(() => + CliRunReceiptSchema.parse({ + sessionId: "cliagent-1", + turnIntentId: "intent-x", + effectiveTurnIntentId: "intent-x", + status: "maybe-running", + duplicate: false, + }) + ).toThrow(); + }); +}); diff --git a/src/api/tauri/rpc/schemas/cli.ts b/src/api/tauri/rpc/schemas/cli.ts index 44866bd7fa..9d118c27bb 100644 --- a/src/api/tauri/rpc/schemas/cli.ts +++ b/src/api/tauri/rpc/schemas/cli.ts @@ -23,7 +23,19 @@ export const CliMessageInputSchema = z.object({ export const CliRunReceiptSchema = z.object({ sessionId: z.string(), turnIntentId: z.string(), - status: z.string(), + effectiveTurnIntentId: z.string().min(1), + status: z.enum([ + "optimistic", + "queued", + "running", + "completed", + "failed", + "cancelled", + "stale", + "coalesced", + "rejected", + ]), + duplicate: z.boolean(), }); export const CliSessionIdInputSchema = z.object({