Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 36 additions & 35 deletions crates/ironclaw_reborn_composition/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -479,12 +479,12 @@ pub struct RebornRuntime {
#[cfg(feature = "root-llm-provider")]
skill_learning_extraction_tasks:
Option<Arc<crate::skill_learning::SkillLearningExtractionTasks>>,
/// Late-binding slot shared with the poller's `PostSubmitHookWrappedSubmitter`.
/// `set_trigger_post_submit_hook` fills this after `build_reborn_runtime` returns.
/// Late-binding dispatcher shared with the trigger poller.
/// `set_trigger_post_submit_hook` installs the hook after
/// `build_reborn_runtime` returns and drains any startup settlements.
/// `None` when the trigger poller is not enabled.
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot:
Option<Arc<std::sync::OnceLock<Arc<dyn crate::slack_delivery::PostSubmitDeliveryHook>>>>,
post_submit_hook_dispatch: Option<Arc<crate::trigger_poller::PostSubmitHookDispatch>>,
#[cfg(any(test, feature = "test-support"))]
trigger_conversation_pairing:
Option<Arc<dyn ironclaw_conversations::ConversationActorPairingService>>,
Expand Down Expand Up @@ -562,13 +562,11 @@ type LocalDevSkillExecutionAdapter =
struct TriggerPollerServices {
materializer: Arc<dyn ironclaw_triggers::TriggerPromptMaterializer>,
trusted_submitter: Arc<dyn ironclaw_triggers::TrustedTriggerFireSubmitter>,
/// Late-binding slot for the post-submit hook. Created here and shared with
/// the poller wrapper; filled later by `RebornRuntime::set_trigger_post_submit_hook`
/// so `build_slack_host_beta_mounts` (called after runtime build) can wire the
/// hook without restarting the poller.
/// Late-binding dispatcher for the post-submit hook. Created here and shared
/// with the poller; `RebornRuntime::set_trigger_post_submit_hook` installs
/// the hook later and drains accepted fires settled during startup.
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot:
Arc<std::sync::OnceLock<Arc<dyn crate::slack_delivery::PostSubmitDeliveryHook>>>,
post_submit_hook_dispatch: Arc<crate::trigger_poller::PostSubmitHookDispatch>,
/// Test-support handle on the SAME conversation services instance the
/// poller-side materializer/submitter use, so integration tests can call
/// the production `pair_external_actor` API to seed the trigger
Expand Down Expand Up @@ -617,7 +615,9 @@ async fn build_trigger_poller_services(
materializer,
trusted_submitter,
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot: Arc::new(std::sync::OnceLock::new()),
post_submit_hook_dispatch: Arc::new(
crate::trigger_poller::PostSubmitHookDispatch::new(),
),
#[cfg(any(test, feature = "test-support"))]
pairing_service,
})
Expand All @@ -644,7 +644,9 @@ async fn build_trigger_poller_services(
materializer,
trusted_submitter,
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot: Arc::new(std::sync::OnceLock::new()),
post_submit_hook_dispatch: Arc::new(
crate::trigger_poller::PostSubmitHookDispatch::new(),
),
#[cfg(any(test, feature = "test-support"))]
pairing_service,
})
Expand Down Expand Up @@ -1326,33 +1328,32 @@ impl RebornRuntime {
/// the hook itself is constructed (e.g. inside
/// [`crate::slack_host_beta::build_slack_host_beta_mounts`]). The hook is
/// idempotent: a second call is silently ignored. Returns `false` when the
/// trigger poller is not enabled (slot is `None`) or the slot is already
/// occupied, `true` on first successful set.
/// trigger poller is not enabled (no dispatcher) or a hook is already
/// installed, `true` on first successful install.
#[cfg(feature = "slack-v2-host-beta")]
pub fn set_trigger_post_submit_hook(
&self,
hook: Arc<dyn crate::slack_delivery::PostSubmitDeliveryHook>,
) -> bool {
let Some(slot) = self.post_submit_hook_slot.as_ref() else {
let Some(dispatch) = self.post_submit_hook_dispatch.as_ref() else {
tracing::debug!("set_trigger_post_submit_hook: trigger poller not enabled, ignoring");
return false;
};
match slot.set(hook) {
Ok(()) => true,
Err(_) => {
tracing::debug!(
"set_trigger_post_submit_hook: slot already occupied, ignoring (idempotent)"
);
false
}
if dispatch.install_hook(hook) {
true
} else {
tracing::debug!(
"set_trigger_post_submit_hook: hook already installed, ignoring (idempotent)"
);
false
}
}

#[cfg(feature = "slack-v2-host-beta")]
pub(crate) fn trigger_post_submit_hook_is_set(&self) -> bool {
self.post_submit_hook_slot
self.post_submit_hook_dispatch
.as_ref()
.is_some_and(|slot| slot.get().is_some())
.is_some_and(|dispatch| dispatch.is_hook_installed())
}

#[cfg(test)]
Expand Down Expand Up @@ -3210,16 +3211,16 @@ pub async fn build_reborn_runtime(
};
services.turn_coordinator = Some(Arc::clone(&planned_turn_coordinator));

// `trigger_poller_handle`, `post_submit_hook_slot`, and the test-support
// `trigger_conversation_pairing_value` are produced atomically inside
// a single `if trigger_poller.enabled` expression. Avoid a
// `trigger_poller_handle`, `post_submit_hook_dispatch`, and the
// test-support `trigger_conversation_pairing_value` are produced atomically
// inside a single `if trigger_poller.enabled` expression. Avoid a
// `let mut … = None` sentinel pattern flagged by code review
// (review f-ptr-3): the `let X;` deferred-init form is single-assign
// per branch and Rust's borrow checker prevents reads before init.
let trigger_poller_handle: Option<TriggerPollerRuntimeHandle>;
#[cfg(feature = "slack-v2-host-beta")]
let runtime_post_submit_hook_slot: Option<
Arc<std::sync::OnceLock<Arc<dyn crate::slack_delivery::PostSubmitDeliveryHook>>>,
let runtime_post_submit_hook_dispatch: Option<
Arc<crate::trigger_poller::PostSubmitHookDispatch>,
>;
#[cfg(any(test, feature = "test-support"))]
let trigger_conversation_pairing_value: Option<
Expand Down Expand Up @@ -3251,10 +3252,10 @@ pub async fn build_reborn_runtime(
Some(Arc::clone(&trigger_poller_services.pairing_service));
}
#[cfg(feature = "slack-v2-host-beta")]
let hook_slot = Arc::clone(&trigger_poller_services.post_submit_hook_slot);
let hook_dispatch = Arc::clone(&trigger_poller_services.post_submit_hook_dispatch);
#[cfg(feature = "slack-v2-host-beta")]
{
runtime_post_submit_hook_slot = Some(Arc::clone(&hook_slot));
runtime_post_submit_hook_dispatch = Some(Arc::clone(&hook_dispatch));
}
trigger_poller_handle = spawn_trigger_poller(
trigger_poller,
Expand All @@ -3264,7 +3265,7 @@ pub async fn build_reborn_runtime(
trusted_submitter: trigger_poller_services.trusted_submitter,
active_run_lookup,
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot: hook_slot,
post_submit_hook_dispatch: hook_dispatch,
},
)
.map_err(|error| RebornRuntimeError::InvalidArgument {
Expand All @@ -3274,7 +3275,7 @@ pub async fn build_reborn_runtime(
trigger_poller_handle = None;
#[cfg(feature = "slack-v2-host-beta")]
{
runtime_post_submit_hook_slot = None;
runtime_post_submit_hook_dispatch = None;
}
#[cfg(any(test, feature = "test-support"))]
{
Expand Down Expand Up @@ -3351,7 +3352,7 @@ pub async fn build_reborn_runtime(
#[cfg(feature = "root-llm-provider")]
skill_learning_extraction_tasks,
#[cfg(feature = "slack-v2-host-beta")]
post_submit_hook_slot: runtime_post_submit_hook_slot,
post_submit_hook_dispatch: runtime_post_submit_hook_dispatch,
#[cfg(any(test, feature = "test-support"))]
trigger_conversation_pairing: trigger_conversation_pairing_value,
outbound_delivery_target_registry,
Expand Down
77 changes: 9 additions & 68 deletions crates/ironclaw_reborn_composition/src/slack_delivery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1804,16 +1804,18 @@ fn is_accepted_auth_denial(envelope: &ProductInboundEnvelope, ack: &ProductInbou
// using the same polling machinery as the observer path (above).

/// Composition-owned hook invoked by the trigger poller after a successful fire
/// submission. The composition root wires a real implementation (behind the
/// `slack-v2-host-beta` feature) or a no-op.
/// submission has been durably settled in trigger storage. The composition root
/// wires a real implementation (behind the `slack-v2-host-beta` feature) or a
/// no-op.
#[async_trait]
pub trait PostSubmitDeliveryHook: Send + Sync {
/// Called with the original trigger fire, the submitted run id, and the
/// turn scope the run was submitted under. The trigger poller owns the
/// non-blocking handoff by invoking this hook from a detached task, so hook
/// latency cannot delay fire settlement. Implementations may still spawn
/// their own longer-lived delivery tasks when they need bounded admission or
/// shutdown tracking.
/// turn scope the run was submitted under. The trigger poller invokes this
/// hook from a detached task after the accepted fire appears as settled, so
/// hook latency cannot delay settlement and delivery cannot precede the
/// persisted run/thread mapping. Implementations may still spawn their own
/// longer-lived delivery tasks when they need bounded admission or shutdown
/// tracking.
async fn on_trigger_submitted(&self, fire: TriggerFire, run_id: TurnRunId, scope: TurnScope);
}

Expand Down Expand Up @@ -2976,8 +2978,6 @@ mod tests {
// the plan explicitly forbids exposing `host_state_filesystem` from
// `RebornRuntime` for that purpose.

use std::sync::OnceLock;

use ironclaw_outbound::{
CommunicationPreferenceRecord, DeliveryDefaultScope, InMemoryDeliveredGateRouteStore,
InMemoryOutboundStateStore, InMemoryTriggeredRunDeliveryStore,
Expand Down Expand Up @@ -5853,65 +5853,6 @@ mod tests {
);
}

// --- 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<OnceLock<Arc<dyn PostSubmitDeliveryHook>>> = 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<Mutex<bool>>);
#[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<OnceLock<Arc<dyn PostSubmitDeliveryHook>>> = Arc::new(OnceLock::new());
let hook_a: Arc<dyn PostSubmitDeliveryHook> = Arc::new(NoopPostSubmitDeliveryHook);
let hook_b: Arc<dyn PostSubmitDeliveryHook> = 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!(
Expand Down
Loading
Loading