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
22 changes: 14 additions & 8 deletions crates/core/src/host/host_controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ use crate::subscription::row_list_builder_pool::BsatnRowListBuilderPool;
use crate::util::asyncify;
use crate::util::jobs::{AllocatedJobCore, JobCores};
use crate::worker_metrics::{
record_module_host_init_attempt, record_module_host_init_failure, ModuleHostInitFailureCause, WORKER_METRICS,
record_module_host_init_attempt, record_module_host_init_failure, record_module_host_unexpected_exit,
ModuleHostInitFailureCause, WORKER_METRICS,
};
use anyhow::{anyhow, bail, Context};
use async_trait::async_trait;
Expand Down Expand Up @@ -569,6 +570,7 @@ impl HostController {
// `HostController::clone` is fast,
// as all of its fields are either `Copy` or wrapped in `Arc`.
let this = self.clone();
let database_identity = database.database_identity;

// `try_init_host` is not cancel safe, as it will spawn other async tasks
// which hold a filesystem lock past when `try_init_host` returns or is cancelled.
Expand Down Expand Up @@ -601,7 +603,7 @@ impl HostController {
program,
policy,
this.energy_monitor.clone(),
this.unregister_fn(replica_id),
this.unregister_fn(replica_id, database_identity),
this.db_cores.take(),
)
.await?;
Expand Down Expand Up @@ -764,11 +766,14 @@ impl HostController {
/// On-panic callback passed to [`ModuleHost`]s created by this controller.
///
/// Removes the module with the given `replica_id` from this controller.
fn unregister_fn(&self, replica_id: u64) -> impl Fn() + Send + Sync + 'static + use<> {
fn unregister_fn(&self, replica_id: u64, database_identity: Identity) -> impl Fn() + Send + Sync + 'static + use<> {
let hosts = Arc::downgrade(&self.hosts);
move || {
if let Some(hosts) = hosts.upgrade() {
hosts.lock().remove(&replica_id);
let unregistered = hosts
.upgrade()
.is_some_and(|hosts| hosts.lock().remove(&replica_id).is_some());
if unregistered {
record_module_host_unexpected_exit(database_identity);
}
}
}
Expand Down Expand Up @@ -1093,6 +1098,7 @@ impl Host {
database: Database,
replica_id: u64,
) -> anyhow::Result<HostInit> {
let database_identity = database.database_identity;
let HostController {
data_dir,
default_config: config,
Expand Down Expand Up @@ -1197,7 +1203,7 @@ impl Host {
database,
replica_id,
program,
on_panic: host_controller.unregister_fn(replica_id),
on_panic: host_controller.unregister_fn(replica_id, database_identity),
relational_db,
energy_monitor: energy_monitor.clone(),
memory_observer: memory_observer.clone(),
Expand Down Expand Up @@ -1227,7 +1233,7 @@ impl Host {
database: database.clone(),
replica_id,
program: program.clone(),
on_panic: host_controller.unregister_fn(replica_id),
on_panic: host_controller.unregister_fn(replica_id, database_identity),
relational_db: relational_db.clone(),
energy_monitor: energy_monitor.clone(),
memory_observer: memory_observer.clone(),
Expand All @@ -1251,7 +1257,7 @@ impl Host {
database,
replica_id,
program: program.clone(),
on_panic: host_controller.unregister_fn(replica_id),
on_panic: host_controller.unregister_fn(replica_id, database_identity),
relational_db: relational_db.clone(),
energy_monitor: energy_monitor.clone(),
memory_observer: memory_observer.clone(),
Expand Down
12 changes: 12 additions & 0 deletions crates/core/src/worker_metrics/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,13 @@ pub fn record_module_host_init_failure(database_identity: Identity, cause: Modul
.inc();
}

pub fn record_module_host_unexpected_exit(database_identity: Identity) {
WORKER_METRICS
.module_host_unexpected_exits
.with_label_values(&database_identity)
.inc();
}

/// Records at most one disconnect cause for a single accepted websocket client.
#[derive(Clone, Debug)]
pub struct ClientDisconnectRecorder {
Expand Down Expand Up @@ -543,6 +550,11 @@ metrics_group!(
#[labels(database_identity: Identity, cause: str)]
pub module_host_init_failures: IntCounterVec,

#[name = spacetime_module_host_unexpected_exits_total]
#[help = "The cumulative number of unexpected module host exits"]
#[labels(database_identity: Identity)]
pub module_host_unexpected_exits: IntCounterVec,

#[name = spacetime_reducer_wait_time_sec]
#[help = "The amount of time (in seconds) a reducer spends in the queue waiting to run"]
#[labels(db: Identity, reducer: str)]
Expand Down
Loading