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
2 changes: 1 addition & 1 deletion crates/nexum-runtime-metrics/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ pub const METRICS: &[Metric] = &[
Metric {
name: "nexum_runtime_dispatch_latency_seconds",
kind: Kind::Histogram,
help: "Wall-clock seconds to dispatch one trigger.",
help: "Wall-clock seconds to dispatch one trigger, by module and outcome.",
},
Metric {
name: "nexum_runtime_log_fields_dropped_total",
Expand Down
93 changes: 70 additions & 23 deletions crates/nexum-runtime-supervisor/src/supervisor/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,19 +312,18 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
"nexum_runtime_dispatch_dropped_total",
"module" => module.name.to_string(),
"trigger_kind" => trigger_kind,
"reason" => "rate_limited",
"reason" => <&str>::from(DispatchOutcome::RateLimited),
)
.increment(1);
return DispatchOutcome::RateLimited;
}
if let Err(e) = module.live.store.set_fuel(module.seed.spec.fuel) {
error!(
module = %module.name,
chain_id,
trigger_kind,
error = %e,
"set_fuel failed - skipping"
);
if !refuel(
&mut module.live.store,
module.seed.spec.fuel,
&module.name,
chain_id,
trigger_kind,
) {
return DispatchOutcome::FuelSetFailed;
}
let start = Instant::now();
Expand All @@ -335,16 +334,16 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
.live
.bindings
.call_on_trigger(&mut module.live.store, trigger);
let outcome = with_dispatch_deadline(deadline, call)
.await
.unwrap_or_else(|exceeded| Err(wasmtime::Error::from(exceeded)));
let bounded = with_dispatch_deadline(deadline, call).await;
let deadline_hit = bounded.is_err();
let result = bounded.unwrap_or_else(|exceeded| Err(wasmtime::Error::from(exceeded)));
// One post-call sample: the trap instant is start plus elapsed, not
// the pre-dispatch `now`. This is the same clock the lifecycle
// reads, so under `start_paused` the latency histogram records the
// virtual elapsed time: do not assert on it from a paused test.
let elapsed = start.elapsed();
let latency_ms = elapsed.as_millis() as u64;
match outcome {
let outcome = match result {
Ok(Ok(())) => {
debug!(
module = %module.name,
Expand All @@ -354,12 +353,6 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
latency_ms,
"dispatch ok"
);
metrics::histogram!(
"nexum_runtime_dispatch_latency_seconds",
"module" => module.name.to_string(),
"trigger_kind" => trigger_kind,
)
.record(elapsed.as_secs_f64());
module.health.dispatch_succeeded();
DispatchOutcome::Ok
}
Expand Down Expand Up @@ -388,6 +381,14 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
let died_at = start + elapsed;
let seed = crate::runtime::restart_policy::jitter_seed(module.name.as_str());
let verdict = module.health.record_trap(died_at, poison_policy, seed);
// One label for one event: the histogram's outcome and this
// counter's kind come from the same variant, so a deadline
// cannot read as a bug on one metric and a limit on the other.
let outcome = if deadline_hit {
DispatchOutcome::Deadline
} else {
DispatchOutcome::Trapped
};
error!(
module = %module.name,
chain_id,
Expand All @@ -402,7 +403,7 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
metrics::counter!(
"nexum_runtime_module_errors_total",
"module" => module.name.clone(),
"error_kind" => "trap",
"error_kind" => <&str>::from(outcome),
)
.increment(1);
// Leave a retrievable panic record on the dead run; the full
Expand All @@ -416,12 +417,52 @@ impl<T: RuntimeTypes<State = HostState<T>>> Supervisor<T> {
if let Some(recent) = verdict.poisoned {
report_poison(&module.name, recent, poison_policy.window, trap.to_string());
}
DispatchOutcome::Trapped
outcome
}
}
};
// Every outcome that entered the guest contributes a sample: a
// success-only histogram reads p95 low exactly when dispatch is
// failing.
metrics::histogram!(
"nexum_runtime_dispatch_latency_seconds",
"module" => module.name.to_string(),
"trigger_kind" => trigger_kind,
"outcome" => <&str>::from(outcome),
)
.record(elapsed.as_secs_f64());
outcome
}
}

/// Grants the dispatch its fuel budget; `false` means the trigger never
/// reached the guest and was counted as a drop.
pub(super) fn refuel<S>(
store: &mut wasmtime::Store<S>,
fuel: u64,
module: &ModuleId,
chain_id: u64,
trigger_kind: &'static str,
) -> bool {
let Err(e) = store.set_fuel(fuel) else {
return true;
};
error!(
module = %module,
chain_id,
trigger_kind,
error = %e,
"set_fuel failed - skipping"
);
metrics::counter!(
"nexum_runtime_dispatch_dropped_total",
"module" => module.clone(),
"trigger_kind" => trigger_kind,
"reason" => <&str>::from(DispatchOutcome::FuelSetFailed),
)
.increment(1);
false
}

/// The poison-transition trio: quarantine warn plus gauge, with the trap
/// that crossed the threshold.
fn report_poison(name: &ModuleId, recent_failures: u32, window: Duration, last_error: String) {
Expand Down Expand Up @@ -459,13 +500,19 @@ pub(super) async fn with_dispatch_deadline<F: std::future::Future>(
.map_err(|_elapsed| DeadlineExceeded(deadline))
}

#[derive(Debug)]
/// The `outcome` label on the latency histogram; the two pre-guest drops carry
/// theirs as the `reason` on `nexum_runtime_dispatch_dropped_total` instead.
#[derive(Debug, Clone, Copy, strum::IntoStaticStr)]
#[strum(serialize_all = "snake_case")]
pub(super) enum DispatchOutcome {
Ok,
/// Guest returned a typed `fault` via WIT, not a trap; the module stays alive.
Fault,
/// Marked dead, maybe quarantined per the poison policy.
#[strum(serialize = "trap")]
Trapped,
/// Cut off at the wall-clock deadline; handled like a trap.
Deadline,
/// `set_fuel` failed before the call; the module stays alive, the trigger skips.
FuelSetFailed,
/// Dropped before the guest runs; liveness untouched.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,7 @@ fn a_boot_refusal_increments_the_counter_under_its_parse_class() {
// Raw TOML: the textual absence of [dependencies] is the fixture.
let manifest = "[component]\nname = \"example\"\n";
let (refusal, samples) = capture_metrics(|| {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime")
.block_on(scenario().module(manifest.to_owned()).expect_refusal())
block_on_current_thread(scenario().module(manifest.to_owned()).expect_refusal())
});
refusal.variant::<BootRefusal>(|e| {
matches!(e, BootRefusal::Manifest(ParseError::MissingCapabilities))
Expand Down
53 changes: 22 additions & 31 deletions crates/nexum-runtime-supervisor/src/supervisor/tests/chain_gate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -244,21 +244,17 @@ fn a_guest_request_for_an_unconfigured_chain_counts_under_the_sentinel() {
// `scenario()` boots over an empty pool while `test_chain_configs()`
// still declares chain 1.
let (dispatched, samples) = capture_metrics(|| {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime")
.block_on(async {
let mut booted = scenario()
.wasm(wasm)
.module(workspace_manifest(
"modules/fixtures/slow-host/component.toml",
))
.boot()
.await
.expect("slow-host boots over the empty pool");
booted.dispatch_block_on(1).await
})
block_on_current_thread(async {
let mut booted = scenario()
.wasm(wasm)
.module(workspace_manifest(
"modules/fixtures/slow-host/component.toml",
))
.boot()
.await
.expect("slow-host boots over the empty pool");
booted.dispatch_block_on(1).await
})
});
assert_eq!(dispatched, 1, "the fixture swallows the request error");
let hits = samples_named(&samples, "nexum_runtime_chain_request_total");
Expand Down Expand Up @@ -288,22 +284,17 @@ fn a_guest_request_for_a_configured_chain_keeps_its_chain_id() {
let node = FakeNode::new();
node.on_method(nexum_world::ChainMethod::EthBlockNumber, "\"0x1\"");
let (dispatched, samples) = capture_metrics(|| {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime")
.block_on(async {
let mut booted =
BootScenario::over(mock_components_from(&node, MockStateStore::new()))
.wasm(wasm)
.module(workspace_manifest(
"modules/fixtures/slow-host/component.toml",
))
.boot()
.await
.expect("slow-host boots over the mocked pool");
booted.dispatch_block_on(1).await
})
block_on_current_thread(async {
let mut booted = BootScenario::over(mock_components_from(&node, MockStateStore::new()))
.wasm(wasm)
.module(workspace_manifest(
"modules/fixtures/slow-host/component.toml",
))
.boot()
.await
.expect("slow-host boots over the mocked pool");
booted.dispatch_block_on(1).await
})
});
assert_eq!(dispatched, 1);
let hits = samples_named(&samples, "nexum_runtime_chain_request_total");
Expand Down
Loading
Loading