diff --git a/crates/ironclaw_engine/src/executor/orchestrator.rs b/crates/ironclaw_engine/src/executor/orchestrator.rs index 7416f4a1aa9..3fcaaeea215 100644 --- a/crates/ironclaw_engine/src/executor/orchestrator.rs +++ b/crates/ironclaw_engine/src/executor/orchestrator.rs @@ -830,7 +830,7 @@ async fn handle_execute_code_step( "had_error": result.failure.is_some(), "pending_gate": result.need_approval.as_ref().map(|na| { match na { - ThreadOutcome::GatePaused { gate_name, action_name, call_id, parameters, resume_kind, resume_output } => serde_json::json!({ + ThreadOutcome::GatePaused { gate_name, action_name, call_id, parameters, resume_kind, resume_output, paused_lease } => serde_json::json!({ "gate_paused": true, "gate_name": gate_name, "action_name": action_name, @@ -838,6 +838,7 @@ async fn handle_execute_code_step( "parameters": parameters, "resume_kind": serde_json::to_value(resume_kind).unwrap_or_default(), "resume_output": resume_output, + "paused_lease": paused_lease, }), _ => serde_json::Value::Null, } @@ -1117,6 +1118,7 @@ async fn handle_execute_action( parameters, resume_kind, resume_output, + paused_lease, }) => { let _ = leases.refund_use(lease.id).await; let output = serde_json::json!({"status": "gate_paused", "gate_name": gate_name}); @@ -1147,6 +1149,7 @@ async fn handle_execute_action( "parameters": parameters, "resume_kind": serde_json::to_value(&*resume_kind).unwrap_or_default(), "resume_output": resume_output, + "paused_lease": paused_lease.as_deref().cloned(), }); ExtFunctionResult::Return(json_to_monty(&result)) } @@ -1624,6 +1627,7 @@ async fn execute_single_action( parameters, resume_kind, resume_output, + paused_lease, }) => { let output = serde_json::json!({"status": "gate_paused", "gate_name": &gate_name}); let event = EventKind::ApprovalRequested { @@ -1646,6 +1650,7 @@ async fn execute_single_action( "parameters": parameters, "resume_kind": serde_json::to_value(&*resume_kind).unwrap_or_default(), "resume_output": resume_output, + "paused_lease": paused_lease.as_deref().cloned(), }); (result_json, event, output) } @@ -2541,6 +2546,10 @@ fn parse_outcome(result: &serde_json::Value) -> ThreadOutcome { .unwrap_or(serde_json::json!({})), resume_kind, resume_output: result.get("resume_output").cloned(), + paused_lease: result + .get("paused_lease") + .cloned() + .and_then(|value| serde_json::from_value(value).ok()), } } _ => ThreadOutcome::Completed { response: None }, @@ -3112,6 +3121,7 @@ mod tests { parameters: serde_json::json!({"cmd":"ls"}), resume_kind: crate::gate::ResumeKind::Approval { allow_always: true }, resume_output: None, + paused_lease: None, }; normalize_pause_outcome(&mut thread, &outcome).unwrap(); assert_eq!(thread.state, ThreadState::Waiting); @@ -3133,18 +3143,38 @@ mod tests { #[test] fn parse_outcome_gate_paused() { + let lease = crate::types::capability::CapabilityLease { + id: crate::types::capability::LeaseId::new(), + thread_id: crate::types::thread::ThreadId::new(), + capability_name: "test-capability".into(), + granted_actions: crate::types::capability::GrantedActions::Specific(vec![ + "shell".into(), + ]), + granted_at: chrono::Utc::now(), + expires_at: None, + max_uses: Some(1), + uses_remaining: Some(1), + revoked: false, + revoked_reason: None, + }; let result = serde_json::json!({ "outcome": "gate_paused", "gate_name": "approval", "action_name": "shell", "call_id": "abc", "parameters": {"cmd": "rm -rf /"}, - "resume_kind": {"Approval": {"allow_always": true}} + "resume_kind": {"Approval": {"allow_always": true}}, + "paused_lease": lease, }); let outcome = parse_outcome(&result); - assert!( - matches!(outcome, ThreadOutcome::GatePaused { action_name, .. } if action_name == "shell") - ); + assert!(matches!( + outcome, + ThreadOutcome::GatePaused { + action_name, + paused_lease: Some(_), + .. + } if action_name == "shell" + )); } #[test] diff --git a/crates/ironclaw_engine/src/executor/scripting.rs b/crates/ironclaw_engine/src/executor/scripting.rs index 39f88f9d804..3be3060696c 100644 --- a/crates/ironclaw_engine/src/executor/scripting.rs +++ b/crates/ironclaw_engine/src/executor/scripting.rs @@ -1186,6 +1186,7 @@ async fn preflight_action( parameters: params.clone(), resume_kind: crate::gate::ResumeKind::Approval { allow_always: true }, resume_output: None, + paused_lease: None, }, ); } diff --git a/crates/ironclaw_engine/src/executor/structured.rs b/crates/ironclaw_engine/src/executor/structured.rs index 6211438420e..e9791cc1ba7 100644 --- a/crates/ironclaw_engine/src/executor/structured.rs +++ b/crates/ironclaw_engine/src/executor/structured.rs @@ -166,6 +166,7 @@ pub async fn execute_action_calls( parameters: call.parameters.clone(), resume_kind: crate::gate::ResumeKind::Approval { allow_always: true }, resume_output: None, + paused_lease: None, }), }); } @@ -319,6 +320,11 @@ pub async fn execute_action_calls( allow_always: false, }), resume_output: result.output.get("resume_output").cloned(), + paused_lease: result + .output + .get("paused_lease") + .cloned() + .and_then(|value| serde_json::from_value(value).ok()), }); } results.push(result); @@ -387,6 +393,7 @@ fn classify_exec_result( parameters, resume_kind, resume_output, + paused_lease, }) => { let _error_msg = format!("gate paused: {gate_name}"); let error_result = ActionResult { @@ -397,6 +404,7 @@ fn classify_exec_result( "gate": gate_name, "resume_kind": serde_json::to_value(&*resume_kind).unwrap_or_default(), "resume_output": resume_output.as_deref().cloned(), + "paused_lease": paused_lease.as_deref().cloned(), }), is_error: true, duration: std::time::Duration::ZERO, @@ -893,6 +901,7 @@ mod tests { auth_url: None, }), resume_output: None, + paused_lease: None, })], )); let leases = Arc::new(LeaseManager::new()); @@ -976,6 +985,7 @@ mod tests { auth_url: None, }), resume_output: None, + paused_lease: None, }), Ok(ActionResult { call_id: String::new(), diff --git a/crates/ironclaw_engine/src/runtime/messaging.rs b/crates/ironclaw_engine/src/runtime/messaging.rs index baea5d8ae9d..57150df6b59 100644 --- a/crates/ironclaw_engine/src/runtime/messaging.rs +++ b/crates/ironclaw_engine/src/runtime/messaging.rs @@ -42,6 +42,11 @@ pub enum ThreadOutcome { /// Completed action output that should be injected on resume instead /// of re-running the action. resume_output: Option, + /// Lease snapshot captured when the gate paused the action. + /// Boxed to keep `ThreadOutcome::GatePaused` under clippy's + /// `large_enum_variant` threshold — `CapabilityLease` is ~360 + /// bytes and would otherwise dominate the whole enum's size. + paused_lease: Option>, }, } diff --git a/crates/ironclaw_engine/src/runtime/mission.rs b/crates/ironclaw_engine/src/runtime/mission.rs index 9610c58be64..6de8820ecf3 100644 --- a/crates/ironclaw_engine/src/runtime/mission.rs +++ b/crates/ironclaw_engine/src/runtime/mission.rs @@ -4151,6 +4151,7 @@ mod tests { allow_always: false, }), resume_output: None, + paused_lease: None, }); } Ok(ActionResult { diff --git a/crates/ironclaw_engine/src/types/error.rs b/crates/ironclaw_engine/src/types/error.rs index 57cb46872a1..7e6f459af2d 100644 --- a/crates/ironclaw_engine/src/types/error.rs +++ b/crates/ironclaw_engine/src/types/error.rs @@ -79,6 +79,7 @@ pub enum EngineError { parameters: Box, resume_kind: Box, resume_output: Option>, + paused_lease: Option>, }, } diff --git a/src/agent/agent_loop.rs b/src/agent/agent_loop.rs index 034920b6412..6944fa179e0 100644 --- a/src/agent/agent_loop.rs +++ b/src/agent/agent_loop.rs @@ -1595,6 +1595,7 @@ impl Agent { submission, Submission::ExecApproval { .. } | Submission::ApprovalResponse { .. } + | Submission::ExternalCallback { .. } | Submission::GateAuthResolution { .. } ) { message diff --git a/src/bridge/effect_adapter.rs b/src/bridge/effect_adapter.rs index 4b5d9304788..d7f9fea48dd 100644 --- a/src/bridge/effect_adapter.rs +++ b/src/bridge/effect_adapter.rs @@ -169,6 +169,7 @@ impl EffectBridgeAdapter { parameters: serde_json::Value, resume_kind: ironclaw_engine::ResumeKind, resume_output: Option, + paused_lease: Option, ) -> EngineError { EngineError::GatePaused { gate_name: gate_name.to_string(), @@ -177,6 +178,7 @@ impl EffectBridgeAdapter { parameters: Box::new(parameters), resume_kind: Box::new(resume_kind), resume_output: resume_output.map(Box::new), + paused_lease: paused_lease.map(Box::new), } } @@ -185,6 +187,7 @@ impl EffectBridgeAdapter { parameters: serde_json::Value, context: &ThreadExecutionContext, output_value: &serde_json::Value, + lease: &CapabilityLease, ) -> Option { let status = output_value.get("status").and_then(|v| v.as_str())?; let name = output_value.get("name").and_then(|v| v.as_str())?; @@ -221,6 +224,7 @@ impl EffectBridgeAdapter { ), }, None, + Some(lease.clone()), )), _ => None, } @@ -726,7 +730,7 @@ impl EffectBridgeAdapter { &self, action_name: &str, parameters: serde_json::Value, - _lease: &CapabilityLease, + lease: &CapabilityLease, context: &ThreadExecutionContext, approval_already_granted: bool, ) -> Result { @@ -826,6 +830,7 @@ impl EffectBridgeAdapter { auth_url: sanitize_auth_url(auth_url.as_deref()), }, None, + Some(lease.clone()), )); } Ok(crate::bridge::auth_manager::LatentActionExecution::NeedsSetup { message }) => { @@ -898,6 +903,7 @@ impl EffectBridgeAdapter { auth_url: sanitize_auth_url(cred.auth_url.as_deref()), }, None, + Some(lease.clone()), )); } AuthCheckResult::Ready => { @@ -939,6 +945,7 @@ impl EffectBridgeAdapter { auth_url: sanitize_auth_url(auth_url.as_deref()), }, None, + Some(lease.clone()), )); } ToolReadiness::NeedsSetup { message } => { @@ -996,6 +1003,7 @@ impl EffectBridgeAdapter { allow_always: false, }, None, + Some(lease.clone()), )); } } @@ -1011,6 +1019,7 @@ impl EffectBridgeAdapter { parameters, ironclaw_engine::ResumeKind::Approval { allow_always: true }, None, + Some(lease.clone()), )); } } @@ -1163,6 +1172,7 @@ impl EffectBridgeAdapter { parameters.clone(), context, &output_value, + lease, ) { return Err(err); @@ -1204,6 +1214,7 @@ impl EffectBridgeAdapter { auth_url: None, }, None, + Some(lease.clone()), )); } diff --git a/src/bridge/router.rs b/src/bridge/router.rs index e157fd7fb24..05d49006338 100644 --- a/src/bridge/router.rs +++ b/src/bridge/router.rs @@ -485,6 +485,7 @@ async fn requeue_auth_pending_gate( expires_at: chrono::Utc::now() + chrono::Duration::minutes(30), original_message: pending.original_message.clone(), resume_output: pending.resume_output.clone(), + paused_lease: pending.paused_lease.clone(), approval_already_granted: pending.approval_already_granted, }; @@ -512,6 +513,7 @@ fn pairing_pending_gate_from_auth(pending: &PendingGate, extension_name: &str) - expires_at: chrono::Utc::now() + chrono::Duration::minutes(30), original_message: pending.original_message.clone(), resume_output: pending.resume_output.clone(), + paused_lease: pending.paused_lease.clone(), approval_already_granted: pending.approval_already_granted, } } @@ -663,6 +665,47 @@ async fn revert_always_allow( } } +/// Validate that a `paused_lease` snapshot recorded when a gate paused +/// still represents a usable lease at resume time. +/// +/// A gate can sit in the pending-gate store for hours or across process +/// restarts; during that window the original lease might have been +/// revoked, expired, or the pending record could have drifted off its +/// original thread. Callers that fail this check must NOT use the +/// snapshot — fall through to `LeaseManager::find_lease_for_action` +/// (which enforces its own scoping) or fail closed. +fn snapshot_lease_still_valid( + lease: &ironclaw_engine::CapabilityLease, + pending: &PendingGate, +) -> bool { + lease.thread_id == pending.thread_id + && lease.granted_actions.covers(&pending.action_name) + && !lease.revoked + && lease + .expires_at + .map(|exp| exp > chrono::Utc::now()) + .unwrap_or(true) +} + +/// Pick the lease to use for resuming a pending gate action. Prefers the +/// `paused_lease` snapshot the gate recorded if it's still valid; falls +/// back to a live lookup in the `LeaseManager`. Returns `None` if neither +/// path yields a lease — the caller maps that to a "no active lease" +/// error. +async fn resume_lease_for_pending_gate( + pending: &PendingGate, + leases: &ironclaw_engine::LeaseManager, +) -> Option { + if let Some(snapshot) = pending.paused_lease.clone() + && snapshot_lease_still_valid(&snapshot, pending) + { + return Some(snapshot); + } + leases + .find_lease_for_action(pending.thread_id, &pending.action_name) + .await +} + async fn execute_pending_gate_action( agent: &Agent, state: &EngineState, @@ -679,10 +722,7 @@ async fn execute_pending_gate_action( .ok_or_else(|| engine_err("load thread", "thread not found"))?; let resolved_call_id = resolved_or_synthetic_call_id_for_pending_action(state, pending).await?; - let lease = state - .thread_manager - .leases - .find_lease_for_action(pending.thread_id, &pending.action_name) + let lease = resume_lease_for_pending_gate(pending, &state.thread_manager.leases) .await .ok_or_else(|| { engine_err( @@ -750,6 +790,7 @@ async fn execute_pending_gate_action( parameters, resume_kind, resume_output, + paused_lease, }) => { let display_parameters = state .effect_adapter @@ -786,6 +827,7 @@ async fn execute_pending_gate_action( .clone() .or_else(|| Some(message.content.clone())), resume_output: resume_output.map(|value| *value), + paused_lease: paused_lease.map(|lease| *lease), approval_already_granted: approval_already_granted || matches!( pending.resume_kind, @@ -3312,6 +3354,7 @@ async fn await_thread_outcome( expires_at: chrono::Utc::now() + chrono::Duration::minutes(30), original_message: Some(message.content.clone()), resume_output: None, + paused_lease: None, approval_already_granted: false, }; let pending_request_id = pending.request_id.to_string(); @@ -3363,6 +3406,7 @@ async fn await_thread_outcome( parameters, resume_kind, resume_output, + paused_lease, } => { use crate::gate::pending::PendingGate; @@ -3397,6 +3441,10 @@ async fn await_thread_outcome( expires_at: chrono::Utc::now() + chrono::Duration::minutes(30), original_message: Some(message.content.clone()), resume_output, + // Unbox: `ThreadOutcome::GatePaused.paused_lease` is + // `Option>` to keep the outcome + // enum compact; `PendingGate` stores it unboxed. + paused_lease: paused_lease.map(|b| *b), approval_already_granted: false, }; @@ -5343,6 +5391,7 @@ mod tests { expires_at: chrono::Utc::now() + chrono::Duration::minutes(30), original_message: None, resume_output: None, + paused_lease: None, approval_already_granted: false, } } @@ -7712,4 +7761,202 @@ mod tests { completed={completed_idx} call={call_idx} gate_paused={gate_paused_idx}" ); } + + // ── resume_lease_for_pending_gate tests ──────────────────── + // + // Pins the PR #2631 review ask: when resuming a paused gate, we + // prefer the `paused_lease` snapshot the gate recorded. The snapshot + // MUST be validated (thread_id match, action coverage, not revoked, + // not expired) before use — a stale snapshot falling through silently + // would bypass revocation semantics that `find_lease_for_action` + // normally enforces. These tests drive the helper end-to-end to pin + // every branch of the decision. + + fn sample_lease_for_pending(pending: &PendingGate) -> ironclaw_engine::CapabilityLease { + ironclaw_engine::CapabilityLease { + id: ironclaw_engine::types::capability::LeaseId::new(), + thread_id: pending.thread_id, + capability_name: "tools".into(), + granted_actions: ironclaw_engine::GrantedActions::Specific(vec![ + pending.action_name.clone(), + ]), + granted_at: chrono::Utc::now(), + expires_at: None, + max_uses: None, + uses_remaining: None, + revoked: false, + revoked_reason: None, + } + } + + #[tokio::test] + async fn resume_lease_prefers_snapshot_even_when_lease_manager_empty() { + // Reproduces the original bug: the LeaseManager has no active + // lease for the paused action (lease evicted or never persisted + // through restart), but the pending gate carries a snapshot. The + // resume must use the snapshot. + let thread_id = ironclaw_engine::ThreadId::new(); + let mut pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let snapshot = sample_lease_for_pending(&pending); + let snapshot_id = snapshot.id; + pending.paused_lease = Some(snapshot); + + let leases = ironclaw_engine::LeaseManager::new(); + // Intentionally empty — no lease for this thread/action. + + let lease = resume_lease_for_pending_gate(&pending, &leases) + .await + .expect("snapshot must be used when LeaseManager has nothing"); + assert_eq!(lease.id, snapshot_id, "should return the snapshot lease"); + } + + #[tokio::test] + async fn resume_lease_rejects_revoked_snapshot_and_falls_back() { + // A revoked snapshot must NOT resume the action. Fall back to + // the LeaseManager; if that has a valid lease, use it. + let thread_id = ironclaw_engine::ThreadId::new(); + let mut pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let mut revoked_snapshot = sample_lease_for_pending(&pending); + revoked_snapshot.revoked = true; + revoked_snapshot.revoked_reason = Some("user revoked mid-pause".into()); + pending.paused_lease = Some(revoked_snapshot); + + let leases = ironclaw_engine::LeaseManager::new(); + let live_lease = leases + .grant( + thread_id, + "tools", + ironclaw_engine::GrantedActions::Specific(vec!["shell".into()]), + None, + None, + ) + .await + .expect("grant live lease"); + + let lease = resume_lease_for_pending_gate(&pending, &leases) + .await + .expect("fallback lease must be found"); + assert_eq!( + lease.id, live_lease.id, + "revoked snapshot must be skipped and LeaseManager fallback used" + ); + } + + #[tokio::test] + async fn resume_lease_rejects_expired_snapshot_and_falls_back() { + // An expired snapshot must not be accepted even if all other + // fields check out. + let thread_id = ironclaw_engine::ThreadId::new(); + let mut pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let mut expired_snapshot = sample_lease_for_pending(&pending); + expired_snapshot.expires_at = Some(chrono::Utc::now() - chrono::Duration::minutes(1)); + pending.paused_lease = Some(expired_snapshot); + + let leases = ironclaw_engine::LeaseManager::new(); + let live_lease = leases + .grant( + thread_id, + "tools", + ironclaw_engine::GrantedActions::Specific(vec!["shell".into()]), + None, + None, + ) + .await + .expect("grant live lease"); + + let lease = resume_lease_for_pending_gate(&pending, &leases) + .await + .expect("fallback lease must be found"); + assert_eq!(lease.id, live_lease.id, "expired snapshot must be skipped"); + } + + #[tokio::test] + async fn resume_lease_rejects_snapshot_with_wrong_thread_id() { + // Defensive: a snapshot whose thread_id doesn't match the pending + // gate's thread_id must never be trusted, even if other fields + // look valid. Guards against pending-gate-store drift or future + // refactors that reuse a snapshot across threads. + let thread_id = ironclaw_engine::ThreadId::new(); + let mut pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let mut mismatched_snapshot = sample_lease_for_pending(&pending); + mismatched_snapshot.thread_id = ironclaw_engine::ThreadId::new(); // different thread + pending.paused_lease = Some(mismatched_snapshot); + + let leases = ironclaw_engine::LeaseManager::new(); + // Empty — no fallback. + let result = resume_lease_for_pending_gate(&pending, &leases).await; + assert!( + result.is_none(), + "mismatched-thread snapshot must be skipped and fallback must fail cleanly" + ); + } + + #[tokio::test] + async fn resume_lease_rejects_snapshot_missing_action_coverage() { + // A snapshot whose granted_actions does not cover the pending + // action name must not be used, even when untargeted. + let thread_id = ironclaw_engine::ThreadId::new(); + let mut pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let mut mismatched_snapshot = sample_lease_for_pending(&pending); + mismatched_snapshot.granted_actions = + ironclaw_engine::GrantedActions::Specific(vec!["unrelated_tool".into()]); + pending.paused_lease = Some(mismatched_snapshot); + + let leases = ironclaw_engine::LeaseManager::new(); + let result = resume_lease_for_pending_gate(&pending, &leases).await; + assert!( + result.is_none(), + "snapshot must be skipped when it doesn't grant the pending action" + ); + } + + #[tokio::test] + async fn resume_lease_returns_none_when_no_snapshot_and_no_active_lease() { + // Sanity: if there's nothing in either path, the helper returns + // None so the caller can map to a clean "no active lease" error. + let thread_id = ironclaw_engine::ThreadId::new(); + let pending = sample_pending_gate( + "alice", + thread_id, + ironclaw_engine::ResumeKind::Approval { + allow_always: false, + }, + ); + let leases = ironclaw_engine::LeaseManager::new(); + assert!( + resume_lease_for_pending_gate(&pending, &leases) + .await + .is_none() + ); + } } diff --git a/src/channels/web/features/oauth/mod.rs b/src/channels/web/features/oauth/mod.rs index 01f2a814545..55b9368c4d6 100644 --- a/src/channels/web/features/oauth/mod.rs +++ b/src/channels/web/features/oauth/mod.rs @@ -31,7 +31,9 @@ use axum::{ use sha2::{Digest, Sha256}; use crate::channels::relay::DEFAULT_RELAY_NAME; -use crate::channels::web::platform::legacy_auth::clear_auth_mode; +use crate::channels::web::platform::legacy_auth::{ + clear_auth_mode, clear_session_auth_mode_for_thread, +}; use crate::channels::web::platform::state::{GatewayState, rate_limit_key_from_headers}; use crate::channels::web::types::AppEvent; use crate::channels::web::util::web_incoming_message; @@ -311,9 +313,17 @@ pub(crate) async fn oauth_callback_handler( } } - // Clear auth mode regardless of outcome so the next user message goes - // through to the LLM instead of being intercepted as a token. - clear_auth_mode(&state, &flow.user_id).await; + // Clear legacy session auth mode regardless of outcome so the next + // user message goes through to the LLM instead of being intercepted + // as a token. + // + // Do NOT clear the engine pending-auth gate here: the successful + // callback path still needs the pending gate so it can resolve and + // replay the paused action (preserving the paused_lease), and failed + // callbacks should leave the gate visible for retry from the UI. + // The gate is cleared by the engine itself when `ExternalCallback` + // resolves (success) or when the user explicitly cancels (failure). + let _ = clear_session_auth_mode_for_thread(&state, &flow.user_id, None).await; // After successful OAuth, auto-activate the extension so it moves // from "Installed (Authenticate)" → "Active" without a second click. diff --git a/src/channels/web/platform/legacy_auth.rs b/src/channels/web/platform/legacy_auth.rs index 3803098c01c..5a6533702ef 100644 --- a/src/channels/web/platform/legacy_auth.rs +++ b/src/channels/web/platform/legacy_auth.rs @@ -21,15 +21,35 @@ use crate::channels::web::types::{ ActionResponse, AppEvent, AuthCancelRequest, AuthTokenRequest, OnboardingStateDto, }; -/// Clear pending auth mode on the active thread for a user. +/// Clear pending auth mode on the active thread for a user (both legacy +/// v1 session-scoped and engine v2 pending-gate). pub async fn clear_auth_mode(state: &GatewayState, user_id: &str) { let _ = clear_auth_mode_for_thread(state, user_id, None).await; } +/// Clear both the legacy v1 session `pending_auth` and the engine v2 +/// pending auth gate. Use this after the user has resolved the credential +/// flow (or explicitly cancelled) and the gate is done with. pub(crate) async fn clear_auth_mode_for_thread( state: &GatewayState, user_id: &str, thread_id: Option<&str>, +) -> Result<(), (StatusCode, String)> { + clear_session_auth_mode_for_thread(state, user_id, thread_id).await?; + crate::bridge::clear_engine_pending_auth(user_id, thread_id).await; + Ok(()) +} + +/// Clear ONLY the legacy v1 session-level `pending_auth` state, leaving +/// the engine v2 pending auth gate intact. Used by the OAuth callback: +/// the successful-callback path still needs the engine gate present so +/// the `ExternalCallback` replay can resolve it (and preserve the +/// paused_lease), and the failed-callback path should leave the gate +/// visible so the user can retry from the UI. +pub(crate) async fn clear_session_auth_mode_for_thread( + state: &GatewayState, + user_id: &str, + thread_id: Option<&str>, ) -> Result<(), (StatusCode, String)> { if let Some(ref sm) = state.session_manager { let session = sm.get_or_create_session(user_id).await; @@ -49,7 +69,6 @@ pub(crate) async fn clear_auth_mode_for_thread( thread.pending_auth = None; } } - crate::bridge::clear_engine_pending_auth(user_id, thread_id).await; Ok(()) } diff --git a/src/gate/pending.rs b/src/gate/pending.rs index 3b3c60bd97c..5ab49608aa6 100644 --- a/src/gate/pending.rs +++ b/src/gate/pending.rs @@ -1,7 +1,7 @@ //! Pending gate state — unified type replacing `PendingApproval` and `PendingAuth`. use chrono::{DateTime, Utc}; -use ironclaw_engine::{ResumeKind, ThreadId}; +use ironclaw_engine::{CapabilityLease, ResumeKind, ThreadId}; use serde::{Deserialize, Serialize}; use uuid::Uuid; @@ -66,6 +66,10 @@ pub struct PendingGate { /// Completed action output to inject on resume after auth finishes. #[serde(default, skip_serializing_if = "Option::is_none")] pub resume_output: Option, + /// Lease snapshot to reuse when resuming a paused action whose original + /// lease use was already consumed before the gate fired. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub paused_lease: Option, /// Whether approval has already been granted for this paused action. #[serde(default)] pub approval_already_granted: bool, @@ -142,6 +146,7 @@ mod tests { expires_at: Utc::now() + Duration::seconds(expires_in_secs), original_message: None, resume_output: None, + paused_lease: None, approval_already_granted: false, } } diff --git a/src/gate/persistence.rs b/src/gate/persistence.rs index 6b413db3af8..c6161fbd562 100644 --- a/src/gate/persistence.rs +++ b/src/gate/persistence.rs @@ -160,6 +160,7 @@ mod tests { expires_at: Utc::now() + Duration::minutes(30), original_message: None, resume_output: None, + paused_lease: None, approval_already_granted: false, } } diff --git a/src/gate/store.rs b/src/gate/store.rs index f0aedcf9713..1548266971d 100644 --- a/src/gate/store.rs +++ b/src/gate/store.rs @@ -326,6 +326,7 @@ mod tests { expires_at: Utc::now() + Duration::seconds(expires_in_secs), original_message: None, resume_output: None, + paused_lease: None, approval_already_granted: false, } } diff --git a/tests/engine_v2_gate_integration.rs b/tests/engine_v2_gate_integration.rs index c3fa2c8b284..bc8bbe00edc 100644 --- a/tests/engine_v2_gate_integration.rs +++ b/tests/engine_v2_gate_integration.rs @@ -219,6 +219,7 @@ impl EffectExecutor for InstallThenAliasEffects { instructions: "Authenticate GitHub".into(), auth_url: None, }), + paused_lease: None, resume_output: None, }); } @@ -283,6 +284,7 @@ impl EffectExecutor for GateMockEffects { call_id: "call_gate_1".into(), parameters: Box::new(parameters), resume_kind: Box::new(ResumeKind::Approval { allow_always: true }), + paused_lease: None, resume_output: None, }); } @@ -298,6 +300,7 @@ impl EffectExecutor for GateMockEffects { instructions: "Authenticate your Notion workspace".into(), auth_url: None, }), + paused_lease: None, resume_output: None, }); } @@ -311,6 +314,7 @@ impl EffectExecutor for GateMockEffects { call_id: "call_gate_1".into(), parameters: Box::new(parameters), resume_kind: Box::new(ResumeKind::Approval { allow_always: true }), + paused_lease: None, resume_output: None, }); } @@ -327,6 +331,7 @@ impl EffectExecutor for GateMockEffects { instructions: "Provide your API key".into(), auth_url: None, }), + paused_lease: None, resume_output: None, }); } @@ -660,6 +665,7 @@ fn sample_pending_gate( created_at: Utc::now(), expires_at: Utc::now() + chrono::Duration::minutes(30), original_message: None, + paused_lease: None, resume_output: None, approval_already_granted: false, }