diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 1004b19d072..617829674bc 100644 --- a/crates/ironclaw_reborn_composition/src/slack_delivery.rs +++ b/crates/ironclaw_reborn_composition/src/slack_delivery.rs @@ -4,6 +4,7 @@ //! observer runs after the workflow accepts an inbound Slack message, waits for //! the submitted run to finish, reads the finalized assistant reply, and sends it //! through the host-mediated product outbound delivery seam. +// arch-exempt: large_file, deferred-busy hint logic stays here to share observer/test fixtures with final-reply delivery; decomposition tracked in #4818. use std::collections::hash_map::DefaultHasher; use std::hash::{Hash, Hasher}; @@ -45,6 +46,7 @@ use ironclaw_turns::{ }; use ironclaw_wasm_product_adapters::ImmediateAckWorkflowObserver; use serde::{Deserialize, Serialize}; +use std::collections::{HashSet, VecDeque}; use tokio::sync::Semaphore; use crate::AuthChallengeProvider; @@ -61,6 +63,13 @@ const SLACK_DELIVERY_TIMEOUT_MESSAGE: &str = "This is taking longer than expected — check the WebUI for the result."; const SLACK_DELIVERY_ERROR_MESSAGE: &str = "Something went wrong delivering the result here. Check the WebUI."; +/// Posted when the blocking run is `BlockedApproval` and no gate_ref is available. +const SLACK_DEFERRED_BUSY_APPROVAL_MESSAGE: &str = "Ironclaw is waiting on a pending approval before taking new messages — reply `approve` or `deny` (or `approve gate:`) to resume."; +/// Posted for any other non-terminal blocking state, or when the state lookup fails. +/// +/// Honest copy: no queue yet — deferred-message drain ships in a separate PR (#4812). +/// Revisit this wording once #4812 lands and automatic retry is available. +const SLACK_DEFERRED_BUSY_GENERIC_MESSAGE: &str = "Ironclaw is still working on a previous message and can't take this one yet — please resend it once the current task finishes."; #[derive(Debug, Clone, PartialEq, Eq)] struct BlockedActionableMarker { @@ -113,10 +122,23 @@ pub struct SlackFinalReplyDeliveryServices { pub auth_challenges: Option>, } +/// Maximum number of (conversation, run_id) pairs remembered for hint dedup. +/// FIFO eviction beyond this cap keeps memory O(1); a false-negative after +/// eviction just means one extra hint, which is harmless. +const HINT_SEEN_CAP: usize = 256; + +type HintSeenKey = (String, TurnRunId); +type HintSeenSet = Mutex<(VecDeque, HashSet)>; + pub struct SlackFinalReplyDeliveryObserver { services: SlackFinalReplyDeliveryServices, settings: SlackFinalReplyDeliverySettings, delivery_permits: Arc, + /// Per-observer throttle: at most one deferred-busy hint per + /// (conversation fingerprint, active_run_id) pair. + /// Caps Slack API usage per blocked run; bounded FIFO eviction keeps memory O(1); + /// a false-negative after eviction just means one extra hint, harmless. + hint_seen: HintSeenSet, } impl SlackFinalReplyDeliveryObserver { @@ -132,6 +154,7 @@ impl SlackFinalReplyDeliveryObserver { services, settings, delivery_permits: Arc::new(Semaphore::new(settings.max_concurrent_deliveries.get())), + hint_seen: Mutex::new((VecDeque::new(), HashSet::new())), } } @@ -946,6 +969,91 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { { return; } + // A2b: DeferredBusy feedback — the user's message was silently dropped + // because a run is blocked on a pending gate. Post a one-shot state-aware + // hint so the user knows to approve/deny/wait rather than being left in + // silence. Same best-effort semantics as A2: post failure → debug! only. + // + // Authorization: only post if the binding lookup succeeds, matching the + // same guard used by `post_rejection_hint_if_authorized`. + // + // Inline await is safe: the protocol ACK already returned before this + // observer runs, and the runner's admission permit in runner_immediate_ack.rs + // bounds the lifetime of this entire post-ACK task. A detached spawn would + // escape `drain_immediate_ack_tasks` shutdown/drain without adding any + // backpressure benefit. + if let Some(active_run_id) = deferred_busy_user_message_run_id(&envelope, &ack) { + // Throttle: at most one hint per (conversation, active_run_id) pair. + // Check before the coordinator call to avoid a round-trip on repeats. + let conv_key = envelope + .external_conversation_ref() + .conversation_fingerprint(); + let throttle_key = (conv_key, active_run_id); + let already_seen = { + let mut guard = self.hint_seen.lock().unwrap_or_else(|e| e.into_inner()); + let (queue, set) = &mut *guard; + if set.contains(&throttle_key) { + true + } else { + // FIFO eviction to keep the set bounded at O(1) memory. + if set.len() >= HINT_SEEN_CAP + && let Some(oldest) = queue.pop_front() + { + set.remove(&oldest); + } + set.insert(throttle_key.clone()); + queue.push_back(throttle_key); + false + } + }; + if already_seen { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + "deferred-busy hint suppressed: already posted for this (conversation, run_id) pair" + ); + return; + } + // Authorization: verify the conversation has a valid binding before + // posting. Skip silently (debug!) when unauthorized. + let binding = match self + .services + .binding_service + .lookup_binding(ResolveBindingRequest::from_envelope(&envelope)) + .await + { + Ok(b) => b, + Err(error) => { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + error = %error, + "skipped deferred-busy hint because the originating conversation was not authorized" + ); + return; + } + }; + // Derive the scope for the active run state lookup. + // Falls back to generic copy on any failure — never skips the hint. + let hint = deferred_busy_hint_from_run_state( + self.services.turn_coordinator.as_ref(), + &binding, + active_run_id, + ) + .await; + if let Err(post_err) = post_slack_message( + self.services.egress.as_ref(), + envelope.external_conversation_ref(), + &hint, + ) + .await + { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + error = %post_err, + "failed to post deferred-busy hint to Slack (best-effort)" + ); + } + return; + } let Ok(_permit) = self.delivery_permits.clone().acquire_owned().await else { tracing::warn!( target = "ironclaw::reborn::slack_delivery", @@ -1151,6 +1259,100 @@ fn rejection_hint_for_resolution( Some(hint) } +/// Returns `Some(active_run_id)` when the ack + payload combination should trigger +/// the deferred-busy hint flow: a plain `DeferredBusy` ack on a `UserMessage` payload. +/// +/// Returns `None` for `Duplicate` acks (the idempotency ledger never settles +/// `DeferredBusy`, so Slack transport retries re-process as a fresh plain +/// `DeferredBusy`, never `Duplicate{DeferredBusy}`; suppressing all `Duplicate` +/// matches the invariant in `rejection_hint_for_resolution`) and for all non-user- +/// message payloads (resolution/control payloads must stay silent). +fn deferred_busy_user_message_run_id( + envelope: &ProductInboundEnvelope, + ack: &ProductInboundAck, +) -> Option { + // All Duplicate acks are suppressed — same pattern as rejection_hint_for_resolution. + if matches!(ack, ProductInboundAck::Duplicate { .. }) { + return None; + } + // Only reply to user messages — resolution/control/noop payloads must stay silent. + if !matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) { + return None; + } + match ack { + ProductInboundAck::DeferredBusy { active_run_id, .. } => Some(*active_run_id), + _ => None, + } +} + +/// Looks up the blocking run's state and returns the appropriate deferred-busy +/// hint copy. +/// +/// - `BlockedApproval` with `Some(gate_ref)` → approval wording with concrete `approve {ref}` command +/// - `BlockedApproval` with `None` gate_ref → approval wording without a specific gate command +/// - `BlockedAuth` with `Some(gate_ref)` → auth wording with concrete `auth deny {ref}` command +/// - `BlockedAuth` with `None` gate_ref → auth wording without the deny command +/// - anything else / lookup failure → generic wording +/// +/// Never returns an error — lookup failures degrade to the generic copy. +async fn deferred_busy_hint_from_run_state( + coordinator: &dyn TurnCoordinator, + binding: &ResolvedBinding, + active_run_id: TurnRunId, +) -> String { + let scope = match (|| -> Result { + let thread_scope = thread_scope_from_binding(binding)?; + turn_scope_from_thread_scope(binding, &thread_scope) + })() { + Ok(s) => s, + Err(err) => { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + error = %err, + "deferred-busy scope derivation failed; using generic copy" + ); + return SLACK_DEFERRED_BUSY_GENERIC_MESSAGE.to_string(); + } + }; + match coordinator + .get_run_state(GetRunStateRequest { + scope, + run_id: active_run_id, + }) + .await + { + Ok(state) => match state.status { + TurnStatus::BlockedApproval => match state.gate_ref.as_ref() { + Some(gate_ref) => format!( + "Ironclaw is waiting on a pending approval before taking new messages \ + — reply `approve {ref}` or `deny` to resume.", + ref = gate_ref.as_str() + ), + None => SLACK_DEFERRED_BUSY_APPROVAL_MESSAGE.to_string(), + }, + TurnStatus::BlockedAuth => match state.gate_ref.as_ref() { + Some(gate_ref) => format!( + "Ironclaw is waiting on an authentication step before taking new messages \ + — complete the authentication prompt (or reply `auth deny {ref}` to decline).", + ref = gate_ref.as_str() + ), + None => "Ironclaw is waiting on an authentication step before taking new messages \ + — complete the authentication prompt to continue." + .to_string(), + }, + _ => SLACK_DEFERRED_BUSY_GENERIC_MESSAGE.to_string(), + }, + Err(err) => { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + error = %err, + "deferred-busy run-state lookup failed; using generic copy" + ); + SLACK_DEFERRED_BUSY_GENERIC_MESSAGE.to_string() + } + } +} + fn rejection_ack_for_workflow_error(error: &ProductAdapterError) -> Option { match error { ProductAdapterError::WorkflowRejected { @@ -2090,7 +2292,8 @@ mod tests { WriteCommunicationPreferenceRequest, }; use ironclaw_product_adapters::{ - FakeOutboundDeliverySink, FakeProtocolHttpEgress, ProductAdapterId, + EgressResponse, FakeOutboundDeliverySink, FakeProtocolHttpEgress, ProductAdapterId, + ProtocolHttpEgressError, }; use ironclaw_slack_v2_adapter::{ SlackV2Adapter, SlackV2AdapterConfig, slack_request_signature_auth_requirement, @@ -3151,41 +3354,70 @@ mod tests { ); } - /// Duplicate { prior: Accepted } → nothing posted (already succeeded). + // ── DeferredBusy ack feedback tests ─────────────────────────────────────── + + fn deferred_busy_ack() -> ProductInboundAck { + ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:deferred").expect("ref"), + active_run_id: TurnRunId::new(), + } + } + + /// DeferredBusy ack + UserMessage payload + BlockedApproval state with gate_ref → + /// exactly one Slack post containing the concrete `approve gate:` command. + /// + /// The hint post is awaited inline; no yield needed before inspecting the egress capture. #[tokio::test] - async fn duplicate_accepted_ack_posts_nothing() { + async fn deferred_busy_ack_with_user_message_posts_hint() { let install = "test-install"; - let run_id = TurnRunId::new(); let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "2000.1"), + )), + ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - // The coordinator is a no-op because Duplicate{Accepted} has no submitted_run_id. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, - )); + // BlockedApproval with concrete gate_ref → hint embeds the actionable command. + let gate_ref_str = "gate:approval-abc123"; + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + TurnStatus::BlockedApproval, + Some(gate_ref_str), + )])); let observer = make_observer(coordinator, egress.clone(), outbound, install); - let env = envelope(scoped_approval_resolution_payload()); - let prior = ProductInboundAck::Accepted { - accepted_message_ref: AcceptedMessageRef::new("slack:prior").expect("ref"), - submitted_run_id: run_id, - }; - let ack = ProductInboundAck::Duplicate { - prior: Box::new(prior), - }; + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage for DeferredBusy + UserMessage" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); assert!( - !calls.iter().any(|c| c.path == "/api/chat.postMessage"), - "no chat.postMessage expected for Duplicate{{Accepted}}" + body.contains("waiting on a pending approval"), + "deferred-busy hint must mention 'waiting on a pending approval', got: {body}" + ); + assert!( + body.contains(gate_ref_str), + "deferred-busy approval hint must embed the concrete gate ref '{gate_ref_str}', got: {body}" ); } - /// Duplicate { prior: Rejected } → nothing posted at observer level. + /// DeferredBusy ack + non-UserMessage payload → no post (resolution payloads + /// already have their own feedback path and must stay silent here). #[tokio::test] - async fn duplicate_rejected_ack_posts_nothing() { + async fn deferred_busy_ack_with_resolution_payload_posts_nothing() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); @@ -3196,59 +3428,26 @@ mod tests { )); let observer = make_observer(coordinator, egress.clone(), outbound, install); let env = envelope(scoped_approval_resolution_payload()); - let ack = ProductInboundAck::Duplicate { - prior: Box::new(rejected_ack( - ironclaw_product_adapters::ProductRejectionKind::BindingRequired, - )), - }; + let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; let calls = egress.calls(); assert!( !calls.iter().any(|c| c.path == "/api/chat.postMessage"), - "no chat.postMessage expected for Duplicate{{Rejected}}" + "no chat.postMessage expected for DeferredBusy with non-UserMessage payload" ); } - /// All `Duplicate` acks suppress the hint, regardless of the prior inside. + /// Duplicate { prior: DeferredBusy } + UserMessage → nothing posted. /// - /// `Duplicate` is keyed on the external event id: transport retries of the - /// same Slack event land as `Duplicate{original}`. The original processing - /// already posted any hint, so replays must not repeat the side effect. - /// A user re-typing "approve" produces a new event id → a fresh `Rejected`, - /// never a `Duplicate`, so suppressing `Duplicate{Rejected}` loses nothing. - #[test] - fn duplicate_acks_produce_no_hint() { - let env = envelope(scoped_approval_resolution_payload()); - - let duplicate_rejected = ProductInboundAck::Duplicate { - prior: Box::new(rejected_ack( - ironclaw_product_adapters::ProductRejectionKind::BindingRequired, - )), - }; - assert!( - rejection_hint_for_resolution(&env, &duplicate_rejected).is_none(), - "Duplicate{{Rejected}} must NOT produce a hint (transport replay)" - ); - - let duplicate_accepted = ProductInboundAck::Duplicate { - prior: Box::new(ProductInboundAck::Accepted { - accepted_message_ref: ironclaw_turns::AcceptedMessageRef::new("slack:prior") - .expect("ref"), - submitted_run_id: ironclaw_turns::TurnRunId::new(), - }), - }; - assert!( - rejection_hint_for_resolution(&env, &duplicate_accepted).is_none(), - "Duplicate{{Accepted}} must not produce a hint" - ); - } - - /// A failed best-effort rejection-hint post must not fall through into generic - /// delivery-error feedback. + /// `should_settle_ack` returns false for DeferredBusy, so the idempotency + /// ledger never settles it. Slack transport retries re-process as a fresh + /// plain DeferredBusy (never Duplicate{DeferredBusy}), making this case + /// unreachable in practice. We still enforce None-for-all-Duplicate for + /// safety, matching the invariant in `rejection_hint_for_resolution`. #[tokio::test] - async fn rejected_resolution_hint_post_failure_is_best_effort() { + async fn duplicate_deferred_busy_with_user_message_posts_nothing() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); @@ -3258,36 +3457,30 @@ mod tests { TurnStatus::Running, )); let observer = make_observer(coordinator, egress.clone(), outbound, install); - let env = envelope(scoped_approval_resolution_payload()); - let ack = rejected_ack(ironclaw_product_adapters::ProductRejectionKind::BindingRequired); + let env = envelope(user_message_payload()); + let ack = ProductInboundAck::Duplicate { + prior: Box::new(deferred_busy_ack()), + }; observer.observe_workflow_ack(env, ack).await; let calls = egress.calls(); - let post_calls: Vec<_> = calls - .iter() - .filter(|c| c.path == "/api/chat.postMessage") - .collect(); - assert_eq!( - post_calls.len(), - 1, - "expected exactly one failed best-effort hint post" + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "Duplicate{{DeferredBusy}} must NOT post a hint (None-for-all-Duplicate invariant)" ); } - /// Delivery error (RunWaitTimedOut) → timeout notice posted to conversation. + /// Two distinct plain DeferredBusy + UserMessage envelopes with different active_run_ids + /// → two posts (throttle is per (conversation, run_id) pair). /// - /// Uses `FakeConversationBindingService` so the binding lookup succeeds and - /// delivery enters the polling loop, which then times out because the - /// coordinator always returns `Running`. + /// Each `deferred_busy_ack()` call creates a fresh `TurnRunId`, so the two acks have + /// distinct throttle keys and each posts exactly one hint. #[tokio::test] - async fn delivery_timeout_posts_timeout_notice_to_conversation() { - use ironclaw_product_workflow::FakeConversationBindingService; - + async fn two_distinct_deferred_busy_user_messages_post_two_hints() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - // Two slots: one for any working-message post, one for the timeout notice. egress.program_response( "slack.com", Ok(EgressResponse::new( @@ -3304,74 +3497,37 @@ mod tests { ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - // Always Running → wait_for_actionable times out after max_wait=1ms. + // BlockedApproval so the state-aware lookup returns the approval copy for both. let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, + TurnStatus::BlockedApproval, )); - // Use FakeConversationBindingService so the binding lookup succeeds and - // the delivery loop can actually reach the timeout. - let binding_service = Arc::new(FakeConversationBindingService::new()); - let thread_service = Arc::new(InMemorySessionThreadService::default()); - let services = SlackFinalReplyDeliveryServices { - binding_service, - thread_service, - turn_coordinator: coordinator, - outbound_store: outbound.clone(), - route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), - communication_preferences: outbound, - adapter: test_adapter(install), - egress: egress.clone(), - delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), - auth_challenges: None, - }; - let settings = SlackFinalReplyDeliverySettings { - poll_interval: std::time::Duration::ZERO, - max_wait: std::time::Duration::from_millis(1), - max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), - max_pending_deliveries: NonZeroUsize::new(8).unwrap(), - }; - let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); - - // Accepted ack for a user message so deliver_final_reply enters the - // polling loop and hits the timeout. - let run_id = TurnRunId::new(); - let env = envelope(user_message_payload()); - let ack = ProductInboundAck::Accepted { - accepted_message_ref: AcceptedMessageRef::new("slack:timeout-test").expect("ref"), - submitted_run_id: run_id, - }; + let observer = make_observer(coordinator, egress.clone(), outbound, install); - observer.observe_workflow_ack(env, ack).await; + // First distinct user message while blocked (fresh active_run_id in each ack). + observer + .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) + .await; + // Second distinct user message while blocked (different active_run_id → different key). + observer + .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) + .await; let calls = egress.calls(); let post_calls: Vec<_> = calls .iter() .filter(|c| c.path == "/api/chat.postMessage") .collect(); - assert!( - !post_calls.is_empty(), - "expected at least one chat.postMessage for timeout notice" - ); - - // The timeout notice must be the final post — a working message may - // precede it, but nothing should be posted after the timeout notice. - let last_body = std::str::from_utf8(&post_calls[post_calls.len() - 1].body).unwrap_or(""); - assert!( - last_body.contains("longer than expected"), - "last chat.postMessage must contain timeout notice text, bodies: {:?}", - post_calls - .iter() - .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) - .collect::>() + assert_eq!( + post_calls.len(), + 2, + "two distinct DeferredBusy + UserMessage deliveries must each post a hint" ); } - /// Accepted ack, then binding lookup fails → generic delivery-error notice - /// posted to the conversation (A3). Drives the observer (the caller), not - /// just `deliver_final_reply`, so the error→feedback mapping in - /// `observe_workflow_ack` is covered. + /// DeferredBusy + UserMessage + BlockedAuth state with gate_ref → auth copy with + /// concrete `auth deny ` command posted. #[tokio::test] - async fn accepted_ack_then_binding_error_posts_delivery_error_notice() { + async fn deferred_busy_blocked_auth_state_posts_auth_hint() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); @@ -3379,41 +3535,20 @@ mod tests { "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D123", "3000.1"), + slack_post_ok_json("D123", "6000.1"), )), ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, - )); - let services = SlackFinalReplyDeliveryServices { - // Errors on lookup_binding, so delivery fails after the Accepted - // ack and before any polling. - binding_service: Arc::new(TestNoopConversationBindingService), - thread_service: Arc::new(InMemorySessionThreadService::default()), - turn_coordinator: coordinator, - outbound_store: outbound.clone(), - route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), - communication_preferences: outbound, - adapter: test_adapter(install), - egress: egress.clone(), - delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), - auth_challenges: None, - }; - let settings = SlackFinalReplyDeliverySettings { - poll_interval: std::time::Duration::ZERO, - max_wait: std::time::Duration::from_millis(1), - max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), - max_pending_deliveries: NonZeroUsize::new(8).unwrap(), - }; - let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); - + // BlockedAuth with concrete gate_ref → hint embeds `auth deny `. + let gate_ref_str = "gate:auth-slack-hint"; + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + TurnStatus::BlockedAuth, + Some(gate_ref_str), + )])); + let observer = make_observer(coordinator, egress.clone(), outbound, install); let env = envelope(user_message_payload()); - let ack = ProductInboundAck::Accepted { - accepted_message_ref: AcceptedMessageRef::new("slack:binding-error-test").expect("ref"), - submitted_run_id: TurnRunId::new(), - }; + let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; @@ -3425,67 +3560,48 @@ mod tests { assert_eq!( post_calls.len(), 1, - "expected exactly one chat.postMessage (the delivery-error notice), bodies: {:?}", - post_calls - .iter() - .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) - .collect::>() + "expected exactly one chat.postMessage for DeferredBusy + BlockedAuth" ); - let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); assert!( - body.contains("Something went wrong delivering the result"), - "post must contain the generic delivery-error notice, body: {body}" + body.contains("authentication step"), + "deferred-busy hint for BlockedAuth must mention 'authentication step', got: {body}" + ); + assert!( + body.contains(gate_ref_str), + "deferred-busy auth hint must embed the concrete gate ref '{gate_ref_str}', got: {body}" + ); + // Must not contain approval command wording. + assert!( + !body.contains("approve") || body.contains("auth deny"), + "auth hint must not contain approval-only wording, got: {body}" ); } - /// Rejected AuthResolution ack → auth-flavored hint posted; approval - /// command text and internal rejection reason must not appear. - /// - /// This is the caller-level regression for `rejection_hint_for_resolution` - /// covering the `ProductInboundPayload::AuthResolution(_)` branch: the hint - /// must come from `user_facing_auth_hint()` (which references `auth deny - /// `), not from `user_facing_hint()` (which references - /// approval commands), and not from the raw internal reason. + /// DeferredBusy + UserMessage + BlockedApproval with no gate_ref → fallback wording + /// without a specific gate command. #[tokio::test] - async fn rejected_auth_resolution_ack_posts_static_hint_not_internal_reason() { + async fn deferred_busy_blocked_approval_no_gate_ref_posts_fallback_hint() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - // Program a success response for the hint post. egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D123", "4000.1"), + slack_post_ok_json("D123", "6001.1"), )), ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, - )); + // BlockedApproval with no gate_ref → static fallback. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + TurnStatus::BlockedApproval, + None, + )])); let observer = make_observer(coordinator, egress.clone(), outbound, install); - - // Build an AuthResolution envelope — this payload kind is included in - // `rejection_hint_for_resolution`'s `is_resolution` match but had no - // caller-level test before this one. - let auth_resolution_payload = ProductInboundPayload::AuthResolution( - ironclaw_product_adapters::AuthResolutionPayload::new( - "gate:auth-hint-test", - ironclaw_product_adapters::AuthResolutionResult::Denied, - ) - .expect("auth resolution payload"), - ); - let env = envelope(auth_resolution_payload); - - // Use a rejection with a distinctive internal reason that must NOT appear - // in the posted message. - let internal_marker = "internal-secret-reason-marker"; - let ack = - ProductInboundAck::Rejected(ironclaw_product_adapters::ProductRejection::permanent( - ironclaw_product_adapters::ProductRejectionKind::BindingRequired, - internal_marker, - )); + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; @@ -3497,41 +3613,18 @@ mod tests { assert_eq!( post_calls.len(), 1, - "expected exactly one chat.postMessage (the hint), bodies: {:?}", - post_calls - .iter() - .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) - .collect::>() - ); - - let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); - - // The posted text must contain the auth-specific hint for BindingRequired, - // not the approval-command variant. - let expected_hint = ironclaw_product_adapters::ProductRejectionKind::BindingRequired - .user_facing_auth_hint(); - assert!( - body.contains(expected_hint), - "post must contain the auth-flavored hint '{expected_hint}', body: {body}" - ); - - // The approval command must NOT appear in an auth-resolution hint. - assert!( - !body.contains("approve gate:"), - "post must not contain approval command 'approve gate:', body: {body}" + "expected one post for fallback approval hint" ); - - // The internal rejection reason must NOT appear in the post. + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); assert!( - !body.contains(internal_marker), - "post must not contain the internal rejection reason '{internal_marker}', body: {body}" + body.contains("waiting on a pending approval"), + "fallback approval hint must still mention pending approval, got: {body}" ); } - /// WorkflowRejected errors after protocol ACK still post resolution hints when - /// the originating conversation is authorized. + /// DeferredBusy + UserMessage + Running state (non-blocked) → generic copy posted. #[tokio::test] - async fn workflow_rejected_resolution_error_posts_authorized_hint() { + async fn deferred_busy_running_state_posts_generic_hint() { let install = "test-install"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); @@ -3539,24 +3632,20 @@ mod tests { "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D123", "4500.1"), + slack_post_ok_json("D123", "6000.2"), )), ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // Running → state-aware lookup returns generic wording. let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( TurnStatus::Running, )); let observer = make_observer(coordinator, egress.clone(), outbound, install); - let env = envelope(scoped_approval_resolution_payload()); - let error = ProductAdapterError::WorkflowRejected { - kind: ProductWorkflowRejectionKind::ScopeNotFound, - status_code: 404, - retryable: false, - reason: ironclaw_product_adapters::RedactedString::new("missing gate"), - }; + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); - observer.observe_workflow_error(env, error).await; + observer.observe_workflow_ack(env, ack).await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3566,108 +3655,34 @@ mod tests { assert_eq!( post_calls.len(), 1, - "expected exactly one hint chat.postMessage call" + "expected exactly one chat.postMessage for DeferredBusy + Running" ); let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); assert!( - body.contains("approve gate:"), - "workflow rejection hint body must contain approval guidance, got: {body}" - ); - assert!( - !body.contains("missing gate"), - "workflow rejection hint must not expose redacted reason, got: {body}" - ); - } - - /// If route/binding authorization fails, rejected-resolution feedback is - /// suppressed instead of posting to an arbitrary shared Slack conversation. - #[tokio::test] - async fn workflow_rejected_resolution_error_without_binding_posts_nothing() { - let install = "test-install"; - let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); - egress.allow_credential_handle("slack_bot_token"); - - let outbound = Arc::new(InMemoryOutboundStateStore::default()); - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, - )); - let thread_service = Arc::new(InMemorySessionThreadService::default()); - let services = make_services( - coordinator, - thread_service, - egress.clone(), - outbound, - install, - ); - let settings = SlackFinalReplyDeliverySettings { - poll_interval: std::time::Duration::ZERO, - max_wait: std::time::Duration::from_millis(1), - max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), - max_pending_deliveries: NonZeroUsize::new(8).unwrap(), - }; - let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); - let env = envelope(scoped_approval_resolution_payload()); - let error = ProductAdapterError::WorkflowRejected { - kind: ProductWorkflowRejectionKind::ScopeNotFound, - status_code: 404, - retryable: false, - reason: ironclaw_product_adapters::RedactedString::new("missing gate"), - }; - - observer.observe_workflow_error(env, error).await; - - let calls = egress.calls(); - assert!( - !calls.iter().any(|c| c.path == "/api/chat.postMessage"), - "no chat.postMessage expected when binding authorization fails" + body.contains("still working on a previous message"), + "deferred-busy hint for Running state must contain generic copy, got: {body}" ); } - /// When a blocked-state notification (approval prompt) was delivered and - /// the subsequent wait times out, no additional timeout notice must be - /// posted to Slack. + /// DeferredBusy + UserMessage + unauthorized binding → no post, no panic. /// - /// This is the caller-level regression for the - /// `RunWaitTimedOutAfterNotification` error variant: `observe_workflow_ack` - /// maps this variant to `feedback = None` so the user is not double-notified - /// after already seeing the approval prompt. + /// Uses `TestNoopConversationBindingService` (always fails lookup_binding) to + /// simulate an unauthorized conversation: the observer must skip the hint post + /// silently rather than posting to an arbitrary channel. #[tokio::test] - async fn timeout_after_blocked_notification_suppresses_timeout_message() { - use ironclaw_product_workflow::FakeConversationBindingService; - + async fn deferred_busy_unauthorized_binding_posts_nothing() { let install = "test-install"; - let gate_ref_str = "gate:approval-timeout-test"; let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - // One programmed response for the approval-prompt postMessage. - // No second response — if the timeout notice were posted, the test - // would fail because `FakeProtocolHttpEgress` returns an error on - // an empty queue. - egress.program_response( - "slack.com", - Ok(EgressResponse::new( - 200, - slack_post_ok_json("D123", "5000.1"), - )), - ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - // Always BlockedApproval with the same gate_ref: the first poll exits - // `wait_for_actionable` immediately (new blocked state, different from - // `delivered_blocked_marker=None`), delivering the approval prompt. - // Subsequent polls (second call to `wait_for_actionable`) return the same - // marker as `delivered_blocked_marker`, so the loop does not exit — it - // times out after `max_wait=1ms` and the error is converted to - // `RunWaitTimedOutAfterNotification`. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( TurnStatus::BlockedApproval, - Some(gate_ref_str), - )])); - - let binding_service = Arc::new(FakeConversationBindingService::new()); + )); + // Override the default `make_observer` so we can inject the no-binding service. let thread_service = Arc::new(InMemorySessionThreadService::default()); let services = SlackFinalReplyDeliveryServices { - binding_service, + binding_service: Arc::new(TestNoopConversationBindingService), thread_service, turn_coordinator: coordinator, outbound_store: outbound.clone(), @@ -3686,72 +3701,177 @@ mod tests { }; let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); - let run_id = TurnRunId::new(); let env = envelope(user_message_payload()); - let ack = ProductInboundAck::Accepted { - accepted_message_ref: AcceptedMessageRef::new("slack:blocked-timeout-test") - .expect("ref"), - submitted_run_id: run_id, - }; + let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; let calls = egress.calls(); - let post_calls: Vec<_> = calls - .iter() - .filter(|c| c.path == "/api/chat.postMessage") - .collect(); - - // Exactly one postMessage: the approval prompt notification. - // No timeout notice must have been posted. - assert_eq!( - post_calls.len(), - 1, - "expected exactly one chat.postMessage (the approval prompt), bodies: {:?}", - post_calls - .iter() - .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) - .collect::>() - ); - - let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); - - // The one message must be the approval prompt, not a timeout notice. - assert!( - !body.contains("longer than expected"), - "timeout notice must not be posted after blocked-notification timeout, body: {body}" - ); - // The approval prompt must reference the gate ref. assert!( - body.contains(gate_ref_str), - "approval prompt must reference the gate ref, body: {body}" + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "no chat.postMessage expected when binding authorization fails" ); } - /// Same suppression as the approval-prompt case, but through the auth prompt - /// payload path. + /// Duplicate { prior: Accepted } → nothing posted (already succeeded). #[tokio::test] - async fn timeout_after_auth_blocked_notification_suppresses_timeout_message() { - use ironclaw_product_workflow::FakeConversationBindingService; - + async fn duplicate_accepted_ack_posts_nothing() { let install = "test-install"; - let gate_ref_str = "gate:auth-timeout-test"; + let run_id = TurnRunId::new(); let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - egress.program_response( - "slack.com", - Ok(EgressResponse::new( - 200, - slack_post_ok_json("D123", "5000.2"), - )), - ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( - TurnStatus::BlockedAuth, - Some(gate_ref_str), - )])); + // The coordinator is a no-op because Duplicate{Accepted} has no submitted_run_id. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + let env = envelope(scoped_approval_resolution_payload()); + let prior = ProductInboundAck::Accepted { + accepted_message_ref: AcceptedMessageRef::new("slack:prior").expect("ref"), + submitted_run_id: run_id, + }; + let ack = ProductInboundAck::Duplicate { + prior: Box::new(prior), + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "no chat.postMessage expected for Duplicate{{Accepted}}" + ); + } + + /// Duplicate { prior: Rejected } → nothing posted at observer level. + #[tokio::test] + async fn duplicate_rejected_ack_posts_nothing() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + let env = envelope(scoped_approval_resolution_payload()); + let ack = ProductInboundAck::Duplicate { + prior: Box::new(rejected_ack( + ironclaw_product_adapters::ProductRejectionKind::BindingRequired, + )), + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "no chat.postMessage expected for Duplicate{{Rejected}}" + ); + } + + /// All `Duplicate` acks suppress the hint, regardless of the prior inside. + /// + /// `Duplicate` is keyed on the external event id: transport retries of the + /// same Slack event land as `Duplicate{original}`. The original processing + /// already posted any hint, so replays must not repeat the side effect. + /// A user re-typing "approve" produces a new event id → a fresh `Rejected`, + /// never a `Duplicate`, so suppressing `Duplicate{Rejected}` loses nothing. + #[test] + fn duplicate_acks_produce_no_hint() { + let env = envelope(scoped_approval_resolution_payload()); + + let duplicate_rejected = ProductInboundAck::Duplicate { + prior: Box::new(rejected_ack( + ironclaw_product_adapters::ProductRejectionKind::BindingRequired, + )), + }; + assert!( + rejection_hint_for_resolution(&env, &duplicate_rejected).is_none(), + "Duplicate{{Rejected}} must NOT produce a hint (transport replay)" + ); + + let duplicate_accepted = ProductInboundAck::Duplicate { + prior: Box::new(ProductInboundAck::Accepted { + accepted_message_ref: ironclaw_turns::AcceptedMessageRef::new("slack:prior") + .expect("ref"), + submitted_run_id: ironclaw_turns::TurnRunId::new(), + }), + }; + assert!( + rejection_hint_for_resolution(&env, &duplicate_accepted).is_none(), + "Duplicate{{Accepted}} must not produce a hint" + ); + } + + /// A failed best-effort rejection-hint post must not fall through into generic + /// delivery-error feedback. + #[tokio::test] + async fn rejected_resolution_hint_post_failure_is_best_effort() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + let env = envelope(scoped_approval_resolution_payload()); + let ack = rejected_ack(ironclaw_product_adapters::ProductRejectionKind::BindingRequired); + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one failed best-effort hint post" + ); + } + + /// Delivery error (RunWaitTimedOut) → timeout notice posted to conversation. + /// + /// Uses `FakeConversationBindingService` so the binding lookup succeeds and + /// delivery enters the polling loop, which then times out because the + /// coordinator always returns `Running`. + #[tokio::test] + async fn delivery_timeout_posts_timeout_notice_to_conversation() { + use ironclaw_product_workflow::FakeConversationBindingService; + + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // Two slots: one for any working-message post, one for the timeout notice. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "2000.1"), + )), + ); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "2000.2"), + )), + ); + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // Always Running → wait_for_actionable times out after max_wait=1ms. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + // Use FakeConversationBindingService so the binding lookup succeeds and + // the delivery loop can actually reach the timeout. let binding_service = Arc::new(FakeConversationBindingService::new()); let thread_service = Arc::new(InMemorySessionThreadService::default()); let services = SlackFinalReplyDeliveryServices { @@ -3774,10 +3894,86 @@ mod tests { }; let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + // Accepted ack for a user message so deliver_final_reply enters the + // polling loop and hits the timeout. + let run_id = TurnRunId::new(); let env = envelope(user_message_payload()); let ack = ProductInboundAck::Accepted { - accepted_message_ref: AcceptedMessageRef::new("slack:auth-blocked-timeout-test") - .expect("ref"), + accepted_message_ref: AcceptedMessageRef::new("slack:timeout-test").expect("ref"), + submitted_run_id: run_id, + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert!( + !post_calls.is_empty(), + "expected at least one chat.postMessage for timeout notice" + ); + + // The timeout notice must be the final post — a working message may + // precede it, but nothing should be posted after the timeout notice. + let last_body = std::str::from_utf8(&post_calls[post_calls.len() - 1].body).unwrap_or(""); + assert!( + last_body.contains("longer than expected"), + "last chat.postMessage must contain timeout notice text, bodies: {:?}", + post_calls + .iter() + .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) + .collect::>() + ); + } + + /// Accepted ack, then binding lookup fails → generic delivery-error notice + /// posted to the conversation (A3). Drives the observer (the caller), not + /// just `deliver_final_reply`, so the error→feedback mapping in + /// `observe_workflow_ack` is covered. + #[tokio::test] + async fn accepted_ack_then_binding_error_posts_delivery_error_notice() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "3000.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let services = SlackFinalReplyDeliveryServices { + // Errors on lookup_binding, so delivery fails after the Accepted + // ack and before any polling. + binding_service: Arc::new(TestNoopConversationBindingService), + thread_service: Arc::new(InMemorySessionThreadService::default()), + turn_coordinator: coordinator, + outbound_store: outbound.clone(), + route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), + communication_preferences: outbound, + adapter: test_adapter(install), + egress: egress.clone(), + delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), + auth_challenges: None, + }; + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_millis(1), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + + let env = envelope(user_message_payload()); + let ack = ProductInboundAck::Accepted { + accepted_message_ref: AcceptedMessageRef::new("slack:binding-error-test").expect("ref"), submitted_run_id: TurnRunId::new(), }; @@ -3791,7 +3987,7 @@ mod tests { assert_eq!( post_calls.len(), 1, - "expected exactly one chat.postMessage (the auth prompt), bodies: {:?}", + "expected exactly one chat.postMessage (the delivery-error notice), bodies: {:?}", post_calls .iter() .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) @@ -3799,435 +3995,1263 @@ mod tests { ); let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); assert!( - !body.contains("longer than expected"), - "timeout notice must not be posted after auth notification timeout, body: {body}" - ); - assert!( - body.contains(gate_ref_str), - "auth prompt must reference the gate ref, body: {body}" + body.contains("Something went wrong delivering the result"), + "post must contain the generic delivery-error notice, body: {body}" ); } - // --- OnceLock slot behaviour tests ---------------------------------------- - + /// Rejected AuthResolution ack → auth-flavored hint posted; approval + /// command text and internal rejection reason must not appear. + /// + /// This is the caller-level regression for `rejection_hint_for_resolution` + /// covering the `ProductInboundPayload::AuthResolution(_)` branch: the hint + /// must come from `user_facing_auth_hint()` (which references `auth deny + /// `), not from `user_facing_hint()` (which references + /// approval commands), and not from the raw internal reason. #[tokio::test] - async fn post_submit_hook_slot_empty_hook_does_not_fire() { - // Verify the contract that `PostSubmitHookWrappedSubmitter` relies on: - // when the OnceLock slot is empty, reading it returns None and no hook fires. - let slot: Arc>> = Arc::new(OnceLock::new()); - + async fn rejected_auth_resolution_ack_posts_static_hint_not_internal_reason() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // Program a success response for the hint post. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "4000.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + + // Build an AuthResolution envelope — this payload kind is included in + // `rejection_hint_for_resolution`'s `is_resolution` match but had no + // caller-level test before this one. + let auth_resolution_payload = ProductInboundPayload::AuthResolution( + ironclaw_product_adapters::AuthResolutionPayload::new( + "gate:auth-hint-test", + ironclaw_product_adapters::AuthResolutionResult::Denied, + ) + .expect("auth resolution payload"), + ); + let env = envelope(auth_resolution_payload); + + // Use a rejection with a distinctive internal reason that must NOT appear + // in the posted message. + let internal_marker = "internal-secret-reason-marker"; + let ack = + ProductInboundAck::Rejected(ironclaw_product_adapters::ProductRejection::permanent( + ironclaw_product_adapters::ProductRejectionKind::BindingRequired, + internal_marker, + )); + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage (the hint), bodies: {:?}", + post_calls + .iter() + .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) + .collect::>() + ); + + let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); + + // The posted text must contain the auth-specific hint for BindingRequired, + // not the approval-command variant. + let expected_hint = ironclaw_product_adapters::ProductRejectionKind::BindingRequired + .user_facing_auth_hint(); + assert!( + body.contains(expected_hint), + "post must contain the auth-flavored hint '{expected_hint}', body: {body}" + ); + + // The approval command must NOT appear in an auth-resolution hint. + assert!( + !body.contains("approve gate:"), + "post must not contain approval command 'approve gate:', body: {body}" + ); + + // The internal rejection reason must NOT appear in the post. + assert!( + !body.contains(internal_marker), + "post must not contain the internal rejection reason '{internal_marker}', body: {body}" + ); + } + + /// WorkflowRejected errors after protocol ACK still post resolution hints when + /// the originating conversation is authorized. + #[tokio::test] + async fn workflow_rejected_resolution_error_posts_authorized_hint() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "4500.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + let env = envelope(scoped_approval_resolution_payload()); + let error = ProductAdapterError::WorkflowRejected { + kind: ProductWorkflowRejectionKind::ScopeNotFound, + status_code: 404, + retryable: false, + reason: ironclaw_product_adapters::RedactedString::new("missing gate"), + }; + + observer.observe_workflow_error(env, error).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one hint chat.postMessage call" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("approve gate:"), + "workflow rejection hint body must contain approval guidance, got: {body}" + ); + assert!( + !body.contains("missing gate"), + "workflow rejection hint must not expose redacted reason, got: {body}" + ); + } + + /// If route/binding authorization fails, rejected-resolution feedback is + /// suppressed instead of posting to an arbitrary shared Slack conversation. + #[tokio::test] + async fn workflow_rejected_resolution_error_without_binding_posts_nothing() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::Running, + )); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + let services = make_services( + coordinator, + thread_service, + egress.clone(), + outbound, + install, + ); + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_millis(1), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + let env = envelope(scoped_approval_resolution_payload()); + let error = ProductAdapterError::WorkflowRejected { + kind: ProductWorkflowRejectionKind::ScopeNotFound, + status_code: 404, + retryable: false, + reason: ironclaw_product_adapters::RedactedString::new("missing gate"), + }; + + observer.observe_workflow_error(env, error).await; + + let calls = egress.calls(); + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "no chat.postMessage expected when binding authorization fails" + ); + } + + /// When a blocked-state notification (approval prompt) was delivered and + /// the subsequent wait times out, no additional timeout notice must be + /// posted to Slack. + /// + /// This is the caller-level regression for the + /// `RunWaitTimedOutAfterNotification` error variant: `observe_workflow_ack` + /// maps this variant to `feedback = None` so the user is not double-notified + /// after already seeing the approval prompt. + #[tokio::test] + async fn timeout_after_blocked_notification_suppresses_timeout_message() { + use ironclaw_product_workflow::FakeConversationBindingService; + + let install = "test-install"; + let gate_ref_str = "gate:approval-timeout-test"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // One programmed response for the approval-prompt postMessage. + // No second response — if the timeout notice were posted, the test + // would fail because `FakeProtocolHttpEgress` returns an error on + // an empty queue. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "5000.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // Always BlockedApproval with the same gate_ref: the first poll exits + // `wait_for_actionable` immediately (new blocked state, different from + // `delivered_blocked_marker=None`), delivering the approval prompt. + // Subsequent polls (second call to `wait_for_actionable`) return the same + // marker as `delivered_blocked_marker`, so the loop does not exit — it + // times out after `max_wait=1ms` and the error is converted to + // `RunWaitTimedOutAfterNotification`. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + TurnStatus::BlockedApproval, + Some(gate_ref_str), + )])); + + let binding_service = Arc::new(FakeConversationBindingService::new()); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + let services = SlackFinalReplyDeliveryServices { + binding_service, + thread_service, + turn_coordinator: coordinator, + outbound_store: outbound.clone(), + route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), + communication_preferences: outbound, + adapter: test_adapter(install), + egress: egress.clone(), + delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), + auth_challenges: None, + }; + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_millis(1), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + + let run_id = TurnRunId::new(); + let env = envelope(user_message_payload()); + let ack = ProductInboundAck::Accepted { + accepted_message_ref: AcceptedMessageRef::new("slack:blocked-timeout-test") + .expect("ref"), + submitted_run_id: run_id, + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + + // Exactly one postMessage: the approval prompt notification. + // No timeout notice must have been posted. + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage (the approval prompt), bodies: {:?}", + post_calls + .iter() + .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) + .collect::>() + ); + + let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); + + // The one message must be the approval prompt, not a timeout notice. + assert!( + !body.contains("longer than expected"), + "timeout notice must not be posted after blocked-notification timeout, body: {body}" + ); + // The approval prompt must reference the gate ref. + assert!( + body.contains(gate_ref_str), + "approval prompt must reference the gate ref, body: {body}" + ); + } + + /// Same suppression as the approval-prompt case, but through the auth prompt + /// payload path. + #[tokio::test] + async fn timeout_after_auth_blocked_notification_suppresses_timeout_message() { + use ironclaw_product_workflow::FakeConversationBindingService; + + let install = "test-install"; + let gate_ref_str = "gate:auth-timeout-test"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "5000.2"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![scripted_state( + TurnStatus::BlockedAuth, + Some(gate_ref_str), + )])); + + let binding_service = Arc::new(FakeConversationBindingService::new()); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + let services = SlackFinalReplyDeliveryServices { + binding_service, + thread_service, + turn_coordinator: coordinator, + outbound_store: outbound.clone(), + route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), + communication_preferences: outbound, + adapter: test_adapter(install), + egress: egress.clone(), + delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), + auth_challenges: None, + }; + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_millis(1), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + + let env = envelope(user_message_payload()); + let ack = ProductInboundAck::Accepted { + accepted_message_ref: AcceptedMessageRef::new("slack:auth-blocked-timeout-test") + .expect("ref"), + submitted_run_id: TurnRunId::new(), + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage (the auth prompt), bodies: {:?}", + post_calls + .iter() + .map(|c| std::str::from_utf8(&c.body).unwrap_or("?")) + .collect::>() + ); + let body = std::str::from_utf8(&post_calls[0].body).unwrap_or(""); + assert!( + !body.contains("longer than expected"), + "timeout notice must not be posted after auth notification timeout, body: {body}" + ); + assert!( + body.contains(gate_ref_str), + "auth prompt must reference the gate ref, body: {body}" + ); + } + + // --- OnceLock slot behaviour tests ---------------------------------------- + + #[tokio::test] + async fn post_submit_hook_slot_empty_hook_does_not_fire() { + // Verify the contract that `PostSubmitHookWrappedSubmitter` relies on: + // when the OnceLock slot is empty, reading it returns None and no hook fires. + let slot: Arc>> = Arc::new(OnceLock::new()); + // Slot is empty — `get()` returns None, so no hook fires. assert!(slot.get().is_none(), "empty slot must return None"); - // When the slot IS occupied the hook IS reachable. - let fired = Arc::new(Mutex::new(false)); - let fired_clone = Arc::clone(&fired); - struct FlagHook(Arc>); - #[async_trait] - impl PostSubmitDeliveryHook for FlagHook { - async fn on_trigger_submitted( - &self, - _fire: TriggerFire, - _run_id: TurnRunId, - _scope: TurnScope, - ) { - *self.0.lock().expect("flag") = true; - } - } - slot.set(Arc::new(FlagHook(fired_clone))) - .unwrap_or_else(|_| panic!("first slot set should succeed")); - if let Some(hook) = slot.get() { - let fire = minimal_trigger_fire(None); - hook.on_trigger_submitted(fire, TurnRunId::new(), personal_turn_scope()) - .await; - } + // When the slot IS occupied the hook IS reachable. + let fired = Arc::new(Mutex::new(false)); + let fired_clone = Arc::clone(&fired); + struct FlagHook(Arc>); + #[async_trait] + impl PostSubmitDeliveryHook for FlagHook { + async fn on_trigger_submitted( + &self, + _fire: TriggerFire, + _run_id: TurnRunId, + _scope: TurnScope, + ) { + *self.0.lock().expect("flag") = true; + } + } + slot.set(Arc::new(FlagHook(fired_clone))) + .unwrap_or_else(|_| panic!("first slot set should succeed")); + if let Some(hook) = slot.get() { + let fire = minimal_trigger_fire(None); + hook.on_trigger_submitted(fire, TurnRunId::new(), personal_turn_scope()) + .await; + } + assert!( + *fired.lock().expect("flag"), + "hook must fire after slot is set" + ); + } + + #[test] + fn post_submit_hook_slot_second_set_is_noop() { + // OnceLock semantics: first set succeeds, second set returns the value. + // This is the contract behind `set_trigger_post_submit_hook` returning false + // on duplicate calls. + let slot: Arc>> = Arc::new(OnceLock::new()); + let hook_a: Arc = Arc::new(NoopPostSubmitDeliveryHook); + let hook_b: Arc = Arc::new(NoopPostSubmitDeliveryHook); + + assert!(slot.set(hook_a).is_ok(), "first set should succeed"); + assert!( + slot.set(hook_b).is_err(), + "second set should fail (slot already occupied)" + ); + assert!( + slot.get().is_some(), + "slot still occupied after duplicate set" + ); + } + + #[test] + fn slack_approval_prompt_offers_always_for_typed_approval_gate() { + let gate_ref = GateRef::new(format!( + "gate:approval-{}", + ironclaw_host_api::ApprovalRequestId::new() + )) + .expect("gate ref"); + + let prompt = slack_approval_gate_prompt_view(TurnRunId::new(), &gate_ref); + + assert_eq!(prompt.gate_ref, gate_ref.as_str()); + assert!(prompt.allow_always); + } + + #[test] + fn slack_approval_prompt_does_not_offer_always_for_generic_gate() { + let gate_ref = GateRef::new("gate:approve-slack").expect("gate ref"); + + let prompt = slack_approval_gate_prompt_view(TurnRunId::new(), &gate_ref); + + assert_eq!(prompt.gate_ref, gate_ref.as_str()); + assert!(!prompt.allow_always); + } + + // --- Bug-fix regression tests: gate-route refs carry team id (space_id) ---- + + /// Test A: triggered approval delivery records a gate-route that includes a + /// posted-message ref whose fingerprint matches an inbound-style ref carrying + /// the Slack team id (space_id = "T123"). + /// + /// The test_slack_binding_ref helper encodes space = "T123", conversation = "D456". + /// After delivery the authority's `resolved_space_id` must be Some("T123"), + /// and the recorded route must contain a ref with space_id = Some("T123"), + /// conversation_id = the channel returned by Slack ("D456"), and thread_id = + /// the ts of the posted message ("1111.2222"). + #[tokio::test] + async fn triggered_approval_route_ref_carries_resolved_space_id() { + let install = "test-install"; + let gate_ref_str = "gate:approval-space-test"; + let scope = personal_turn_scope(); + let run_id = TurnRunId::new(); + let binding_ref = + test_slack_binding_ref(install, scope.agent_id.as_ref().expect("agent").as_str()); + + // First poll → BlockedApproval; second poll → Completed. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![ + scripted_state(TurnStatus::BlockedApproval, Some(gate_ref_str)), + scripted_state(TurnStatus::Completed, None), + ])); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + seed_finalized_assistant_message( + &thread_service, + &scope, + run_id, + "Run complete after approval.", + ) + .await; + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + seed_personal_preference(&outbound, &scope, binding_ref).await; + + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // Approval-prompt response: channel D456, ts 1111.2222. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D456", "1111.2222"), + )), + ); + // Final-reply response. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D456", "3333.4444"), + )), + ); + + let delivery_store = Arc::new(InMemoryTriggeredRunDeliveryStore::default()); + let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); + let services = make_services( + coordinator, + thread_service, + egress.clone(), + outbound, + install, + ); + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_secs(5), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let driver = TriggeredRunDeliveryDriver::with_settings( + services, + settings, + delivery_store.clone(), + route_store.clone(), + scope.agent_id.clone().expect("test scope has agent"), + ); + + let fire = minimal_trigger_fire(None); + driver + .on_trigger_submitted(fire, run_id, scope.clone()) + .await; + wait_for_delivery_record(&delivery_store, run_id).await; + + let creator = ironclaw_host_api::UserId::new("creator-user").expect("user id"); + let route = route_store + .load_delivered_gate_route(&scope.tenant_id, &creator, gate_ref_str) + .await + .expect("load route") + .expect("gate route was recorded"); + + // The binding ref encodes space = "T123", so the recorded route must + // contain a ref that fingerprint-matches an inbound-style ref with + // space_id = Some("T123"), conversation_id = "D456", thread_id = "1111.2222". + let expected_inbound_ref = ironclaw_conversations::ExternalConversationRef::new( + Some("T123"), + "D456", + Some("1111.2222"), + None, + ) + .expect("expected inbound ref"); + let expected_fingerprint = expected_inbound_ref.conversation_fingerprint(); + assert!( - *fired.lock().expect("flag"), - "hook must fire after slot is set" + route + .delivered_conversation_fingerprints + .iter() + .any(|fingerprint| fingerprint == &expected_fingerprint), + "recorded route must include space_id=T123 fingerprint; fingerprints={:?}", + route.delivered_conversation_fingerprints, ); } - #[test] - fn post_submit_hook_slot_second_set_is_noop() { - // OnceLock semantics: first set succeeds, second set returns the value. - // This is the contract behind `set_trigger_post_submit_hook` returning false - // on duplicate calls. - let slot: Arc>> = Arc::new(OnceLock::new()); - let hook_a: Arc = Arc::new(NoopPostSubmitDeliveryHook); - let hook_b: Arc = Arc::new(NoopPostSubmitDeliveryHook); + /// Test B: triggered auth (BlockedAuth) delivery records a gate-route keyed + /// by the auth gate_ref. This is the Bug-2 regression: previously + /// `gate_ref_for_routing` was `None` for BlockedAuth so no route was + /// recorded, causing a `MissingGate` when the user replied "approve". + #[tokio::test] + async fn triggered_auth_delivery_records_gate_route() { + let install = "test-install"; + let gate_ref_str = "gate:auth-route-regression"; + let scope = personal_turn_scope(); + let run_id = TurnRunId::new(); + let binding_ref = + test_slack_binding_ref(install, scope.agent_id.as_ref().expect("agent").as_str()); + + // First poll → BlockedAuth with gate_ref; second poll → Completed. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![ + scripted_state(TurnStatus::BlockedAuth, Some(gate_ref_str)), + scripted_state(TurnStatus::Completed, None), + ])); + let thread_service = Arc::new(InMemorySessionThreadService::default()); + seed_finalized_assistant_message( + &thread_service, + &scope, + run_id, + "Run complete after auth.", + ) + .await; + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + seed_personal_preference(&outbound, &scope, binding_ref).await; + + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // Auth-prompt delivery response. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D456", "9999.1111"), + )), + ); + // Final-reply response (after Completed). + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D456", "9999.2222"), + )), + ); + // Auth message is deleted after final; need a delete response. + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + serde_json::json!({"ok": true}).to_string().into_bytes(), + )), + ); + + let delivery_store = Arc::new(InMemoryTriggeredRunDeliveryStore::default()); + let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); + let services = make_services( + coordinator, + thread_service, + egress.clone(), + outbound, + install, + ); + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_secs(5), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let driver = TriggeredRunDeliveryDriver::with_settings( + services, + settings, + delivery_store.clone(), + route_store.clone(), + scope.agent_id.clone().expect("test scope has agent"), + ); + + let fire = minimal_trigger_fire(None); + driver + .on_trigger_submitted(fire, run_id, scope.clone()) + .await; + wait_for_delivery_record(&delivery_store, run_id).await; + + let creator = ironclaw_host_api::UserId::new("creator-user").expect("user id"); + // A gate route for the auth gate_ref must have been recorded. + let route = route_store + .load_delivered_gate_route(&scope.tenant_id, &creator, gate_ref_str) + .await + .expect("load route") + .expect("auth gate route must be recorded (Bug 2 regression)"); + + assert_eq!( + route.gate_ref, gate_ref_str, + "recorded route must reference the auth gate_ref" + ); + assert_eq!(route.run_id, run_id); + } + + /// Test C: the `record_gate_route_if_needed` helper, called from the + /// observer path, stores a posted-message ref whose fingerprint matches an + /// inbound-style ref that carries the envelope's team id (space_id). + /// + /// This tests the observer call site directly (without driving the full + /// observer loop) to verify that `envelope_space_id` is extracted and + /// passed through. + #[tokio::test] + async fn observer_approval_route_ref_carries_envelope_space_id() { + let tenant_id = ironclaw_host_api::TenantId::new("test-tenant").expect("tenant"); + let user_id = ironclaw_host_api::UserId::new("user-obs").expect("user"); + let run_id = TurnRunId::new(); + let gate_ref_str = "gate:observer-space-test"; + let agent = ironclaw_host_api::AgentId::new("obs-agent").expect("agent"); + let thread = ironclaw_host_api::ThreadId::new("obs-thread").expect("thread"); + let scope = TurnScope::new_with_owner(tenant_id.clone(), Some(agent), None, thread, None); + + // Simulate a posted message: channel D789, ts 5555.6666. + let posted = vec![PostedSlackMessage { + channel: "D789".to_string(), + ts: "5555.6666".to_string(), + }]; + + // Envelope conv ref carries space_id = "T999" (the team id). + let envelope_conv_ref = + ExternalConversationRef::new(Some("T999"), "D789", Some("5555.6666"), None) + .expect("envelope conv ref"); + + let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); + + // Derive space_id from the envelope ref — mirrors the observer call site. + let envelope_space_id = conversations_ref_from_product_ref(&envelope_conv_ref) + .ok() + .and_then(|r| r.space_id().map(str::to_string)); + assert_eq!( + envelope_space_id.as_deref(), + Some("T999"), + "space_id must be extracted from envelope ref" + ); + + record_gate_route_if_needed( + route_store.as_ref(), + run_id, + &tenant_id, + &user_id, + gate_ref_str, + &scope, + &posted, + Some(&envelope_conv_ref), + envelope_space_id.as_deref(), + ) + .await; + + let route = route_store + .load_delivered_gate_route(&tenant_id, &user_id, gate_ref_str) + .await + .expect("load route") + .expect("route was recorded"); + + // Must contain a ref with space_id = "T999" matching the inbound fingerprint. + let expected_inbound_ref = ironclaw_conversations::ExternalConversationRef::new( + Some("T999"), + "D789", + Some("5555.6666"), + None, + ) + .expect("inbound ref"); + let expected_fingerprint = expected_inbound_ref.conversation_fingerprint(); - assert!(slot.set(hook_a).is_ok(), "first set should succeed"); assert!( - slot.set(hook_b).is_err(), - "second set should fail (slot already occupied)" + route + .delivered_conversation_fingerprints + .iter() + .any(|fingerprint| fingerprint == &expected_fingerprint), + "recorded route must include space_id=T999 fingerprint; fingerprints={:?}", + route.delivered_conversation_fingerprints, + ); + + // Also verify that the no-space fallback variant is present (inbound + // events without team_id must still match). + let no_space_ref = ironclaw_conversations::ExternalConversationRef::new( + None, + "D789", + Some("5555.6666"), + None, + ) + .expect("no-space ref"); + let no_space_fingerprint = no_space_ref.conversation_fingerprint(); + assert!( + route + .delivered_conversation_fingerprints + .iter() + .any(|fingerprint| fingerprint == &no_space_fingerprint), + "recorded route must include the no-space fallback fingerprint; fingerprints={:?}", + route.delivered_conversation_fingerprints, + ); + + let channel_root_ref = + ironclaw_conversations::ExternalConversationRef::new(Some("T999"), "D789", None, None) + .expect("channel root ref"); + let channel_root_fingerprint = channel_root_ref.conversation_fingerprint(); + assert!( + route + .delivered_conversation_fingerprints + .iter() + .any(|fingerprint| fingerprint == &channel_root_fingerprint), + "recorded route must include the space-qualified channel-root fingerprint for bare replies; fingerprints={:?}", + route.delivered_conversation_fingerprints, ); + + let no_space_channel_root_ref = + ironclaw_conversations::ExternalConversationRef::new(None, "D789", None, None) + .expect("no-space channel root ref"); + let no_space_channel_root_fingerprint = + no_space_channel_root_ref.conversation_fingerprint(); assert!( - slot.get().is_some(), - "slot still occupied after duplicate set" + route + .delivered_conversation_fingerprints + .iter() + .any(|fingerprint| fingerprint == &no_space_channel_root_fingerprint), + "recorded route must include the no-space channel-root fingerprint for bare replies; fingerprints={:?}", + route.delivered_conversation_fingerprints, ); } - #[test] - fn slack_approval_prompt_offers_always_for_typed_approval_gate() { - let gate_ref = GateRef::new(format!( - "gate:approval-{}", - ironclaw_host_api::ApprovalRequestId::new() - )) - .expect("gate ref"); + // ── Extra DeferredBusy coverage tests ───────────────────────────────────── - let prompt = slack_approval_gate_prompt_view(TurnRunId::new(), &gate_ref); + /// Binding service that always returns a binding with `agent_id = None`. + /// + /// Used to exercise the scope-derivation fallback in + /// `deferred_busy_hint_from_run_state`: when `thread_scope_from_binding` fails + /// because `agent_id` is missing, the hint must still be posted using the + /// generic copy rather than being silently dropped. + struct NoAgentConversationBindingService; - assert_eq!(prompt.gate_ref, gate_ref.as_str()); - assert!(prompt.allow_always); + #[async_trait] + impl ConversationBindingService for NoAgentConversationBindingService { + async fn resolve_binding( + &self, + request: ResolveBindingRequest, + ) -> Result { + Ok(ResolvedBinding { + tenant_id: ironclaw_host_api::TenantId::new("tenant:test").expect("tenant"), + actor_user_id: ironclaw_host_api::UserId::new(format!( + "user:{}", + request.external_actor_ref.id() + )) + .expect("user"), + subject_user_id: None, + thread_id: ironclaw_host_api::ThreadId::new("thread:test").expect("thread"), + agent_id: None, // deliberately no agent — triggers scope derivation failure + project_id: None, + }) + } + + async fn lookup_binding( + &self, + request: ResolveBindingRequest, + ) -> Result { + self.resolve_binding(request).await + } } - #[test] - fn slack_approval_prompt_does_not_offer_always_for_generic_gate() { - let gate_ref = GateRef::new("gate:approve-slack").expect("gate ref"); + /// A `TurnCoordinator` double whose `get_run_state` always returns `Err`. + struct ErroringTurnCoordinator; - let prompt = slack_approval_gate_prompt_view(TurnRunId::new(), &gate_ref); + #[async_trait] + impl TurnCoordinator for ErroringTurnCoordinator { + async fn prepare_turn(&self, _scope: TurnScope) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } - assert_eq!(prompt.gate_ref, gate_ref.as_str()); - assert!(!prompt.allow_always); - } + async fn submit_turn( + &self, + _request: SubmitTurnRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } - // --- Bug-fix regression tests: gate-route refs carry team id (space_id) ---- + async fn resume_turn( + &self, + _request: ResumeTurnRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } - /// Test A: triggered approval delivery records a gate-route that includes a - /// posted-message ref whose fingerprint matches an inbound-style ref carrying - /// the Slack team id (space_id = "T123"). + async fn get_run_state( + &self, + _request: GetRunStateRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "simulated run-state lookup failure".to_string(), + }) + } + + async fn cancel_run( + &self, + _request: ironclaw_turns::CancelRunRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } + } + + /// Binding with no `agent_id` → scope derivation fails → generic copy posted. /// - /// The test_slack_binding_ref helper encodes space = "T123", conversation = "D456". - /// After delivery the authority's `resolved_space_id` must be Some("T123"), - /// and the recorded route must contain a ref with space_id = Some("T123"), - /// conversation_id = the channel returned by Slack ("D456"), and thread_id = - /// the ts of the posted message ("1111.2222"). + /// `deferred_busy_hint_from_run_state` calls `thread_scope_from_binding` which + /// returns `Err` when `agent_id` is `None`. The code must fall back to + /// `SLACK_DEFERRED_BUSY_GENERIC_MESSAGE` and still post the hint. #[tokio::test] - async fn triggered_approval_route_ref_carries_resolved_space_id() { + async fn deferred_busy_missing_agent_binding_posts_generic_hint() { let install = "test-install"; - let gate_ref_str = "gate:approval-space-test"; - let scope = personal_turn_scope(); - let run_id = TurnRunId::new(); - let binding_ref = - test_slack_binding_ref(install, scope.agent_id.as_ref().expect("agent").as_str()); - - // First poll → BlockedApproval; second poll → Completed. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![ - scripted_state(TurnStatus::BlockedApproval, Some(gate_ref_str)), - scripted_state(TurnStatus::Completed, None), - ])); - let thread_service = Arc::new(InMemorySessionThreadService::default()); - seed_finalized_assistant_message( - &thread_service, - &scope, - run_id, - "Run complete after approval.", - ) - .await; - - let outbound = Arc::new(InMemoryOutboundStateStore::default()); - seed_personal_preference(&outbound, &scope, binding_ref).await; - let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - // Approval-prompt response: channel D456, ts 1111.2222. egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D456", "1111.2222"), + slack_post_ok_json("D123", "7000.1"), )), ); - // Final-reply response. + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::BlockedApproval, + )); + + // Wire the no-agent binding service directly. + let thread_service = Arc::new(InMemorySessionThreadService::default()); + let services = SlackFinalReplyDeliveryServices { + binding_service: Arc::new(NoAgentConversationBindingService), + thread_service, + turn_coordinator: coordinator, + outbound_store: outbound.clone(), + route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), + communication_preferences: outbound, + adapter: test_adapter(install), + egress: egress.clone(), + delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), + auth_challenges: None, + }; + let settings = SlackFinalReplyDeliverySettings { + poll_interval: std::time::Duration::ZERO, + max_wait: std::time::Duration::from_millis(1), + max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), + max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + }; + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); + + observer.observe_workflow_ack(env, ack).await; + tokio::task::yield_now().await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage even when agent_id is missing" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("still working on a previous message"), + "hint must fall back to generic copy when scope derivation fails, got: {body}" + ); + } + + /// Run-state lookup returns `Err` → generic copy posted. + /// + /// `deferred_busy_hint_from_run_state` swallows `TurnError` from + /// `get_run_state` and degrades to `SLACK_DEFERRED_BUSY_GENERIC_MESSAGE`. + #[tokio::test] + async fn deferred_busy_run_state_lookup_error_posts_generic_hint() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D456", "3333.4444"), + slack_post_ok_json("D123", "7000.2"), )), ); - let delivery_store = Arc::new(InMemoryTriggeredRunDeliveryStore::default()); - let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); - let services = make_services( - coordinator, + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + + use ironclaw_product_workflow::FakeConversationBindingService; + let thread_service = Arc::new(InMemorySessionThreadService::default()); + let services = SlackFinalReplyDeliveryServices { + binding_service: Arc::new(FakeConversationBindingService::new()), thread_service, - egress.clone(), - outbound, - install, - ); + // ErroringTurnCoordinator: get_run_state always returns Err. + turn_coordinator: Arc::new(ErroringTurnCoordinator), + outbound_store: outbound.clone(), + route_store: Arc::new(InMemoryDeliveredGateRouteStore::default()), + communication_preferences: outbound, + adapter: test_adapter(install), + egress: egress.clone(), + delivery_sink: Arc::new(FakeOutboundDeliverySink::default()), + auth_challenges: None, + }; let settings = SlackFinalReplyDeliverySettings { poll_interval: std::time::Duration::ZERO, - max_wait: std::time::Duration::from_secs(5), + max_wait: std::time::Duration::from_millis(1), max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), max_pending_deliveries: NonZeroUsize::new(8).unwrap(), }; - let driver = TriggeredRunDeliveryDriver::with_settings( - services, - settings, - delivery_store.clone(), - route_store.clone(), - scope.agent_id.clone().expect("test scope has agent"), + let observer = SlackFinalReplyDeliveryObserver::with_settings(services, settings); + + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "expected exactly one chat.postMessage even when run-state lookup fails" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("still working on a previous message"), + "hint must fall back to generic copy when run-state lookup fails, got: {body}" ); + } - let fire = minimal_trigger_fire(None); - driver - .on_trigger_submitted(fire, run_id, scope.clone()) - .await; - wait_for_delivery_record(&delivery_store, run_id).await; + /// Slack post failure → no panic, ack path unaffected. + /// + /// The post is best-effort; a transport error must be swallowed with debug! + /// and the observer must return normally. + #[tokio::test] + async fn deferred_busy_post_failure_is_best_effort() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + // Program a transport failure for the hint post. + egress.program_response("slack.com", Err(ProtocolHttpEgressError::Timeout)); - let creator = ironclaw_host_api::UserId::new("creator-user").expect("user id"); - let route = route_store - .load_delivered_gate_route(&scope.tenant_id, &creator, gate_ref_str) - .await - .expect("load route") - .expect("gate route was recorded"); + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::BlockedApproval, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); - // The binding ref encodes space = "T123", so the recorded route must - // contain a ref that fingerprint-matches an inbound-style ref with - // space_id = Some("T123"), conversation_id = "D456", thread_id = "1111.2222". - let expected_inbound_ref = ironclaw_conversations::ExternalConversationRef::new( - Some("T123"), - "D456", - Some("1111.2222"), - None, - ) - .expect("expected inbound ref"); - let expected_fingerprint = expected_inbound_ref.conversation_fingerprint(); + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); - assert!( - route - .delivered_conversation_fingerprints - .iter() - .any(|fingerprint| fingerprint == &expected_fingerprint), - "recorded route must include space_id=T123 fingerprint; fingerprints={:?}", - route.delivered_conversation_fingerprints, + // Must not panic regardless of the egress failure. + observer.observe_workflow_ack(env, ack).await; + + // The call was recorded even though the programmed result was an error. + let calls = egress.calls(); + assert!( + calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "egress must have been called even when the hint post fails (best-effort)" ); } - /// Test B: triggered auth (BlockedAuth) delivery records a gate-route keyed - /// by the auth gate_ref. This is the Bug-2 regression: previously - /// `gate_ref_for_routing` was `None` for BlockedAuth so no route was - /// recorded, causing a `MissingGate` when the user replied "approve". + /// Two DeferredBusy acks with the same conversation + same active_run_id + /// → exactly one post (throttle suppresses the second). #[tokio::test] - async fn triggered_auth_delivery_records_gate_route() { + async fn deferred_busy_same_conversation_same_run_id_posts_once() { let install = "test-install"; - let gate_ref_str = "gate:auth-route-regression"; - let scope = personal_turn_scope(); - let run_id = TurnRunId::new(); - let binding_ref = - test_slack_binding_ref(install, scope.agent_id.as_ref().expect("agent").as_str()); - - // First poll → BlockedAuth with gate_ref; second poll → Completed. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_states(vec![ - scripted_state(TurnStatus::BlockedAuth, Some(gate_ref_str)), - scripted_state(TurnStatus::Completed, None), - ])); - let thread_service = Arc::new(InMemorySessionThreadService::default()); - seed_finalized_assistant_message( - &thread_service, - &scope, - run_id, - "Run complete after auth.", - ) - .await; - - let outbound = Arc::new(InMemoryOutboundStateStore::default()); - seed_personal_preference(&outbound, &scope, binding_ref).await; - let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); egress.allow_credential_handle("slack_bot_token"); - // Auth-prompt delivery response. egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D456", "9999.1111"), + slack_post_ok_json("D123", "9001.1"), )), ); - // Final-reply response (after Completed). + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let run_id = TurnRunId::new(); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::BlockedApproval, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + + let make_ack = || ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:deferred-same").expect("ref"), + active_run_id: run_id, + }; + + // First delivery: posts. + observer + .observe_workflow_ack(envelope(user_message_payload()), make_ack()) + .await; + // Second delivery: same (conversation, run_id) pair → throttled, no second post. + observer + .observe_workflow_ack(envelope(user_message_payload()), make_ack()) + .await; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 1, + "throttle must suppress the second hint for the same (conversation, run_id)" + ); + } + + /// Two DeferredBusy acks with the same conversation but different active_run_ids + /// → two posts (distinct throttle keys). + #[tokio::test] + async fn deferred_busy_same_conversation_different_run_id_posts_twice() { + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - slack_post_ok_json("D456", "9999.2222"), + slack_post_ok_json("D123", "9002.1"), )), ); - // Auth message is deleted after final; need a delete response. egress.program_response( "slack.com", Ok(EgressResponse::new( 200, - serde_json::json!({"ok": true}).to_string().into_bytes(), + slack_post_ok_json("D123", "9002.2"), )), ); - let delivery_store = Arc::new(InMemoryTriggeredRunDeliveryStore::default()); - let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); - let services = make_services( - coordinator, - thread_service, - egress.clone(), - outbound, - install, - ); - let settings = SlackFinalReplyDeliverySettings { - poll_interval: std::time::Duration::ZERO, - max_wait: std::time::Duration::from_secs(5), - max_concurrent_deliveries: NonZeroUsize::new(1).unwrap(), - max_pending_deliveries: NonZeroUsize::new(8).unwrap(), + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::BlockedApproval, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + + let make_ack = |run_id: TurnRunId| ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:deferred-diff").expect("ref"), + active_run_id: run_id, }; - let driver = TriggeredRunDeliveryDriver::with_settings( - services, - settings, - delivery_store.clone(), - route_store.clone(), - scope.agent_id.clone().expect("test scope has agent"), - ); - let fire = minimal_trigger_fire(None); - driver - .on_trigger_submitted(fire, run_id, scope.clone()) + // First delivery: run_id A. + observer + .observe_workflow_ack(envelope(user_message_payload()), make_ack(TurnRunId::new())) + .await; + // Second delivery: run_id B (different key) → separate hint posted. + observer + .observe_workflow_ack(envelope(user_message_payload()), make_ack(TurnRunId::new())) .await; - wait_for_delivery_record(&delivery_store, run_id).await; - - let creator = ironclaw_host_api::UserId::new("creator-user").expect("user id"); - // A gate route for the auth gate_ref must have been recorded. - let route = route_store - .load_delivered_gate_route(&scope.tenant_id, &creator, gate_ref_str) - .await - .expect("load route") - .expect("auth gate route must be recorded (Bug 2 regression)"); + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); assert_eq!( - route.gate_ref, gate_ref_str, - "recorded route must reference the auth gate_ref" + post_calls.len(), + 2, + "different run_ids must produce separate hints even for the same conversation" ); - assert_eq!(route.run_id, run_id); } - /// Test C: the `record_gate_route_if_needed` helper, called from the - /// observer path, stores a posted-message ref whose fingerprint matches an - /// inbound-style ref that carries the envelope's team id (space_id). - /// - /// This tests the observer call site directly (without driving the full - /// observer loop) to verify that `envelope_space_id` is extracted and - /// passed through. + /// deferred_busy_uses_ack_active_run_id_and_binding_scope_for_state_lookup: + /// the GetRunStateRequest forwarded to the coordinator must carry the + /// active_run_id from the DeferredBusy ack and the TurnScope derived from the + /// conversation binding. #[tokio::test] - async fn observer_approval_route_ref_carries_envelope_space_id() { - let tenant_id = ironclaw_host_api::TenantId::new("test-tenant").expect("tenant"); - let user_id = ironclaw_host_api::UserId::new("user-obs").expect("user"); - let run_id = TurnRunId::new(); - let gate_ref_str = "gate:observer-space-test"; - let agent = ironclaw_host_api::AgentId::new("obs-agent").expect("agent"); - let thread = ironclaw_host_api::ThreadId::new("obs-thread").expect("thread"); - let scope = TurnScope::new_with_owner(tenant_id.clone(), Some(agent), None, thread, None); - - // Simulate a posted message: channel D789, ts 5555.6666. - let posted = vec![PostedSlackMessage { - channel: "D789".to_string(), - ts: "5555.6666".to_string(), - }]; - - // Envelope conv ref carries space_id = "T999" (the team id). - let envelope_conv_ref = - ExternalConversationRef::new(Some("T999"), "D789", Some("5555.6666"), None) - .expect("envelope conv ref"); + async fn deferred_busy_uses_ack_active_run_id_and_binding_scope_for_state_lookup() { + use std::sync::Mutex as StdMutex; - let route_store = Arc::new(InMemoryDeliveredGateRouteStore::default()); - - // Derive space_id from the envelope ref — mirrors the observer call site. - let envelope_space_id = conversations_ref_from_product_ref(&envelope_conv_ref) - .ok() - .and_then(|r| r.space_id().map(str::to_string)); - assert_eq!( - envelope_space_id.as_deref(), - Some("T999"), - "space_id must be extracted from envelope ref" + let install = "test-install"; + let egress = Arc::new(FakeProtocolHttpEgress::new(vec!["slack.com".to_string()])); + egress.allow_credential_handle("slack_bot_token"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "9003.1"), + )), ); - record_gate_route_if_needed( - route_store.as_ref(), - run_id, - &tenant_id, - &user_id, - gate_ref_str, - &scope, - &posted, - Some(&envelope_conv_ref), - envelope_space_id.as_deref(), - ) - .await; + let outbound = Arc::new(InMemoryOutboundStateStore::default()); - let route = route_store - .load_delivered_gate_route(&tenant_id, &user_id, gate_ref_str) - .await - .expect("load route") - .expect("route was recorded"); + // Recording coordinator: captures every GetRunStateRequest it receives. + struct RecordingTurnCoordinator { + inner: ScriptedTurnCoordinator, + recorded: StdMutex>, + } + #[async_trait] + impl TurnCoordinator for RecordingTurnCoordinator { + async fn prepare_turn(&self, scope: TurnScope) -> Result { + self.inner.prepare_turn(scope).await + } + async fn submit_turn( + &self, + req: SubmitTurnRequest, + ) -> Result { + self.inner.submit_turn(req).await + } + async fn resume_turn( + &self, + req: ResumeTurnRequest, + ) -> Result { + self.inner.resume_turn(req).await + } + async fn get_run_state( + &self, + request: GetRunStateRequest, + ) -> Result { + self.recorded + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(request.clone()); + self.inner.get_run_state(request).await + } + async fn cancel_run( + &self, + req: ironclaw_turns::CancelRunRequest, + ) -> Result { + self.inner.cancel_run(req).await + } + } - // Must contain a ref with space_id = "T999" matching the inbound fingerprint. - let expected_inbound_ref = ironclaw_conversations::ExternalConversationRef::new( - Some("T999"), - "D789", - Some("5555.6666"), - None, - ) - .expect("inbound ref"); - let expected_fingerprint = expected_inbound_ref.conversation_fingerprint(); + let active_run_id = TurnRunId::new(); + let recording_coordinator = Arc::new(RecordingTurnCoordinator { + inner: ScriptedTurnCoordinator::with_single_status(TurnStatus::BlockedApproval), + recorded: StdMutex::new(Vec::new()), + }); - assert!( - route - .delivered_conversation_fingerprints - .iter() - .any(|fingerprint| fingerprint == &expected_fingerprint), - "recorded route must include space_id=T999 fingerprint; fingerprints={:?}", - route.delivered_conversation_fingerprints, + let observer = make_observer( + recording_coordinator.clone(), + egress.clone(), + outbound, + install, ); - // Also verify that the no-space fallback variant is present (inbound - // events without team_id must still match). - let no_space_ref = ironclaw_conversations::ExternalConversationRef::new( - None, - "D789", - Some("5555.6666"), - None, - ) - .expect("no-space ref"); - let no_space_fingerprint = no_space_ref.conversation_fingerprint(); - assert!( - route - .delivered_conversation_fingerprints - .iter() - .any(|fingerprint| fingerprint == &no_space_fingerprint), - "recorded route must include the no-space fallback fingerprint; fingerprints={:?}", - route.delivered_conversation_fingerprints, - ); + let ack = ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:scope-check").expect("ref"), + active_run_id, + }; + let env = envelope(user_message_payload()); - let channel_root_ref = - ironclaw_conversations::ExternalConversationRef::new(Some("T999"), "D789", None, None) - .expect("channel root ref"); - let channel_root_fingerprint = channel_root_ref.conversation_fingerprint(); - assert!( - route - .delivered_conversation_fingerprints - .iter() - .any(|fingerprint| fingerprint == &channel_root_fingerprint), - "recorded route must include the space-qualified channel-root fingerprint for bare replies; fingerprints={:?}", - route.delivered_conversation_fingerprints, + observer.observe_workflow_ack(env.clone(), ack).await; + + let recorded = recording_coordinator + .recorded + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone(); + assert_eq!(recorded.len(), 1, "expected exactly one GetRunStateRequest"); + assert_eq!( + recorded[0].run_id, active_run_id, + "run_id in GetRunStateRequest must equal the DeferredBusy ack's active_run_id" ); - let no_space_channel_root_ref = - ironclaw_conversations::ExternalConversationRef::new(None, "D789", None, None) - .expect("no-space channel root ref"); - let no_space_channel_root_fingerprint = - no_space_channel_root_ref.conversation_fingerprint(); - assert!( - route - .delivered_conversation_fingerprints - .iter() - .any(|fingerprint| fingerprint == &no_space_channel_root_fingerprint), - "recorded route must include the no-space channel-root fingerprint for bare replies; fingerprints={:?}", - route.delivered_conversation_fingerprints, + // Derive the expected scope from the same binding the observer resolves + // (FakeConversationBindingService is deterministic), through the same + // production helpers, and require an exact match. + let binding = ironclaw_product_workflow::FakeConversationBindingService::new() + .lookup_binding(ResolveBindingRequest::from_envelope(&env)) + .await + .expect("fake binding service resolves test envelope"); + let thread_scope = + thread_scope_from_binding(&binding).expect("test binding derives thread scope"); + let expected_scope = turn_scope_from_thread_scope(&binding, &thread_scope) + .expect("test binding derives turn scope"); + assert_eq!( + recorded[0].scope, expected_scope, + "GetRunStateRequest scope must be derived from the authorized binding" ); } }