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
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,12 @@ const WS0_COMPOSITION_SHARE_BP: usize = 658;
/// assembly) with `bash scripts/ci/check-composition-budget.sh` -> 41_509.
/// Paired with `[gate].loc_ceiling` in scripts/ci/composition-budget.toml --
/// this ratchet fails when the two disagree.
const COMPOSITION_ABSOLUTE_SRC_LOC: usize = 41_509;
/// ✎ Union re-measured 41_509 → 41_582 on 2026-08-10 after merging main into
/// implement-issue-6896-fix: #7131's run-failure settlement observer adds its
/// production wiring and inline regression coverage to the current main tree.
/// Measured with `bash scripts/ci/check-composition-budget.sh --print`; the
/// manifest ceiling and observed value move with this record.
const COMPOSITION_ABSOLUTE_SRC_LOC: usize = 41_582;

/// Composition dispatch, from the same `--print` run: "composition dispatch:
/// 827 Arc<dyn> (governed prod, excl slack/extension_host)".
Expand Down
5 changes: 4 additions & 1 deletion crates/app/ironclaw_composition/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -189,7 +189,10 @@ tokio = { version = "1", features = ["macros", "net", "rt", "rt-multi-thread", "
# proper WS client.
tokio-tungstenite = { version = "0.29", default-features = false, features = ["connect"] }
tower = { version = "0.5", features = ["util"] }
tracing-test = "0.2"
# `no-env-filter` captures every event regardless of its metadata target —
# the observer test asserts a `target: "ironclaw::reborn::…"` warning, which
# the default `{crate}=trace` filter would drop.
tracing-test = { version = "0.2", features = ["no-env-filter"] }

[[test]]
name = "production_runtime_automations"
Expand Down
62 changes: 61 additions & 1 deletion crates/app/ironclaw_composition/src/automation/trigger_poller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use ironclaw_triggers::{
};
use ironclaw_triggers::{
TriggerAcceptedFireSettlement, TriggerFailedFireSettlement, TriggerFireSettlementObserver,
TriggerRunFailureSettlement,
};
use rand::RngExt;
use tokio::task::JoinHandle;
Expand Down Expand Up @@ -236,6 +237,21 @@ impl TriggerFireSettlementObserver for PostSubmitHookObserver {
};
spawn_post_submit_delivery(hook, TriggerFireSettlement::Failed(event));
}

async fn on_run_failure_settled(&self, event: TriggerRunFailureSettlement) {
// The accepted fire's run reached a terminal failure state. The
// canonical delivery watcher owns the creator-facing notice; this
// observer only emits structured automation-health telemetry and
// must never mint a replacement run or bypass the delivery path.
tracing::warn!(
target: "ironclaw::reborn::trigger_poller",
tenant_id = %event.tenant_id,
trigger_id = %event.trigger_id,
fire_slot = %event.fire_slot,
run_id = %event.run_id,
"accepted trigger fire settled with a failed run"
);
}
}

async fn run_trigger_poller(
Expand Down Expand Up @@ -352,11 +368,12 @@ mod tests {
use ironclaw_triggers::{
TriggerAcceptedFireSettlement, TriggerFailedFireSettlement, TriggerFire,
TriggerFireIdentity, TriggerFireSettlementObserver, TriggerId,
TriggerPollerFailureReason,
TriggerPollerFailureReason, TriggerRunFailureSettlement,
};
use ironclaw_turns::{TurnRunId, TurnScope};
use tokio::sync::Notify;
use tokio_util::sync::CancellationToken;
use tracing_test::traced_test;

use super::super::{POST_SUBMIT_HOOK_PENDING_CAPACITY, PostSubmitHookObserver};

Expand Down Expand Up @@ -486,6 +503,49 @@ mod tests {
}
}

fn run_failure_settlement_event(run_id: TurnRunId) -> TriggerRunFailureSettlement {
TriggerRunFailureSettlement {
tenant_id: observer_tenant(),
trigger_id: TriggerId::new(),
fire_slot: Utc::now(),
run_id,
}
}

#[tokio::test]
#[traced_test]
async fn run_failure_settlement_emits_health_warning_without_delivery() {
let hook_slot = Arc::new(std::sync::OnceLock::new());
let recording = Arc::new(RecordingHook::default());
let observer = PostSubmitHookObserver::new(hook_slot.clone(), CancellationToken::new());
hook_slot.set(recording.clone()).ok().expect("hook install");

observer
.on_run_failure_settled(run_failure_settlement_event(TurnRunId::new()))
.await;

assert!(
logs_contain("accepted trigger fire settled with a failed run"),
"production observer must emit the automation-health signal; buffer: {:?}",
String::from_utf8(
tracing_test::internal::global_buf()
.lock()
.expect("buf")
.to_vec()
)
.expect("utf8")
);
assert!(
recording.calls().is_empty()
&& recording
.failed_calls
.lock()
.unwrap_or_else(|p| p.into_inner())
.is_empty(),
"run-failure observation must not bypass the canonical delivery watcher"
);
}

#[tokio::test]
async fn uninstalled_hook_buffers_until_hook_is_installed() {
let run_id = TurnRunId::new();
Expand Down
6 changes: 3 additions & 3 deletions crates/domains/ironclaw_triggers/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1344,9 +1344,9 @@ pub use worker::{
TriggerActiveRunState, TriggerActiveRunStateRequest, TriggerFailedFireSettlement,
TriggerFireSettlementObserver, TriggerPollerFailureReason, TriggerPollerFireOutcome,
TriggerPollerFireReport, TriggerPollerTickReport, TriggerPollerWorker,
TriggerPollerWorkerConfig, TriggerPollerWorkerDeps, TrustedTriggerFireSubmitOutcome,
TrustedTriggerFireSubmitter, TrustedTriggerSubmitRequest, active_hold_projection,
active_holds_for_records,
TriggerPollerWorkerConfig, TriggerPollerWorkerDeps, TriggerRunFailureSettlement,
TrustedTriggerFireSubmitOutcome, TrustedTriggerFireSubmitter, TrustedTriggerSubmitRequest,
active_hold_projection, active_holds_for_records,
};

#[derive(Clone, Default)]
Expand Down
5 changes: 3 additions & 2 deletions crates/domains/ironclaw_triggers/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,9 @@ pub use ports::{
ActiveHoldReason, BlockedActiveRunKind, MissingTriggerActiveRunLookup,
NoopTriggerFireSettlementObserver, TriggerAcceptedFireSettlement, TriggerActiveRunLookup,
TriggerActiveRunState, TriggerActiveRunStateRequest, TriggerFailedFireSettlement,
TriggerFireSettlementObserver, TrustedTriggerFireSubmitOutcome, TrustedTriggerFireSubmitter,
TrustedTriggerSubmitRequest, active_hold_projection, active_holds_for_records,
TriggerFireSettlementObserver, TriggerRunFailureSettlement, TrustedTriggerFireSubmitOutcome,
TrustedTriggerFireSubmitter, TrustedTriggerSubmitRequest, active_hold_projection,
active_holds_for_records,
};
pub use report::{
TriggerPollerFailureReason, TriggerPollerFireOutcome, TriggerPollerFireReport,
Expand Down
31 changes: 27 additions & 4 deletions crates/domains/ironclaw_triggers/src/worker/active_cleanup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use crate::{
use super::{
TriggerActiveRunState, TriggerActiveRunStateRequest, TriggerPollerFailureReason,
TriggerPollerFireOutcome, TriggerPollerFireReport, TriggerPollerTickReport,
TriggerPollerWorker, failure::classify_failure,
TriggerPollerWorker, TriggerRunFailureSettlement, failure::classify_failure,
};

struct ActiveLookupItem {
Expand Down Expand Up @@ -145,7 +145,7 @@ impl TriggerPollerWorker {
};
match clear {
Some((cleared_outcome, status)) => {
let outcome = if self
let cleared = self
.deps
.repository
.clear_active_fire(ClearActiveFireRequest {
Expand All @@ -156,8 +156,31 @@ impl TriggerPollerWorker {
status,
})
.await?
.is_some()
{
.is_some();
// A terminal failure (run ended `Failed` / `Cancelled` /
// `RecoveryRequired`, cleared as `Error`) previously had no
// settlement failure hook — the active-cleanup sweep cleared
// the slot with no observability. Surface it to the
// settlement observer so automation-health telemetry can see
Comment thread
serrrfirat marked this conversation as resolved.
// post-accept failures (#6896). `Ok`/`Running` and
// already-cleared fires do not fire the hook — the former
// are not failures, and the latter already did (or never
// reached terminal here). The observer is invoked inline
// per the `TriggerFireSettlementObserver` contract: it
// must be cheap and non-blocking, detaching any heavy
// work internally.
if cleared && status == TriggerRunHistoryStatus::Error {
self.deps
.fire_settlement_observer
.on_run_failure_settled(TriggerRunFailureSettlement {
tenant_id: record.tenant_id.clone(),
trigger_id: record.trigger_id,
fire_slot,
run_id,
})
.await;

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Low — Inline await of on_failed_fire_settled stalls active-cleanup sweep.

The new failure-settlement hook is awaited inline inside the per-record sweep loop, between the durable clear_active_fire and report/cursor advancement. Today's impls are cheap (production logs one warn; Noop does nothing; tests push to a Mutex), and the loop already awaits backend I/O per record, so marginal latency is negligible. But the trait is async and open (Send+Sync, any implementor); the sibling on_accepted_fire_settled path deliberately decouples via bounded spawn_post_submit_delivery + bounded buffer, while this path couples sweep latency and scan-cursor progress to observer behavior. A future heavy observer (e.g., network telemetry sink) would stall active-fire cleanup for the whole tick. No lock is held across the await, so no deadlock is introduced.

Fix: Detach the observer call with a bounded spawn (mirroring on_accepted_fire_settled), or harden the trait contract docs to require cheap non-blocking observers.

}
let outcome = if cleared {
cleared_outcome
} else {
TriggerPollerFireOutcome::SkippedAlreadyCleared { run_id }
Expand Down
31 changes: 31 additions & 0 deletions crates/domains/ironclaw_triggers/src/worker/ports.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,11 +136,40 @@ pub struct TriggerFailedFireSettlement {
pub reason: TriggerPollerFailureReason,
}

/// A previously-accepted fire whose run reached a terminal *failure* state
/// observed by the active-cleanup sweep. Carries enough identity to let an
/// observer correlate it back to the originating fire for automation-health
/// telemetry without re-deriving it from display strings.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TriggerRunFailureSettlement {
pub tenant_id: TenantId,
pub trigger_id: TriggerId,
pub fire_slot: Timestamp,
pub run_id: TurnRunId,
}

#[async_trait]
pub trait TriggerFireSettlementObserver: Send + Sync {
/// The worker invokes these hooks **inline** while processing poller
/// work: `on_accepted_fire_settled` from the per-fire submit path,
/// `on_failed_fire_settled` from the claim-time failure path, and
/// `on_run_failure_settled` from the active-cleanup sweep. Implementors
/// MUST be cheap and non-blocking; any heavy work (delivery, telemetry
/// egress) must be detached internally (for example a bounded spawn, as
/// the composition `PostSubmitHookObserver` does for accepted-fire
/// delivery) rather than awaited through this call.
async fn on_accepted_fire_settled(&self, event: TriggerAcceptedFireSettlement);

async fn on_failed_fire_settled(&self, _event: TriggerFailedFireSettlement) {}

/// A previously-accepted fire settled into a terminal failure state.
///
/// Implementors must handle this explicitly so production composition
/// cannot silently discard the automation-health signal. The callback is
/// strictly observational — it must not mint runs or mutate trigger state.
/// Invoked inline from the active-cleanup sweep; see the trait contract
/// above: keep this implementation cheap and non-blocking.
async fn on_run_failure_settled(&self, event: TriggerRunFailureSettlement);
}

#[derive(Debug, Default)]
Expand All @@ -149,6 +178,8 @@ pub struct NoopTriggerFireSettlementObserver;
#[async_trait]
impl TriggerFireSettlementObserver for NoopTriggerFireSettlementObserver {
async fn on_accepted_fire_settled(&self, _event: TriggerAcceptedFireSettlement) {}

async fn on_run_failure_settled(&self, _event: TriggerRunFailureSettlement) {}
}

#[derive(Debug, Clone, PartialEq, Eq)]
Expand Down
Loading
Loading