use std::fmt::{Display, Formatter};
use serde::{Deserialize, Serialize};
use crate::mob_handle_runtime::MobRuntimeError;
use crate::runtime::{
NormalizationError, RuntimeRouteMutationError, RuntimeShutdownReport, SubscribeError,
};
use super::edge_types::{DesiredPeerEdge, EdgeReconcileFailure};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct UnifiedRuntimeReconcileEdgesReport {
pub desired_edges: Vec<DesiredPeerEdge>,
pub wired_edges: Vec<DesiredPeerEdge>,
pub unwired_edges: Vec<DesiredPeerEdge>,
pub retained_edges: Vec<DesiredPeerEdge>,
pub preexisting_edges: Vec<DesiredPeerEdge>,
pub skipped_missing_members: Vec<DesiredPeerEdge>,
pub pruned_stale_managed_edges: Vec<DesiredPeerEdge>,
#[serde(default)]
pub failures: Vec<EdgeReconcileFailure>,
}
impl UnifiedRuntimeReconcileEdgesReport {
pub fn is_complete(&self) -> bool {
self.failures.is_empty() && self.skipped_missing_members.is_empty()
}
}
#[derive(Debug)]
pub enum UnifiedRuntimeBootstrapError {
Mob(MobRuntimeError),
Module(crate::runtime::MobkitRuntimeError),
ModuleStartupThreadPanicked,
ModuleStartupRollbackFailed {
startup_error: Box<UnifiedRuntimeBootstrapError>,
rollback_error: MobRuntimeError,
},
PreSpawnHook(String),
IdentityFirst(String),
Topology(String),
}
impl Display for UnifiedRuntimeBootstrapError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::Mob(err) => write!(f, "failed to bootstrap mob runtime: {err}"),
Self::Module(err) => write!(f, "failed to bootstrap module runtime: {err:?}"),
Self::ModuleStartupThreadPanicked => {
write!(
f,
"failed to bootstrap module runtime: startup thread panicked"
)
}
Self::PreSpawnHook(err) => {
write!(f, "pre-spawn hook failed: {err}")
}
Self::IdentityFirst(err) => {
write!(f, "identity-first bootstrap failed: {err}")
}
Self::Topology(err) => write!(f, "topology-control bootstrap failed: {err}"),
Self::ModuleStartupRollbackFailed {
startup_error,
rollback_error,
} => {
write!(
f,
"failed to bootstrap unified runtime: startup error ({startup_error}) and rollback failed: {rollback_error}"
)
}
}
}
}
impl std::error::Error for UnifiedRuntimeBootstrapError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum UnifiedRuntimeBuilderField {
MobSpec,
ModuleConfig,
Timeout,
}
#[derive(Debug)]
pub enum UnifiedRuntimeBuilderError {
MissingRequiredField(UnifiedRuntimeBuilderField),
Bootstrap(UnifiedRuntimeBootstrapError),
Io(String),
DefinitionLoad(String),
ConflictingConfiguration(String),
StorageLayout(crate::storage_layout::StorageLayoutError),
StorageProvider(crate::storage_provider::MobKitStorageProviderError),
}
impl From<crate::storage_layout::StorageLayoutError> for UnifiedRuntimeBuilderError {
fn from(error: crate::storage_layout::StorageLayoutError) -> Self {
Self::StorageLayout(error)
}
}
impl From<crate::storage_provider::MobKitStorageProviderError> for UnifiedRuntimeBuilderError {
fn from(error: crate::storage_provider::MobKitStorageProviderError) -> Self {
Self::StorageProvider(error)
}
}
impl Display for UnifiedRuntimeBuilderError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::MissingRequiredField(UnifiedRuntimeBuilderField::MobSpec) => {
write!(f, "missing required builder field: mob_spec or definition")
}
Self::MissingRequiredField(UnifiedRuntimeBuilderField::ModuleConfig) => {
write!(f, "missing required builder field: module_config")
}
Self::MissingRequiredField(UnifiedRuntimeBuilderField::Timeout) => {
write!(f, "missing required builder field: timeout")
}
Self::Bootstrap(err) => write!(f, "{err}"),
Self::Io(msg) => write!(f, "{msg}"),
Self::DefinitionLoad(msg) => write!(f, "{msg}"),
Self::ConflictingConfiguration(msg) => write!(f, "conflicting configuration: {msg}"),
Self::StorageLayout(err) => write!(f, "{err}"),
Self::StorageProvider(err) => write!(f, "{err}"),
}
}
}
impl std::error::Error for UnifiedRuntimeBuilderError {}
#[derive(Debug)]
pub enum UnifiedRuntimeError {
Normalize(NormalizationError),
Subscribe(SubscribeError),
RuntimeShuttingDown,
}
impl Display for UnifiedRuntimeError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::Normalize(err) => write!(f, "failed to normalize unified event: {err:?}"),
Self::Subscribe(err) => write!(f, "failed to subscribe to unified events: {err:?}"),
Self::RuntimeShuttingDown => write!(f, "unified runtime is shutting down"),
}
}
}
impl std::error::Error for UnifiedRuntimeError {}
impl From<NormalizationError> for UnifiedRuntimeError {
fn from(value: NormalizationError) -> Self {
Self::Normalize(value)
}
}
impl From<SubscribeError> for UnifiedRuntimeError {
fn from(value: SubscribeError) -> Self {
Self::Subscribe(value)
}
}
#[derive(Debug)]
pub enum MobStopOutcome {
Stopped,
ProceededWithoutInterrupt {
waited_ms: u64,
member: Option<String>,
error: String,
},
Failed(crate::mob_handle_runtime::MobRuntimeError),
}
impl MobStopOutcome {
pub fn teardown_may_proceed(&self) -> bool {
matches!(self, Self::Stopped | Self::ProceededWithoutInterrupt { .. })
}
pub fn stopped_cleanly(&self) -> bool {
matches!(self, Self::Stopped)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IdentityAuthorityReleaseOutcome {
NotConfigured,
Released { grant_count: usize },
Failed { error: String },
SkippedResetCleanupFailed { error: String },
SkippedMobStopFailed,
}
#[derive(Debug)]
pub struct UnifiedRuntimeShutdownReport {
pub drain: ShutdownDrainReport,
pub module_shutdown: RuntimeShutdownReport,
pub mob_stop: Result<(), MobRuntimeError>,
pub identity_authority_release: IdentityAuthorityReleaseOutcome,
}
impl UnifiedRuntimeShutdownReport {
pub fn cleanup_completed(&self) -> bool {
!self.drain.timed_out
&& self.mob_stop.is_ok()
&& matches!(
&self.identity_authority_release,
IdentityAuthorityReleaseOutcome::NotConfigured
| IdentityAuthorityReleaseOutcome::Released { .. }
)
&& self.module_shutdown.orphan_processes == 0
}
}
#[derive(Debug)]
pub struct UnifiedRuntimeRunReport {
pub serve_result: std::io::Result<()>,
pub shutdown: UnifiedRuntimeShutdownReport,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RediscoverReport {
pub spawned: Vec<String>,
pub edges: UnifiedRuntimeReconcileEdgesReport,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnifiedRuntimeReconcileRoutingReport {
pub router_module_loaded: bool,
pub active_members: Vec<String>,
pub added_route_keys: Vec<String>,
pub removed_route_keys: Vec<String>,
}
pub use meerkat_contracts::MobReconcileFailureWire as MobReconcileFailure;
pub use meerkat_contracts::MobReconcileReportWire as MobReconcileReport;
pub fn meerkat_reconcile_report_to_wire(
mob_id: &str,
report: meerkat_mob::runtime::reconcile::ReconcileReport,
) -> MobReconcileReport {
use meerkat_contracts::{MobSpawnReceiptWire, WireMemberRef};
let alias_of =
|id: &str| -> String { crate::member_comms_id::runtime_alias_str(id).into_owned() };
MobReconcileReport {
desired: report
.desired
.into_iter()
.map(|id| alias_of(id.as_str()))
.collect(),
retained: report
.retained
.into_iter()
.map(|id| alias_of(id.as_str()))
.collect(),
spawned: report
.spawned
.into_iter()
.map(|receipt| {
let identity_str = alias_of(receipt.agent_identity.as_str());
MobSpawnReceiptWire {
member_ref: WireMemberRef::encode(mob_id, &identity_str),
agent_identity: identity_str,
}
})
.collect(),
retired: report
.retired
.into_iter()
.map(|id| alias_of(id.as_str()))
.collect(),
failures: report
.failures
.into_iter()
.map(|failure| MobReconcileFailure {
agent_identity: alias_of(failure.agent_identity.as_str()),
stage: match failure.stage {
meerkat_mob::runtime::reconcile::ReconcileStage::Spawn => {
meerkat_contracts::WireMobReconcileStage::Spawn
}
meerkat_mob::runtime::reconcile::ReconcileStage::Retire => {
meerkat_contracts::WireMobReconcileStage::Retire
}
},
error: meerkat_contracts::WireMobError {
code: meerkat_mob::mob_error_wire_code(&failure.error),
message: failure.error.to_string(),
},
})
.collect(),
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct UnifiedRuntimeReconcileReport {
pub mob: MobReconcileReport,
pub edges: UnifiedRuntimeReconcileEdgesReport,
pub routing: UnifiedRuntimeReconcileRoutingReport,
}
#[derive(Debug)]
pub enum UnifiedRuntimeReconcileError {
Mob(MobRuntimeError),
RouteMutation(RuntimeRouteMutationError),
PartialFailure(Box<UnifiedRuntimeReconcileReport>),
}
impl Display for UnifiedRuntimeReconcileError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::Mob(err) => write!(f, "failed to reconcile mob roster: {err}"),
Self::RouteMutation(err) => {
write!(f, "failed to reconcile routing wiring: {err:?}")
}
Self::PartialFailure(report) => {
write!(
f,
"reconcile completed with {} per-identity failure(s): {:?}",
report.mob.failures.len(),
report.mob.failures
)
}
}
}
}
impl std::error::Error for UnifiedRuntimeReconcileError {}
#[derive(Debug)]
pub struct ShutdownDrainReport {
pub drained_count: usize,
pub timed_out: bool,
pub drain_duration_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CompactionPreservedHistoryFit {
Unclassified,
StillFits,
OverWindow,
}
impl From<meerkat_core::event::CompactionPreservedHistoryFit> for CompactionPreservedHistoryFit {
fn from(fit: meerkat_core::event::CompactionPreservedHistoryFit) -> Self {
use meerkat_core::event::CompactionPreservedHistoryFit as Upstream;
match fit {
Upstream::StillFits => Self::StillFits,
Upstream::OverWindow => Self::OverWindow,
_ => Self::Unclassified,
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "category", rename_all = "snake_case")]
pub enum ErrorEvent {
SpawnFailure {
member_id: String,
profile: String,
error: String,
},
ReconcileIncomplete {
failures: usize,
skipped: usize,
},
CheckpointFailure {
session_id: String,
error: String,
},
CompactionPersistenceRejected {
identity: String,
session_id: String,
error: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
preserved_history: Option<CompactionPreservedHistoryFit>,
#[serde(default, skip_serializing_if = "Option::is_none")]
attempted_entries: Option<usize>,
},
ActorLoopStalled {
probe_waited_secs: u64,
detail: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
stall_id: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
prior_resolved_stalls: Option<u64>,
},
ActorLoopRecovered {
stall_id: u64,
stalled_for_secs: u64,
},
HostLoopCrash {
member_id: String,
error: String,
},
RediscoverFailure {
error: String,
},
EventLogFlushFailure {
error: String,
},
IdentityMaterializationFailure {
identity: String,
initiator: Option<String>,
operation: String,
error: String,
},
MobStopProceededWithoutInterrupt {
waited_ms: u64,
member: Option<String>,
error: String,
},
}
impl Display for ErrorEvent {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::SpawnFailure {
member_id, error, ..
} => {
write!(f, "spawn_failure: {member_id}: {error}")
}
Self::ReconcileIncomplete { failures, skipped } => {
write!(
f,
"reconcile_incomplete: {failures} failures, {skipped} skipped"
)
}
Self::CheckpointFailure { session_id, error } => {
write!(f, "checkpoint_failure: {session_id}: {error}")
}
Self::CompactionPersistenceRejected {
identity,
session_id,
error,
..
} => {
write!(
f,
"compaction_persistence_rejected: {identity} ({session_id}): {error}"
)
}
Self::ActorLoopStalled {
probe_waited_secs,
detail,
..
} => {
write!(
f,
"actor_loop_stalled: probe unanswered for {probe_waited_secs}s: {detail}"
)
}
Self::ActorLoopRecovered {
stall_id,
stalled_for_secs,
} => {
write!(
f,
"actor_loop_recovered: stall {stall_id} resolved after {stalled_for_secs}s"
)
}
Self::HostLoopCrash { member_id, error } => {
write!(f, "host_loop_crash: {member_id}: {error}")
}
Self::RediscoverFailure { error } => {
write!(f, "rediscover_failure: {error}")
}
Self::EventLogFlushFailure { error } => {
write!(f, "event_log_flush_failure: {error}")
}
Self::IdentityMaterializationFailure {
identity,
initiator,
operation,
error,
} => {
if let Some(initiator) = initiator {
write!(
f,
"identity_materialization_failure: {identity} for {initiator} during {operation}: {error}"
)
} else {
write!(
f,
"identity_materialization_failure: {identity} during {operation}: {error}"
)
}
}
Self::MobStopProceededWithoutInterrupt {
waited_ms,
member,
error,
} => {
let subject = match member {
Some(member) => format!(" on {member}"),
None => String::new(),
};
write!(
f,
"mob_stop_proceeded_without_interrupt: teardown proceeded after {waited_ms}ms \
WITHOUT a successful interrupt{subject}; a turn may have been admitted and \
may still be running: {error}"
)
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn completed_shutdown_report() -> UnifiedRuntimeShutdownReport {
UnifiedRuntimeShutdownReport {
drain: ShutdownDrainReport {
drained_count: 1,
timed_out: false,
drain_duration_ms: 2,
},
module_shutdown: RuntimeShutdownReport {
terminated_modules: vec!["router".to_string()],
orphan_processes: 0,
},
mob_stop: Ok(()),
identity_authority_release: IdentityAuthorityReleaseOutcome::NotConfigured,
}
}
#[test]
fn shutdown_cleanup_attestation_requires_every_authority_boundary() {
let mut report = completed_shutdown_report();
assert!(report.cleanup_completed());
report.identity_authority_release =
IdentityAuthorityReleaseOutcome::Released { grant_count: 1 };
assert!(report.cleanup_completed());
let mut report = completed_shutdown_report();
report.drain.timed_out = true;
assert!(!report.cleanup_completed());
let mut report = completed_shutdown_report();
report.mob_stop = Err(MobRuntimeError::InvalidConfig(
"mob stop failed".to_string(),
));
assert!(!report.cleanup_completed());
let mut report = completed_shutdown_report();
report.identity_authority_release = IdentityAuthorityReleaseOutcome::Failed {
error: "provider release failed".to_string(),
};
assert!(!report.cleanup_completed());
let mut report = completed_shutdown_report();
report.identity_authority_release = IdentityAuthorityReleaseOutcome::SkippedMobStopFailed;
assert!(!report.cleanup_completed());
let mut report = completed_shutdown_report();
report.identity_authority_release =
IdentityAuthorityReleaseOutcome::SkippedResetCleanupFailed {
error: "superseded member retained".to_string(),
};
assert!(!report.cleanup_completed());
let mut report = completed_shutdown_report();
report.module_shutdown.orphan_processes = 1;
assert!(!report.cleanup_completed());
}
}