From 700d647896ed4246c9c94bffa1760a9344565669 Mon Sep 17 00:00:00 2001 From: Henry Park Date: Thu, 11 Jun 2026 23:06:49 -0700 Subject: [PATCH 1/5] feat(reborn): post Slack feedback when a message is deferred behind a pending gate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When a user sends a Slack message while a run is blocked on a pending approval/auth gate, the inbound ack is DeferredBusy and the delivery observer silently returned — the user got no indication why their message was ignored. Add a best-effort one-shot hint posted before acquiring the delivery semaphore (same pattern as the A2 rejection-hint block). The post only fires for UserMessage payloads — resolution/control payloads stay silent. Duplicate{DeferredBusy} unwraps to the inner ack so Slack transport retries also post the hint exactly once. Co-Authored-By: Claude Fable 5 --- .../src/slack_delivery.rs | 167 ++++++++++++++++++ 1 file changed, 167 insertions(+) diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 1004b19d072..9cb8acbdda3 100644 --- a/crates/ironclaw_reborn_composition/src/slack_delivery.rs +++ b/crates/ironclaw_reborn_composition/src/slack_delivery.rs @@ -61,6 +61,7 @@ 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."; +const SLACK_DEFERRED_BUSY_MESSAGE: &str = "I'm waiting on a pending approval before I can take new messages — reply `approve` or `deny` (or `approve gate:`) to resume."; #[derive(Debug, Clone, PartialEq, Eq)] struct BlockedActionableMarker { @@ -946,6 +947,26 @@ 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 hint so the + // user knows to approve/deny rather than being left in silence. Same + // best-effort semantics as A2: post failure → debug! only, no recursion. + if let Some(hint) = deferred_busy_hint_for_user_message(&envelope, &ack) { + 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 +1172,31 @@ fn rejection_hint_for_resolution( Some(hint) } +/// Returns the user-facing hint to post when an inbound user message is deferred +/// because a run for the same conversation is blocked on a pending gate. +/// +/// Returns `None` for non-UserMessage payloads (resolution/control payloads must +/// stay silent) and for `Duplicate` acks (transport retries must not double-post). +fn deferred_busy_hint_for_user_message( + envelope: &ProductInboundEnvelope, + ack: &ProductInboundAck, +) -> Option<&'static str> { + // Unwrap one layer of Duplicate so a Slack transport retry of the same + // event that originally produced DeferredBusy does not re-post the hint. + let effective = match ack { + ProductInboundAck::Duplicate { prior } => prior.as_ref(), + other => other, + }; + if !matches!(effective, ProductInboundAck::DeferredBusy { .. }) { + return None; + } + // Only reply to user messages — resolution/control/noop payloads must stay silent. + if !matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) { + return None; + } + Some(SLACK_DEFERRED_BUSY_MESSAGE) +} + fn rejection_ack_for_workflow_error(error: &ProductAdapterError) -> Option { match error { ProductAdapterError::WorkflowRejected { @@ -3151,6 +3197,127 @@ mod tests { ); } + // ── 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 → exactly one Slack post containing + /// "waiting on a pending approval". + #[tokio::test] + async fn deferred_busy_ack_with_user_message_posts_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", "2000.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(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!( + body.contains("waiting on a pending approval"), + "deferred-busy hint must mention 'waiting on a pending approval', got: {body}" + ); + } + + /// DeferredBusy ack + non-UserMessage payload → no post (resolution payloads + /// already have their own feedback path and must stay silent here). + #[tokio::test] + 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"); + + 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 = 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 DeferredBusy with non-UserMessage payload" + ); + } + + /// Duplicate { prior: DeferredBusy } + UserMessage → exactly one post + /// (transport retry of the original DeferredBusy event behaves like the inner + /// ack so the user still gets feedback). + #[tokio::test] + async fn duplicate_deferred_busy_with_user_message_posts_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", "2000.2"), + )), + ); + + 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(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, + "Duplicate{{DeferredBusy}} + UserMessage must produce exactly one hint post" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("waiting on a pending approval"), + "hint must mention 'waiting on a pending approval', got: {body}" + ); + } + /// Duplicate { prior: Accepted } → nothing posted (already succeeded). #[tokio::test] async fn duplicate_accepted_ack_posts_nothing() { From c8df68f736704f74b67dddad560c651e49e1d89e Mon Sep 17 00:00:00 2001 From: Henry Park Date: Fri, 12 Jun 2026 01:55:47 -0700 Subject: [PATCH 2/5] fix(reborn): correct DeferredBusy hint duplicate-ack semantics and message voice MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Return None for ALL Duplicate acks in deferred_busy_hint_for_user_message (DeferredBusy is never settled by the idempotency ledger, so Duplicate{DeferredBusy} is unreachable; None-for-all-Duplicate is the safe invariant matching rejection_hint_for_resolution) - Document that each distinct plain DeferredBusy delivery posts a hint (desired feedback for users sending multiple messages while blocked; rare transport retries may double-post — accepted as benign best-effort) - Replace duplicate_deferred_busy_with_user_message_posts_hint test with duplicate_deferred_busy_with_user_message_posts_nothing asserting the new None-for-all-Duplicate behavior - Add two_distinct_deferred_busy_user_messages_post_two_hints test asserting the deliberate per-delivery retry semantics - Change SLACK_DEFERRED_BUSY_MESSAGE to impersonal voice matching sibling constants ("Ironclaw is waiting..." instead of "I'm waiting...") Co-Authored-By: Claude Fable 5 --- .../src/slack_delivery.rs | 96 ++++++++++++++----- 1 file changed, 71 insertions(+), 25 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 9cb8acbdda3..2aef39defcd 100644 --- a/crates/ironclaw_reborn_composition/src/slack_delivery.rs +++ b/crates/ironclaw_reborn_composition/src/slack_delivery.rs @@ -61,7 +61,7 @@ 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."; -const SLACK_DEFERRED_BUSY_MESSAGE: &str = "I'm waiting on a pending approval before I can take new messages — reply `approve` or `deny` (or `approve gate:`) to resume."; +const SLACK_DEFERRED_BUSY_MESSAGE: &str = "Ironclaw is waiting on a pending approval before taking new messages — reply `approve` or `deny` (or `approve gate:`) to resume."; #[derive(Debug, Clone, PartialEq, Eq)] struct BlockedActionableMarker { @@ -1176,24 +1176,29 @@ fn rejection_hint_for_resolution( /// because a run for the same conversation is blocked on a pending gate. /// /// Returns `None` for non-UserMessage payloads (resolution/control payloads must -/// stay silent) and for `Duplicate` acks (transport retries must not double-post). +/// stay silent) and for ALL `Duplicate` acks (the idempotency ledger never settles +/// `DeferredBusy`, so Slack transport retries re-process as a fresh plain +/// `DeferredBusy`, never `Duplicate{DeferredBusy}`; returning `None` for every +/// `Duplicate` is the safe invariant that matches `rejection_hint_for_resolution`). +/// Returns `Some` only for a plain `DeferredBusy` ack + `UserMessage` payload. fn deferred_busy_hint_for_user_message( envelope: &ProductInboundEnvelope, ack: &ProductInboundAck, ) -> Option<&'static str> { - // Unwrap one layer of Duplicate so a Slack transport retry of the same - // event that originally produced DeferredBusy does not re-post the hint. - let effective = match ack { - ProductInboundAck::Duplicate { prior } => prior.as_ref(), - other => other, - }; - if !matches!(effective, ProductInboundAck::DeferredBusy { .. }) { + // All Duplicate acks are suppressed — same pattern as rejection_hint_for_resolution. + if matches!(ack, ProductInboundAck::Duplicate { .. }) { + return None; + } + if !matches!(ack, ProductInboundAck::DeferredBusy { .. }) { return None; } // Only reply to user messages — resolution/control/noop payloads must stay silent. if !matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) { return None; } + // Each delivery of a DeferredBusy ack posts the hint: a user sending multiple + // messages while blocked gets a hint for each one (desired feedback). Rare + // transport retries may double-post — accepted as benign best-effort. Some(SLACK_DEFERRED_BUSY_MESSAGE) } @@ -3273,14 +3278,57 @@ mod tests { ); } - /// Duplicate { prior: DeferredBusy } + UserMessage → exactly one post - /// (transport retry of the original DeferredBusy event behaves like the inner - /// ack so the user still gets feedback). + /// Duplicate { prior: DeferredBusy } + UserMessage → nothing posted. + /// + /// `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 duplicate_deferred_busy_with_user_message_posts_hint() { + 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"); + + 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(user_message_payload()); + let ack = ProductInboundAck::Duplicate { + prior: Box::new(deferred_busy_ack()), + }; + + observer.observe_workflow_ack(env, ack).await; + + let calls = egress.calls(); + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "Duplicate{{DeferredBusy}} must NOT post a hint (None-for-all-Duplicate invariant)" + ); + } + + /// Two distinct plain DeferredBusy + UserMessage envelopes → two posts. + /// + /// Each delivery of a DeferredBusy ack posts the hint: a user sending + /// multiple messages while blocked gets a hint for each one (desired + /// feedback). Rare transport retries may double-post — accepted as benign + /// best-effort. This is a deliberate design choice: feedback density beats + /// silent suppression for the common case. + #[tokio::test] + 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"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "2000.1"), + )), + ); egress.program_response( "slack.com", Ok(EgressResponse::new( @@ -3294,12 +3342,15 @@ mod tests { TurnStatus::Running, )); let observer = make_observer(coordinator, egress.clone(), outbound, install); - let env = envelope(user_message_payload()); - let ack = ProductInboundAck::Duplicate { - prior: Box::new(deferred_busy_ack()), - }; - observer.observe_workflow_ack(env, ack).await; + // First distinct user message while blocked. + observer + .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) + .await; + // Second distinct user message while blocked (different event, new ack). + observer + .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) + .await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3308,13 +3359,8 @@ mod tests { .collect(); assert_eq!( post_calls.len(), - 1, - "Duplicate{{DeferredBusy}} + UserMessage must produce exactly one hint post" - ); - let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); - assert!( - body.contains("waiting on a pending approval"), - "hint must mention 'waiting on a pending approval', got: {body}" + 2, + "two distinct DeferredBusy + UserMessage deliveries must each post a hint" ); } From b6a23491fca5ace0c8c6f59bdc3f9a5b2385489f Mon Sep 17 00:00:00 2001 From: Henry Park Date: Fri, 12 Jun 2026 08:03:35 -0700 Subject: [PATCH 3/5] fix(reborn): state-aware deferred-busy hint with authorization and detached post MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Look up the blocking run's TurnStatus via turn_coordinator.get_run_state before choosing copy for the DeferredBusy hint: - BlockedApproval → approval wording (approve / deny) - BlockedAuth → auth wording (complete auth or auth deny) - anything else / lookup failure → generic queued-task wording Add binding authorization check (mirrors post_rejection_hint_if_authorized) so the hint is silently skipped for unauthorized conversations. Spawn the Slack post in a detached task so a burst of deferred messages does not pin post-ack workers on external I/O. Failure handling is debug!-only inside the spawned task. Replace SLACK_DEFERRED_BUSY_MESSAGE with three typed constants. Replace deferred_busy_hint_for_user_message with is_deferred_busy_user_message predicate and new async deferred_busy_hint_from_run_state helper. Tests: update 2 existing DeferredBusy tests to use BlockedApproval state and yield_now drain; add 3 new tests: BlockedAuth copy, Running/generic copy, and unauthorized-binding suppression. Co-Authored-By: Claude Fable 5 --- .../src/slack_delivery.rs | 313 +++++++++++++++--- 1 file changed, 275 insertions(+), 38 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 2aef39defcd..22aa8678574 100644 --- a/crates/ironclaw_reborn_composition/src/slack_delivery.rs +++ b/crates/ironclaw_reborn_composition/src/slack_delivery.rs @@ -61,7 +61,12 @@ 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."; -const SLACK_DEFERRED_BUSY_MESSAGE: &str = "Ironclaw is waiting on a pending approval before taking new messages — reply `approve` or `deny` (or `approve gate:`) to resume."; +/// Posted when the blocking run is `BlockedApproval`. +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 when the blocking run is `BlockedAuth`. +const SLACK_DEFERRED_BUSY_AUTH_MESSAGE: &str = "Ironclaw is waiting on an authentication step before taking new messages — complete the authentication prompt (or reply `auth deny ` to decline)."; +/// Posted for any other non-terminal blocking state, or when the state lookup fails. +const SLACK_DEFERRED_BUSY_GENERIC_MESSAGE: &str = "Ironclaw is still working on a previous message — this one is queued and will run when the current task finishes."; #[derive(Debug, Clone, PartialEq, Eq)] struct BlockedActionableMarker { @@ -948,23 +953,59 @@ 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 hint so the - // user knows to approve/deny rather than being left in silence. Same - // best-effort semantics as A2: post failure → debug! only, no recursion. - if let Some(hint) = deferred_busy_hint_for_user_message(&envelope, &ack) { - if let Err(post_err) = post_slack_message( - self.services.egress.as_ref(), - envelope.external_conversation_ref(), - hint, - ) - .await + // 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`. + // + // Backpressure: the hint post is spawned detached so a burst of deferred + // messages does not pin post-ack workers on Slack I/O. + if is_deferred_busy_user_message(&envelope, &ack) { + let active_run_id = match &ack { + ProductInboundAck::DeferredBusy { active_run_id, .. } => *active_run_id, + _ => unreachable!("is_deferred_busy_user_message ensures DeferredBusy"), + }; + // 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 { - tracing::debug!( - target = "ironclaw::reborn::slack_delivery", - error = %post_err, - "failed to post deferred-busy hint to Slack (best-effort)" - ); - } + 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; + let egress = Arc::clone(&self.services.egress); + let conversation = envelope.external_conversation_ref().clone(); + tokio::spawn(async move { + if let Err(post_err) = + post_slack_message(egress.as_ref(), &conversation, 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 { @@ -1172,34 +1213,77 @@ fn rejection_hint_for_resolution( Some(hint) } -/// Returns the user-facing hint to post when an inbound user message is deferred -/// because a run for the same conversation is blocked on a pending gate. +/// Returns `true` when the ack + payload combination should trigger the deferred-busy +/// hint flow: a plain `DeferredBusy` ack on a `UserMessage` payload. /// -/// Returns `None` for non-UserMessage payloads (resolution/control payloads must -/// stay silent) and for ALL `Duplicate` acks (the idempotency ledger never settles +/// Returns `false` for `Duplicate` acks (the idempotency ledger never settles /// `DeferredBusy`, so Slack transport retries re-process as a fresh plain -/// `DeferredBusy`, never `Duplicate{DeferredBusy}`; returning `None` for every -/// `Duplicate` is the safe invariant that matches `rejection_hint_for_resolution`). -/// Returns `Some` only for a plain `DeferredBusy` ack + `UserMessage` payload. -fn deferred_busy_hint_for_user_message( +/// `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 is_deferred_busy_user_message( envelope: &ProductInboundEnvelope, ack: &ProductInboundAck, -) -> Option<&'static str> { +) -> bool { // All Duplicate acks are suppressed — same pattern as rejection_hint_for_resolution. if matches!(ack, ProductInboundAck::Duplicate { .. }) { - return None; + return false; } if !matches!(ack, ProductInboundAck::DeferredBusy { .. }) { - return None; + return false; } // Only reply to user messages — resolution/control/noop payloads must stay silent. - if !matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) { - return None; + matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) +} + +/// Looks up the blocking run's state and returns the appropriate deferred-busy +/// hint copy. +/// +/// - `BlockedApproval` → approval wording +/// - `BlockedAuth` → auth wording +/// - 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, +) -> &'static str { + 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; + } + }; + match coordinator + .get_run_state(GetRunStateRequest { + scope, + run_id: active_run_id, + }) + .await + { + Ok(state) => match state.status { + TurnStatus::BlockedApproval => SLACK_DEFERRED_BUSY_APPROVAL_MESSAGE, + TurnStatus::BlockedAuth => SLACK_DEFERRED_BUSY_AUTH_MESSAGE, + _ => SLACK_DEFERRED_BUSY_GENERIC_MESSAGE, + }, + 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 + } } - // Each delivery of a DeferredBusy ack posts the hint: a user sending multiple - // messages while blocked gets a hint for each one (desired feedback). Rare - // transport retries may double-post — accepted as benign best-effort. - Some(SLACK_DEFERRED_BUSY_MESSAGE) } fn rejection_ack_for_workflow_error(error: &ProductAdapterError) -> Option { @@ -3211,8 +3295,11 @@ mod tests { } } - /// DeferredBusy ack + UserMessage payload → exactly one Slack post containing - /// "waiting on a pending approval". + /// DeferredBusy ack + UserMessage payload + BlockedApproval state → exactly one + /// Slack post containing approval wording. + /// + /// The hint post is spawned detached; `yield_now` lets it execute before we + /// inspect the egress capture. #[tokio::test] async fn deferred_busy_ack_with_user_message_posts_hint() { let install = "test-install"; @@ -3227,14 +3314,17 @@ mod tests { ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // BlockedApproval → state-aware lookup returns approval wording. let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, + TurnStatus::BlockedApproval, )); let observer = make_observer(coordinator, egress.clone(), outbound, install); let env = envelope(user_message_payload()); let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; + // Drain the spawned hint post. + tokio::task::yield_now().await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3317,6 +3407,8 @@ mod tests { /// feedback). Rare transport retries may double-post — accepted as benign /// best-effort. This is a deliberate design choice: feedback density beats /// silent suppression for the common case. + /// + /// Hint posts are spawned detached; `yield_now` drains both before asserting. #[tokio::test] async fn two_distinct_deferred_busy_user_messages_post_two_hints() { let install = "test-install"; @@ -3338,8 +3430,9 @@ mod tests { ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // BlockedApproval so the state-aware lookup returns the approval copy for both. let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( - TurnStatus::Running, + TurnStatus::BlockedApproval, )); let observer = make_observer(coordinator, egress.clone(), outbound, install); @@ -3351,6 +3444,8 @@ mod tests { observer .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) .await; + // Drain both spawned hint posts. + tokio::task::yield_now().await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3364,6 +3459,148 @@ mod tests { ); } + /// DeferredBusy + UserMessage + BlockedAuth state → auth copy posted. + #[tokio::test] + 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"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "6000.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // BlockedAuth → state-aware lookup returns auth wording. + let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + TurnStatus::BlockedAuth, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + 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 for DeferredBusy + BlockedAuth" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("authentication step"), + "deferred-busy hint for BlockedAuth must mention 'authentication step', 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}" + ); + } + + /// DeferredBusy + UserMessage + Running state (non-blocked) → generic copy posted. + #[tokio::test] + 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"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + 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(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 for DeferredBusy + Running" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("still working on a previous message"), + "deferred-busy hint for Running state must contain generic copy, got: {body}" + ); + } + + /// DeferredBusy + UserMessage + unauthorized binding → no post, no panic. + /// + /// 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 deferred_busy_unauthorized_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::BlockedApproval, + )); + // 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: Arc::new(TestNoopConversationBindingService), + thread_service, + turn_coordinator: coordinator, + outbound_store: outbound.clone(), + 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; + // No spawn should have been issued (authorization failed synchronously). + // A yield is still safe to add as a hedge. + tokio::task::yield_now().await; + + let calls = egress.calls(); + assert!( + !calls.iter().any(|c| c.path == "/api/chat.postMessage"), + "no chat.postMessage expected when binding authorization fails" + ); + } + /// Duplicate { prior: Accepted } → nothing posted (already succeeded). #[tokio::test] async fn duplicate_accepted_ack_posts_nothing() { From fb7f982dc691d91bc67d3f07493fbf04d4a0453c Mon Sep 17 00:00:00 2001 From: Henry Park Date: Fri, 12 Jun 2026 09:54:57 -0700 Subject: [PATCH 4/5] fix(reborn): bounded hint-posts, honest copy, no unreachable!, three new tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. Unbounded detached posts: add `hint_post_permits` semaphore (MAX_CONCURRENT_HINT_POSTS = 4) to `SlackFinalReplyDeliveryObserver`. `try_acquire_owned()` before spawning; drop hint with debug! on saturation so a burst cannot create unbounded tasks or bypass max_concurrent_deliveries. Permit held for the full duration of the Slack I/O inside the spawn. 2. Generic copy honesty: rewrite SLACK_DEFERRED_BUSY_GENERIC_MESSAGE to "...can't take this one yet — please resend it once the current task finishes." (no false queue promise); add comment to revisit when deferred-drain PR #4812 lands. 3. Helper API: replace `is_deferred_busy_user_message(…) -> bool` + unreachable! extraction with `deferred_busy_user_message_run_id(…) -> Option`; caller uses `if let Some(run_id) = …`. 4. New test: deferred_busy_missing_agent_binding_posts_generic_hint — binding with agent_id=None → scope derivation fails → generic copy. 5. New test: deferred_busy_run_state_lookup_error_posts_generic_hint — ErroringTurnCoordinator always returns Err → generic copy. 6. New test: deferred_busy_post_failure_is_best_effort — egress programmed with Err(Timeout) → no panic, ack path returns normally. Bonus: deferred_busy_semaphore_saturated_drops_hint_without_panic — all permits pre-held → no post, no panic. Co-Authored-By: Claude Fable 5 --- .../src/slack_delivery.rs | 360 +++++++++++++++++- 1 file changed, 342 insertions(+), 18 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 22aa8678574..20df72da308 100644 --- a/crates/ironclaw_reborn_composition/src/slack_delivery.rs +++ b/crates/ironclaw_reborn_composition/src/slack_delivery.rs @@ -52,6 +52,13 @@ use crate::auth_prompt::auth_prompt_view_for_blocked_auth; use crate::slack_outbound_targets::slack_conversation_id_from_reply_target_binding_ref; const MAX_SLACK_RUN_POLL_INTERVAL: Duration = Duration::from_secs(5); +/// Maximum number of concurrently in-flight deferred-busy hint posts. +/// +/// Hint posts are fire-and-forget best-effort; capping to 4 prevents a burst +/// from creating unbounded tasks or overwhelming Slack egress. Bursts beyond +/// this limit drop the hint with a `debug!` log — the user already received +/// hints for earlier messages in the burst. +const MAX_CONCURRENT_HINT_POSTS: usize = 4; const SLACK_RUN_POLL_JITTER_BUCKETS: u32 = 5; const SLACK_API_HOST: &str = "slack.com"; const SLACK_BOT_TOKEN_HANDLE: &str = "slack_bot_token"; @@ -66,7 +73,10 @@ const SLACK_DEFERRED_BUSY_APPROVAL_MESSAGE: &str = "Ironclaw is waiting on a pen /// Posted when the blocking run is `BlockedAuth`. const SLACK_DEFERRED_BUSY_AUTH_MESSAGE: &str = "Ironclaw is waiting on an authentication step before taking new messages — complete the authentication prompt (or reply `auth deny ` to decline)."; /// Posted for any other non-terminal blocking state, or when the state lookup fails. -const SLACK_DEFERRED_BUSY_GENERIC_MESSAGE: &str = "Ironclaw is still working on a previous message — this one is queued and will run when the current task finishes."; +/// +/// 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 { @@ -123,6 +133,9 @@ pub struct SlackFinalReplyDeliveryObserver { services: SlackFinalReplyDeliveryServices, settings: SlackFinalReplyDeliverySettings, delivery_permits: Arc, + /// Bounds concurrent in-flight deferred-busy hint posts. + /// `try_acquire_owned` before spawning; drop hint (debug!) on saturation. + hint_post_permits: Arc, } impl SlackFinalReplyDeliveryObserver { @@ -138,6 +151,7 @@ impl SlackFinalReplyDeliveryObserver { services, settings, delivery_permits: Arc::new(Semaphore::new(settings.max_concurrent_deliveries.get())), + hint_post_permits: Arc::new(Semaphore::new(MAX_CONCURRENT_HINT_POSTS)), } } @@ -961,12 +975,10 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { // same guard used by `post_rejection_hint_if_authorized`. // // Backpressure: the hint post is spawned detached so a burst of deferred - // messages does not pin post-ack workers on Slack I/O. - if is_deferred_busy_user_message(&envelope, &ack) { - let active_run_id = match &ack { - ProductInboundAck::DeferredBusy { active_run_id, .. } => *active_run_id, - _ => unreachable!("is_deferred_busy_user_message ensures DeferredBusy"), - }; + // messages does not pin post-ack workers on Slack I/O. A small semaphore + // (MAX_CONCURRENT_HINT_POSTS) caps concurrent in-flight spawns so a burst + // does not create unbounded tasks or bypass max_concurrent_deliveries. + if let Some(active_run_id) = deferred_busy_user_message_run_id(&envelope, &ack) { // Authorization: verify the conversation has a valid binding before // posting. Skip silently (debug!) when unauthorized. let binding = match self @@ -985,6 +997,15 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { return; } }; + // Backpressure: try_acquire_owned is non-blocking so the ack path is + // never stalled waiting for a permit. + let Ok(permit) = Arc::clone(&self.hint_post_permits).try_acquire_owned() else { + tracing::debug!( + target = "ironclaw::reborn::slack_delivery", + "deferred-busy hint dropped: hint-post semaphore saturated (best-effort)" + ); + 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( @@ -996,6 +1017,9 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { let egress = Arc::clone(&self.services.egress); let conversation = envelope.external_conversation_ref().clone(); tokio::spawn(async move { + // Hold the permit until the post completes so the in-flight count + // is accurate for the full duration of the Slack I/O. + let _permit = permit; if let Err(post_err) = post_slack_message(egress.as_ref(), &conversation, hint).await { @@ -1213,27 +1237,30 @@ fn rejection_hint_for_resolution( Some(hint) } -/// Returns `true` when the ack + payload combination should trigger the deferred-busy -/// hint flow: a plain `DeferredBusy` ack on a `UserMessage` payload. +/// 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 `false` for `Duplicate` acks (the idempotency ledger never settles +/// 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 is_deferred_busy_user_message( +fn deferred_busy_user_message_run_id( envelope: &ProductInboundEnvelope, ack: &ProductInboundAck, -) -> bool { +) -> Option { // All Duplicate acks are suppressed — same pattern as rejection_hint_for_resolution. if matches!(ack, ProductInboundAck::Duplicate { .. }) { - return false; - } - if !matches!(ack, ProductInboundAck::DeferredBusy { .. }) { - return false; + return None; } // Only reply to user messages — resolution/control/noop payloads must stay silent. - matches!(envelope.payload(), ProductInboundPayload::UserMessage(_)) + 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 @@ -2225,7 +2252,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, @@ -3572,6 +3600,7 @@ mod tests { 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(), @@ -4680,4 +4709,299 @@ mod tests { route.delivered_conversation_fingerprints, ); } + + // ── Extra DeferredBusy coverage tests ───────────────────────────────────── + + /// 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; + + #[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 + } + } + + /// A `TurnCoordinator` double whose `get_run_state` always returns `Err`. + struct ErroringTurnCoordinator; + + #[async_trait] + impl TurnCoordinator for ErroringTurnCoordinator { + async fn prepare_turn(&self, _scope: TurnScope) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } + + async fn submit_turn( + &self, + _request: SubmitTurnRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } + + async fn resume_turn( + &self, + _request: ResumeTurnRequest, + ) -> Result { + Err(TurnError::Unavailable { + reason: "ErroringTurnCoordinator".to_string(), + }) + } + + 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. + /// + /// `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 deferred_busy_missing_agent_binding_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("D123", "7000.1"), + )), + ); + + 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("D123", "7000.2"), + )), + ); + + 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, + // 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_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 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}" + ); + } + + /// Slack post inside the spawned hint task fails → no panic, ack path unaffected. + /// + /// The hint-post semaphore must be released and the outer call must return + /// normally. The `FakeProtocolHttpEgress` returns an error response when we + /// program `Err(...)` into the queue. + #[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 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 env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); + + // Must not panic regardless of the egress failure. + observer.observe_workflow_ack(env, ack).await; + tokio::task::yield_now().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)" + ); + } + + /// Saturated hint-post semaphore → hint dropped, no post, no panic. + /// + /// Acquire all MAX_CONCURRENT_HINT_POSTS permits from the observer's semaphore + /// before the ack arrives. The observer must silently drop the hint instead of + /// spawning an unbounded task. + #[tokio::test] + async fn deferred_busy_semaphore_saturated_drops_hint_without_panic() { + 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::BlockedApproval, + )); + let observer = make_observer(coordinator, egress.clone(), outbound, install); + + // Saturate the hint-post semaphore: acquire all MAX_CONCURRENT_HINT_POSTS permits. + let mut permits = Vec::new(); + for _ in 0..MAX_CONCURRENT_HINT_POSTS { + let p = Arc::clone(&observer.hint_post_permits) + .try_acquire_owned() + .expect("permit must be available before saturation"); + permits.push(p); + } + + let env = envelope(user_message_payload()); + let ack = deferred_busy_ack(); + + // Must not panic, must not post. + 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 when hint-post semaphore is saturated" + ); + + // Release permits explicitly so the test is self-contained. + drop(permits); + } } From 33b0c7aed68f656add644c229c382e75f5e50dc2 Mon Sep 17 00:00:00 2001 From: Henry Park Date: Fri, 12 Jun 2026 13:50:02 -0700 Subject: [PATCH 5/5] fix(reborn): address five review findings on deferred-busy hint feature (#4811) - Fix 1: remove detached tokio::spawn + hint_post_permits semaphore; await post_slack_message inline in observe_workflow_ack; update block comment to explain why inline is safe (protocol ACK already returned, runner admission permit bounds the whole post-ACK task per runner_immediate_ack.rs contract); delete the now-meaningless semaphore-saturation test - Fix 2: add per-observer in-memory throttle (hint_seen: Mutex<...>) bounded at HINT_SEEN_CAP=256 entries with FIFO eviction; check before binding lookup to avoid coordinator round-trip on repeats; add tests deferred_busy_same_conversation_same_run_id_posts_once and deferred_busy_same_conversation_different_run_id_posts_twice - Fix 3: change deferred_busy_hint_from_run_state to return String; format concrete gate ref into BlockedApproval and BlockedAuth hints (`approve {ref}` / `auth deny {ref}`); fall back to static copy when gate_ref is None; update existing BlockedApproval and BlockedAuth tests to use scripted gate refs and assert concrete text; add deferred_busy_blocked_approval_no_gate_ref_posts_fallback_hint for None path - Fix 4: add deferred_busy_uses_ack_active_run_id_and_binding_scope_for_state_lookup using RecordingTurnCoordinator to assert the forwarded GetRunStateRequest carries the ack's active_run_id and a non-empty tenant scope - Fix 5: add arch-exempt: large_file annotation near top of file referencing decomposition issue #4818 Co-Authored-By: Claude Fable 5 --- .../src/slack_delivery.rs | 470 ++++++++++++++---- 1 file changed, 360 insertions(+), 110 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/slack_delivery.rs b/crates/ironclaw_reborn_composition/src/slack_delivery.rs index 20df72da308..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; @@ -52,13 +54,6 @@ use crate::auth_prompt::auth_prompt_view_for_blocked_auth; use crate::slack_outbound_targets::slack_conversation_id_from_reply_target_binding_ref; const MAX_SLACK_RUN_POLL_INTERVAL: Duration = Duration::from_secs(5); -/// Maximum number of concurrently in-flight deferred-busy hint posts. -/// -/// Hint posts are fire-and-forget best-effort; capping to 4 prevents a burst -/// from creating unbounded tasks or overwhelming Slack egress. Bursts beyond -/// this limit drop the hint with a `debug!` log — the user already received -/// hints for earlier messages in the burst. -const MAX_CONCURRENT_HINT_POSTS: usize = 4; const SLACK_RUN_POLL_JITTER_BUCKETS: u32 = 5; const SLACK_API_HOST: &str = "slack.com"; const SLACK_BOT_TOKEN_HANDLE: &str = "slack_bot_token"; @@ -68,10 +63,8 @@ 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`. +/// 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 when the blocking run is `BlockedAuth`. -const SLACK_DEFERRED_BUSY_AUTH_MESSAGE: &str = "Ironclaw is waiting on an authentication step before taking new messages — complete the authentication prompt (or reply `auth deny ` to decline)."; /// 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). @@ -129,13 +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, - /// Bounds concurrent in-flight deferred-busy hint posts. - /// `try_acquire_owned` before spawning; drop hint (debug!) on saturation. - hint_post_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 { @@ -151,7 +154,7 @@ impl SlackFinalReplyDeliveryObserver { services, settings, delivery_permits: Arc::new(Semaphore::new(settings.max_concurrent_deliveries.get())), - hint_post_permits: Arc::new(Semaphore::new(MAX_CONCURRENT_HINT_POSTS)), + hint_seen: Mutex::new((VecDeque::new(), HashSet::new())), } } @@ -974,11 +977,42 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { // Authorization: only post if the binding lookup succeeds, matching the // same guard used by `post_rejection_hint_if_authorized`. // - // Backpressure: the hint post is spawned detached so a burst of deferred - // messages does not pin post-ack workers on Slack I/O. A small semaphore - // (MAX_CONCURRENT_HINT_POSTS) caps concurrent in-flight spawns so a burst - // does not create unbounded tasks or bypass max_concurrent_deliveries. + // 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 @@ -997,15 +1031,6 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { return; } }; - // Backpressure: try_acquire_owned is non-blocking so the ack path is - // never stalled waiting for a permit. - let Ok(permit) = Arc::clone(&self.hint_post_permits).try_acquire_owned() else { - tracing::debug!( - target = "ironclaw::reborn::slack_delivery", - "deferred-busy hint dropped: hint-post semaphore saturated (best-effort)" - ); - 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( @@ -1014,22 +1039,19 @@ impl ImmediateAckWorkflowObserver for SlackFinalReplyDeliveryObserver { active_run_id, ) .await; - let egress = Arc::clone(&self.services.egress); - let conversation = envelope.external_conversation_ref().clone(); - tokio::spawn(async move { - // Hold the permit until the post completes so the in-flight count - // is accurate for the full duration of the Slack I/O. - let _permit = permit; - if let Err(post_err) = - post_slack_message(egress.as_ref(), &conversation, hint).await - { - tracing::debug!( - target = "ironclaw::reborn::slack_delivery", - error = %post_err, - "failed to post deferred-busy hint to Slack (best-effort)" - ); - } - }); + 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 { @@ -1266,16 +1288,18 @@ fn deferred_busy_user_message_run_id( /// Looks up the blocking run's state and returns the appropriate deferred-busy /// hint copy. /// -/// - `BlockedApproval` → approval wording -/// - `BlockedAuth` → auth wording -/// - anything else / lookup failure → generic wording +/// - `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, -) -> &'static str { +) -> String { let scope = match (|| -> Result { let thread_scope = thread_scope_from_binding(binding)?; turn_scope_from_thread_scope(binding, &thread_scope) @@ -1287,7 +1311,7 @@ async fn deferred_busy_hint_from_run_state( error = %err, "deferred-busy scope derivation failed; using generic copy" ); - return SLACK_DEFERRED_BUSY_GENERIC_MESSAGE; + return SLACK_DEFERRED_BUSY_GENERIC_MESSAGE.to_string(); } }; match coordinator @@ -1298,9 +1322,25 @@ async fn deferred_busy_hint_from_run_state( .await { Ok(state) => match state.status { - TurnStatus::BlockedApproval => SLACK_DEFERRED_BUSY_APPROVAL_MESSAGE, - TurnStatus::BlockedAuth => SLACK_DEFERRED_BUSY_AUTH_MESSAGE, - _ => SLACK_DEFERRED_BUSY_GENERIC_MESSAGE, + 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!( @@ -1308,7 +1348,7 @@ async fn deferred_busy_hint_from_run_state( error = %err, "deferred-busy run-state lookup failed; using generic copy" ); - SLACK_DEFERRED_BUSY_GENERIC_MESSAGE + SLACK_DEFERRED_BUSY_GENERIC_MESSAGE.to_string() } } } @@ -3323,11 +3363,10 @@ mod tests { } } - /// DeferredBusy ack + UserMessage payload + BlockedApproval state → exactly one - /// Slack post containing approval wording. + /// DeferredBusy ack + UserMessage payload + BlockedApproval state with gate_ref → + /// exactly one Slack post containing the concrete `approve gate:` command. /// - /// The hint post is spawned detached; `yield_now` lets it execute before we - /// inspect the egress capture. + /// The hint post is awaited inline; no yield needed before inspecting the egress capture. #[tokio::test] async fn deferred_busy_ack_with_user_message_posts_hint() { let install = "test-install"; @@ -3342,17 +3381,17 @@ mod tests { ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - // BlockedApproval → state-aware lookup returns approval wording. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + // 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(user_message_payload()); let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; - // Drain the spawned hint post. - tokio::task::yield_now().await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3369,6 +3408,10 @@ mod tests { 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}" + ); } /// DeferredBusy ack + non-UserMessage payload → no post (resolution payloads @@ -3428,15 +3471,11 @@ mod tests { ); } - /// Two distinct plain DeferredBusy + UserMessage envelopes → two posts. - /// - /// Each delivery of a DeferredBusy ack posts the hint: a user sending - /// multiple messages while blocked gets a hint for each one (desired - /// feedback). Rare transport retries may double-post — accepted as benign - /// best-effort. This is a deliberate design choice: feedback density beats - /// silent suppression for the common case. + /// Two distinct plain DeferredBusy + UserMessage envelopes with different active_run_ids + /// → two posts (throttle is per (conversation, run_id) pair). /// - /// Hint posts are spawned detached; `yield_now` drains both before asserting. + /// 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 two_distinct_deferred_busy_user_messages_post_two_hints() { let install = "test-install"; @@ -3464,16 +3503,14 @@ mod tests { )); let observer = make_observer(coordinator, egress.clone(), outbound, install); - // First distinct user message while blocked. + // 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 event, new ack). + // Second distinct user message while blocked (different active_run_id → different key). observer .observe_workflow_ack(envelope(user_message_payload()), deferred_busy_ack()) .await; - // Drain both spawned hint posts. - tokio::task::yield_now().await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3487,7 +3524,8 @@ mod tests { ); } - /// DeferredBusy + UserMessage + BlockedAuth state → auth copy posted. + /// DeferredBusy + UserMessage + BlockedAuth state with gate_ref → auth copy with + /// concrete `auth deny ` command posted. #[tokio::test] async fn deferred_busy_blocked_auth_state_posts_auth_hint() { let install = "test-install"; @@ -3502,16 +3540,17 @@ mod tests { ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); - // BlockedAuth → state-aware lookup returns auth wording. - let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( + // 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 = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; - tokio::task::yield_now().await; let calls = egress.calls(); let post_calls: Vec<_> = calls @@ -3528,6 +3567,10 @@ mod tests { 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"), @@ -3535,6 +3578,50 @@ mod tests { ); } + /// DeferredBusy + UserMessage + BlockedApproval with no gate_ref → fallback wording + /// without a specific gate command. + #[tokio::test] + 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"); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "6001.1"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + // 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); + 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 one post for fallback approval hint" + ); + let body = std::str::from_utf8(&post_calls[0].body).expect("utf8 body"); + assert!( + body.contains("waiting on a pending approval"), + "fallback approval hint must still mention pending approval, got: {body}" + ); + } + /// DeferredBusy + UserMessage + Running state (non-blocked) → generic copy posted. #[tokio::test] async fn deferred_busy_running_state_posts_generic_hint() { @@ -3559,7 +3646,6 @@ mod tests { 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 @@ -3619,9 +3705,6 @@ mod tests { let ack = deferred_busy_ack(); observer.observe_workflow_ack(env, ack).await; - // No spawn should have been issued (authorization failed synchronously). - // A yield is still safe to add as a hedge. - tokio::task::yield_now().await; let calls = egress.calls(); assert!( @@ -4910,7 +4993,6 @@ mod tests { 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 @@ -4929,11 +5011,10 @@ mod tests { ); } - /// Slack post inside the spawned hint task fails → no panic, ack path unaffected. + /// Slack post failure → no panic, ack path unaffected. /// - /// The hint-post semaphore must be released and the outer call must return - /// normally. The `FakeProtocolHttpEgress` returns an error response when we - /// program `Err(...)` into the queue. + /// 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"; @@ -4953,7 +5034,6 @@ mod tests { // Must not panic regardless of the egress failure. observer.observe_workflow_ack(env, ack).await; - tokio::task::yield_now().await; // The call was recorded even though the programmed result was an error. let calls = egress.calls(); @@ -4963,16 +5043,75 @@ mod tests { ); } - /// Saturated hint-post semaphore → hint dropped, no post, no panic. - /// - /// Acquire all MAX_CONCURRENT_HINT_POSTS permits from the observer's semaphore - /// before the ack arrives. The observer must silently drop the hint instead of - /// spawning an unbounded task. + /// Two DeferredBusy acks with the same conversation + same active_run_id + /// → exactly one post (throttle suppresses the second). + #[tokio::test] + async fn deferred_busy_same_conversation_same_run_id_posts_once() { + 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", "9001.1"), + )), + ); + + 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_semaphore_saturated_drops_hint_without_panic() { + 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("D123", "9002.1"), + )), + ); + egress.program_response( + "slack.com", + Ok(EgressResponse::new( + 200, + slack_post_ok_json("D123", "9002.2"), + )), + ); let outbound = Arc::new(InMemoryOutboundStateStore::default()); let coordinator = Arc::new(ScriptedTurnCoordinator::with_single_status( @@ -4980,28 +5119,139 @@ mod tests { )); let observer = make_observer(coordinator, egress.clone(), outbound, install); - // Saturate the hint-post semaphore: acquire all MAX_CONCURRENT_HINT_POSTS permits. - let mut permits = Vec::new(); - for _ in 0..MAX_CONCURRENT_HINT_POSTS { - let p = Arc::clone(&observer.hint_post_permits) - .try_acquire_owned() - .expect("permit must be available before saturation"); - permits.push(p); + let make_ack = |run_id: TurnRunId| ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:deferred-diff").expect("ref"), + active_run_id: run_id, + }; + + // 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; + + let calls = egress.calls(); + let post_calls: Vec<_> = calls + .iter() + .filter(|c| c.path == "/api/chat.postMessage") + .collect(); + assert_eq!( + post_calls.len(), + 2, + "different run_ids must produce separate hints even for the same conversation" + ); + } + + /// 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 deferred_busy_uses_ack_active_run_id_and_binding_scope_for_state_lookup() { + use std::sync::Mutex as StdMutex; + + 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"), + )), + ); + + let outbound = Arc::new(InMemoryOutboundStateStore::default()); + + // 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 + } } + let active_run_id = TurnRunId::new(); + let recording_coordinator = Arc::new(RecordingTurnCoordinator { + inner: ScriptedTurnCoordinator::with_single_status(TurnStatus::BlockedApproval), + recorded: StdMutex::new(Vec::new()), + }); + + let observer = make_observer( + recording_coordinator.clone(), + egress.clone(), + outbound, + install, + ); + + let ack = ProductInboundAck::DeferredBusy { + accepted_message_ref: AcceptedMessageRef::new("slack:scope-check").expect("ref"), + active_run_id, + }; let env = envelope(user_message_payload()); - let ack = deferred_busy_ack(); - // Must not panic, must not post. - observer.observe_workflow_ack(env, ack).await; + observer.observe_workflow_ack(env.clone(), ack).await; - let calls = egress.calls(); - assert!( - !calls.iter().any(|c| c.path == "/api/chat.postMessage"), - "no chat.postMessage expected when hint-post semaphore is saturated" + 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" ); - // Release permits explicitly so the test is self-contained. - drop(permits); + // 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" + ); } }