use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::time::Duration;
use futures::stream::{BoxStream, SelectAll, StreamExt};
use meerkat_core::comms::EventStream;
use meerkat_core::event::{AgentEvent, agent_event_type};
use meerkat_mob::{
AgentIdentity, AgentRuntimeId, AttributedEvent, FenceToken, MobError, MobHandle,
MobMemberStatus, MobState, ProfileName, SpawnMemberSpec,
};
use tokio::sync::mpsc::{Receiver, Sender};
use tokio::task::JoinHandle;
pub(crate) use self::console_events::ConsoleEventStore;
use self::mob_events::MobEventsStore;
use crate::console_aggregator::{ConsoleLogStore, InMemoryConsoleLogStore};
use crate::mob_handle_runtime::{MobBootstrapSpec, MobRuntime, MobRuntimeError};
use crate::runtime::{
InMemoryMetadataStore, MetadataScope, MobkitRuntimeHandle, PersistentMetadataStore,
RuntimeMetadataTable, RuntimeOptions, start_mobkit_runtime_with_options,
};
use crate::types::{
AgentDiscoverySpec, EventEnvelope, MobKitConfig, MobStructuralEventEnvelope, UnifiedEvent,
};
pub mod builder;
pub(crate) mod console_events;
pub mod cross_mob;
pub mod edge_reconcile;
pub mod edge_types;
pub mod event_log;
pub mod http;
pub(crate) mod implicit_delegate_retirement;
pub mod lifecycle;
pub mod mob_events;
pub mod mob_ops;
pub mod module_ops;
pub mod types;
pub use crate::identity_first::IdentityBootstrapMode;
pub use builder::UnifiedRuntimeBuilder;
pub use edge_types::{
DesiredPeerEdge, DesiredPeerEdgeError, Discovery, EdgeDiscovery, EdgeReconcileFailure,
PreSpawnContext, PreSpawnHook,
};
pub use event_log::{
EventLogConfig, EventLogError, EventLogStore, EventQuery, NullEventLogStore, PersistedEvent,
};
pub use http::DEFAULT_REFERENCE_APP_MAX_CONCURRENT_REQUESTS;
pub use mob_ops::MemberTurnAdmission;
pub use types::{
CompactionPreservedHistoryFit, ErrorEvent, IdentityAuthorityReleaseOutcome, MobStopOutcome,
RediscoverReport, ShutdownDrainReport, UnifiedRuntimeBootstrapError,
UnifiedRuntimeBuilderError, UnifiedRuntimeBuilderField, UnifiedRuntimeError,
UnifiedRuntimeReconcileEdgesReport, UnifiedRuntimeReconcileError,
UnifiedRuntimeReconcileReport, UnifiedRuntimeReconcileRoutingReport, UnifiedRuntimeRunReport,
UnifiedRuntimeShutdownReport,
};
pub type PostSpawnHook =
Arc<dyn Fn(Vec<String>) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>;
pub type PostReconcileHook = Arc<
dyn Fn(UnifiedRuntimeReconcileReport) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync,
>;
pub type ErrorHook =
Arc<dyn Fn(ErrorEvent) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync>;
type SharedErrorHook = Arc<std::sync::RwLock<Option<ErrorHook>>>;
fn current_error_hook(slot: &SharedErrorHook) -> Option<ErrorHook> {
slot.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) fn log_error_event(event: &ErrorEvent, hook_registered: bool) {
if matches!(event, ErrorEvent::ActorLoopRecovered { .. }) {
tracing::info!(
error_event = ?event,
hook_registered,
"mobkit runtime error event resolved: {event}"
);
} else {
tracing::error!(
error_event = ?event,
hook_registered,
"mobkit runtime error event: {event}"
);
}
if !hook_registered {
warn_error_hook_absent_once();
}
}
pub(crate) fn emit_error_hook_absent_notice() {
tracing::warn!(
"no error hook is registered, so runtime error events reach logs only; \
register one with UnifiedRuntimeBuilder::on_error (or \
UnifiedRuntime::set_error_hook) to route them to paging"
);
}
fn warn_error_hook_absent_once() {
static NOTICED: std::sync::Once = std::sync::Once::new();
NOTICED.call_once(emit_error_hook_absent_notice);
}
fn fire_error_hook(slot: &SharedErrorHook, event: ErrorEvent) {
let hook = current_error_hook(slot);
log_error_event(&event, hook.is_some());
if let Some(hook) = hook {
tokio::spawn(async move {
let () = hook(event).await;
});
}
}
const ROSTER_ROUTE_PREFIX: &str = "mob.member.";
const ROSTER_ROUTE_CHANNEL: &str = "notification";
const ROSTER_ROUTE_SINK: &str = "mob_member";
const ROSTER_ROUTE_TARGET_MODULE: &str = "delivery";
const DEFAULT_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
pub fn discovery_spec_to_spawn_spec(spec: &AgentDiscoverySpec) -> SpawnMemberSpec {
let resume_session_id = spec
.resume_session_id
.as_deref()
.and_then(|s| meerkat_core::types::SessionId::parse(s).ok());
let additional_instructions = if spec.additional_instructions.is_empty() {
None
} else {
Some(spec.additional_instructions.clone())
};
let mut spawn = SpawnMemberSpec::new(
meerkat_mob::ProfileName::from(spec.profile.as_str()),
meerkat_mob::ids::AgentIdentity::from(spec.meerkat_id.as_str()),
);
if let Some(context) = spec.context.clone() {
spawn = spawn.with_context(context);
}
if let Some(labels) = spec.labels.clone() {
spawn = spawn.with_labels(labels);
}
if let Some(sid) = resume_session_id {
spawn = spawn.with_resume_bridge_session_id(sid);
}
if let Some(instructions) = additional_instructions {
spawn = spawn.with_additional_instructions(instructions);
}
spawn
}
pub struct UnifiedRuntime {
mob_runtime: MobRuntime,
post_spawn_hook: Option<PostSpawnHook>,
post_reconcile_hook: Option<PostReconcileHook>,
error_hook: SharedErrorHook,
drain_timeout: Duration,
discovery: Option<Box<dyn Discovery>>,
edge_discovery: Option<Arc<dyn EdgeDiscovery>>,
module_runtime: Arc<tokio::sync::Mutex<MobkitRuntimeHandle>>,
managed_dynamic_edges: Arc<tokio::sync::RwLock<BTreeSet<(String, String)>>>,
shutting_down: AtomicBool,
mob_event_ingress: tokio::sync::Mutex<Option<MobEventIngress>>,
bootstrap_edges_report: tokio::sync::RwLock<Option<UnifiedRuntimeReconcileEdgesReport>>,
event_log: Option<event_log::EventLogHandle>,
console_log_store: Arc<dyn ConsoleLogStore>,
console_events: ConsoleEventStore,
mob_events: MobEventsStore,
mob_events_subscriber_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
actor_loop_probe_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
implicit_delegate_retirement_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
implicit_delegate_identity_runtime:
Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
identity_lease_renewal_task:
tokio::sync::Mutex<Option<crate::identity_first::runtime::TrackedLeaseRenewalTask>>,
identity_continuity_repair_task:
tokio::sync::Mutex<Option<crate::identity_first::runtime::TrackedContinuityRepairTask>>,
agent_memory_observer_task:
tokio::sync::Mutex<Option<crate::memory::taint::TaintObserverGuard>>,
agent_memory_steward_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
contact_directory: Option<crate::contact_directory::ContactDirectory>,
peer_mob_handles: tokio::sync::RwLock<BTreeMap<String, cross_mob::PeerMobAuthority>>,
cross_mob_control_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
cross_mob_control_advertised: std::sync::RwLock<Option<String>>,
gateway_peer_keys: crate::runtime::cross_mob_control::ControlSignerSlot,
remote_host_facts: crate::runtime::cross_mob_control::HostFactsProviderSlot,
remote_host_lifecycle:
std::sync::RwLock<Option<Arc<crate::runtime::remote_host::RemoteHostLifecycle>>>,
remote_host_lifecycle_error:
std::sync::RwLock<Option<crate::runtime::remote_host::HostPairingError>>,
remote_host_reconnect_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
session_bridge: Option<Arc<dyn crate::identity_first::bridge::SessionBridge>>,
identity_first_context: Option<Arc<crate::identity_first::IdentityFirstRuntimeContext>>,
access_controller: Option<crate::access::AccessController>,
topology_controller: crate::topology_control::TopologyController,
memory_panel_store:
std::sync::RwLock<Option<Arc<dyn crate::memory::capabilities::MemoryPanelStore>>>,
job_health_projection: Arc<std::sync::RwLock<Option<serde_json::Value>>>,
workgraph_service: Option<meerkat::WorkGraphService>,
workgraph_fact_hub: Option<crate::workgraph_events::WorkGraphFactHub>,
workgraph_fact_tail_task: tokio::sync::Mutex<WorkGraphFactTailTask>,
console_identity_roster:
std::sync::RwLock<Option<Arc<crate::identity_first::MutableRosterProvider>>>,
console_operator_resolver: std::sync::RwLock<
Option<Arc<crate::memory::coordinator::ConsolePrincipalOperatorResolver>>,
>,
metadata_table: Arc<RuntimeMetadataTable>,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
}
enum MobEventIngress {
Forwarder(MobEventForwarder),
}
struct MobEventForwarder {
event_rx: Receiver<ForwardedMemberEvent>,
task: JoinHandle<()>,
identity_stream_health_task: JoinHandle<()>,
}
struct ForwardedMemberEvent {
envelope: EventEnvelope<UnifiedEvent>,
alert: Option<ErrorEvent>,
}
struct WorkGraphFactTailTask(Option<JoinHandle<()>>);
impl WorkGraphFactTailTask {
fn take(&mut self) -> Option<JoinHandle<()>> {
self.0.take()
}
}
impl Drop for WorkGraphFactTailTask {
fn drop(&mut self) {
if let Some(task) = self.0.take() {
task.abort();
}
}
}
impl UnifiedRuntime {
pub fn builder() -> UnifiedRuntimeBuilder {
UnifiedRuntimeBuilder::default()
}
#[allow(
unknown_lints,
clippy::unused_async_trait_impl,
reason = "preserve the async construction seam used by runtime bootstrap"
)]
pub(crate) async fn from_parts(
mob_runtime: MobRuntime,
module_runtime: MobkitRuntimeHandle,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
) -> Self {
let metadata_table = Arc::new(RuntimeMetadataTable::new());
let mob_events_store = MobEventsStore::new().with_metadata_table(metadata_table.clone());
let identity_runtime_authority = Arc::new(std::sync::RwLock::new(None));
let mob_event_ingress = Some(Self::create_event_ingress(
mob_runtime.handle(),
mob_runtime.agent_mob_mcp_state(),
mob_events_store.clone(),
Arc::clone(&identity_runtime_authority),
));
let mob_events_task = Self::spawn_mob_events_subscriber(
mob_runtime.handle(),
mob_events_store.clone(),
persistent_metadata.clone(),
);
let error_hook: SharedErrorHook = Arc::new(std::sync::RwLock::new(None));
let actor_loop_probe_task =
Self::spawn_actor_loop_probe(mob_runtime.handle(), Arc::clone(&error_hook));
let console_events = ConsoleEventStore::new();
mob_runtime.install_console_spawn_sink(crate::console_spawn::ConsoleSpawnSink::new(
console_events.clone(),
));
let workgraph_service = mob_runtime.workgraph_service();
let (workgraph_fact_hub, workgraph_fact_tail_task) =
if let Some(service) = workgraph_service.clone() {
let hub = crate::workgraph_events::WorkGraphFactHub::new();
let task = crate::workgraph_events::spawn_workgraph_fact_tail(
service,
hub.clone(),
crate::workgraph_events::WorkGraphFactTailOptions::default(),
);
(Some(hub), Some(task))
} else {
(None, None)
};
let definition_edge_discovery =
edge_reconcile::DefinitionWiringEdgeDiscovery::from_definition(
mob_runtime.handle().definition(),
)
.map(|policy| Arc::new(policy) as Arc<dyn EdgeDiscovery>);
Self {
mob_runtime,
post_spawn_hook: None,
post_reconcile_hook: None,
error_hook,
drain_timeout: DEFAULT_DRAIN_TIMEOUT,
discovery: None,
edge_discovery: definition_edge_discovery,
module_runtime: Arc::new(tokio::sync::Mutex::new(module_runtime)),
managed_dynamic_edges: Arc::new(tokio::sync::RwLock::new(BTreeSet::new())),
shutting_down: AtomicBool::new(false),
mob_event_ingress: tokio::sync::Mutex::new(mob_event_ingress),
bootstrap_edges_report: tokio::sync::RwLock::new(None),
event_log: None,
console_log_store: Arc::new(InMemoryConsoleLogStore::new()),
console_events,
mob_events: mob_events_store,
mob_events_subscriber_task: tokio::sync::Mutex::new(mob_events_task),
actor_loop_probe_task: tokio::sync::Mutex::new(actor_loop_probe_task),
implicit_delegate_retirement_task: tokio::sync::Mutex::new(None),
implicit_delegate_identity_runtime: identity_runtime_authority,
identity_lease_renewal_task: tokio::sync::Mutex::new(None),
identity_continuity_repair_task: tokio::sync::Mutex::new(None),
agent_memory_observer_task: tokio::sync::Mutex::new(None),
agent_memory_steward_task: tokio::sync::Mutex::new(None),
contact_directory: None,
peer_mob_handles: tokio::sync::RwLock::new(BTreeMap::new()),
cross_mob_control_task: tokio::sync::Mutex::new(None),
cross_mob_control_advertised: std::sync::RwLock::new(None),
gateway_peer_keys: crate::runtime::cross_mob_control::unsigned_control_signer(),
remote_host_facts: crate::runtime::cross_mob_control::empty_host_facts_provider(),
remote_host_lifecycle: std::sync::RwLock::new(None),
remote_host_lifecycle_error: std::sync::RwLock::new(None),
remote_host_reconnect_task: tokio::sync::Mutex::new(None),
session_bridge: None,
identity_first_context: None,
access_controller: None,
topology_controller: crate::topology_control::TopologyController::default(),
memory_panel_store: std::sync::RwLock::new(None),
job_health_projection: Arc::new(std::sync::RwLock::new(None)),
workgraph_service,
workgraph_fact_hub,
workgraph_fact_tail_task: tokio::sync::Mutex::new(WorkGraphFactTailTask(
workgraph_fact_tail_task,
)),
console_identity_roster: std::sync::RwLock::new(None),
console_operator_resolver: std::sync::RwLock::new(None),
metadata_table,
persistent_metadata,
}
}
fn spawn_mob_events_subscriber(
handle: MobHandle,
store: MobEventsStore,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
) -> Option<JoinHandle<()>> {
let runtime_handle = tokio::runtime::Handle::try_current().ok()?;
Some(runtime_handle.spawn(run_mob_events_subscription(
handle,
store,
persistent_metadata,
)))
}
fn spawn_actor_loop_probe(
handle: MobHandle,
error_hook: SharedErrorHook,
) -> Option<JoinHandle<()>> {
let runtime_handle = tokio::runtime::Handle::try_current().ok()?;
Some(runtime_handle.spawn(run_actor_loop_probe(
move || {
let handle = handle.clone();
async move { handle.status().await }
},
error_hook,
actor_loop_probe_interval(),
actor_loop_probe_budget(),
)))
}
pub async fn bootstrap(
mob_spec: MobBootstrapSpec,
module_config: MobKitConfig,
timeout: Duration,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
Box::pin(Self::bootstrap_with_options(
mob_spec,
module_config,
Vec::new(),
timeout,
RuntimeOptions::default(),
Arc::new(InMemoryMetadataStore::new()),
))
.await
}
pub async fn bootstrap_with_options(
mob_spec: MobBootstrapSpec,
module_config: MobKitConfig,
module_agent_events: Vec<EventEnvelope<UnifiedEvent>>,
timeout: Duration,
options: RuntimeOptions,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
Self::bootstrap_with_options_and_topology(
mob_spec,
module_config,
module_agent_events,
timeout,
options,
persistent_metadata,
crate::topology_control::TopologyBootstrapConfig::default(),
)
.await
}
pub async fn bootstrap_with_topology(
mob_spec: MobBootstrapSpec,
module_config: MobKitConfig,
timeout: Duration,
topology: crate::topology_control::TopologyBootstrapConfig,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
Self::bootstrap_with_options_and_topology(
mob_spec,
module_config,
Vec::new(),
timeout,
RuntimeOptions::default(),
Arc::new(InMemoryMetadataStore::new()),
topology,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn bootstrap_with_options_and_topology(
mob_spec: MobBootstrapSpec,
module_config: MobKitConfig,
module_agent_events: Vec<EventEnvelope<UnifiedEvent>>,
timeout: Duration,
options: RuntimeOptions,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
topology: crate::topology_control::TopologyBootstrapConfig,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
let topology_authority = mob_spec.definition.id.to_string();
let topology_controller = match topology.state_path {
Some(path) => {
crate::topology_control::TopologyController::load_or_default(topology.policy, path)
}
None => crate::topology_control::TopologyController::new(topology.policy),
}
.map_err(|error| UnifiedRuntimeBootstrapError::Topology(error.to_string()))?;
topology_controller
.bind_authority(topology_authority)
.await
.map_err(|error| UnifiedRuntimeBootstrapError::Topology(error.to_string()))?;
let mob_runtime = MobRuntime::bootstrap(mob_spec)
.await
.map_err(UnifiedRuntimeBootstrapError::Mob)?;
let runtime_options = options.clone();
let module_start_result = std::thread::spawn(move || {
start_mobkit_runtime_with_options(module_config, module_agent_events, timeout, options)
})
.join();
match module_start_result {
Ok(Ok(module_runtime)) => {
let mut runtime =
Self::from_parts(mob_runtime, module_runtime, persistent_metadata).await;
runtime.topology_controller = topology_controller;
runtime
.configure_implicit_delegate_retirement(&runtime_options)
.await;
if runtime.edge_discovery.is_some()
|| runtime.topology_controller.revision().await > 0
|| runtime.topology_controller.has_pending().await
{
let report = runtime.reconcile_edges().await;
*runtime.bootstrap_edges_report.write().await = Some(report);
}
Ok(runtime)
}
Ok(Err(error)) => {
let startup_error = UnifiedRuntimeBootstrapError::Module(error);
Self::rollback_mob_runtime(mob_runtime, startup_error).await
}
Err(_) => {
let startup_error = UnifiedRuntimeBootstrapError::ModuleStartupThreadPanicked;
Self::rollback_mob_runtime(mob_runtime, startup_error).await
}
}
}
pub async fn bootstrap_edges_report(&self) -> Option<UnifiedRuntimeReconcileEdgesReport> {
self.bootstrap_edges_report.read().await.clone()
}
pub fn set_error_hook(&mut self, hook: ErrorHook) {
*self
.error_hook
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(hook.clone());
if let Some(identity_runtime) = self.identity_runtime() {
identity_runtime.set_error_hook(Some(hook));
}
}
pub fn start_event_log(&mut self, config: EventLogConfig) {
let handle = event_log::start_event_log(config, current_error_hook(&self.error_hook));
self.event_log = Some(handle);
}
pub(crate) fn console_events(&self) -> ConsoleEventStore {
self.console_events.clone()
}
pub fn memory_event_sink(&self) -> Arc<dyn crate::memory::events::MemoryEventSink> {
Arc::new(ConsoleMemoryEventSink {
store: self.console_events(),
handle: tokio::runtime::Handle::current(),
})
}
pub async fn register_gating_resolution_observer(
&self,
observer: Arc<dyn crate::runtime::GatingResolutionObserver>,
) {
self.module_runtime
.lock()
.await
.register_gating_resolution_observer(observer);
}
pub(crate) fn mob_events_store(&self) -> MobEventsStore {
self.mob_events.clone()
}
pub fn binary_blob_store(&self) -> Option<Arc<dyn crate::blob_store::BinaryBlobStore>> {
self.mob_runtime.binary_blob_store()
}
pub fn set_job_health_projection(&self, projection: Option<serde_json::Value>) {
*self
.job_health_projection
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = projection;
}
pub fn job_health_projection(&self) -> Option<serde_json::Value> {
self.job_health_projection
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) fn module_runtime_handle(&self) -> Arc<tokio::sync::Mutex<MobkitRuntimeHandle>> {
Arc::clone(&self.module_runtime)
}
pub(crate) fn mobpack_runtime_catalog_state_snapshot(
&self,
) -> crate::mobpack::MobpackRuntimeCatalogState {
let loaded_modules = self
.module_runtime
.try_lock()
.map(|runtime| runtime.loaded_modules())
.unwrap_or_default();
let has_peer_mob_handles = self
.peer_mob_handles
.try_read()
.map(|handles| !handles.is_empty())
.unwrap_or(false);
let mut runtime_methods = vec![
"mobkit/capabilities".to_string(),
"mobkit/models/catalog".to_string(),
"mobkit/spawn_member".to_string(),
"mobkit/list_members".to_string(),
"mobkit/get_member".to_string(),
"mobkit/run_flow".to_string(),
"mobkit/list_flows".to_string(),
"mobkit/list_runs".to_string(),
];
runtime_methods.extend(
crate::rpc::MOBPACK_AUTHORING_METHODS
.iter()
.map(std::string::ToString::to_string),
);
if self.has_contact_directory() {
runtime_methods.push("mobkit/cross_mob/directory".to_string());
}
if (has_peer_mob_handles && self.has_inproc_contacts()) || self.has_remote_contacts() {
runtime_methods.extend([
"mobkit/cross_mob/wire".to_string(),
"mobkit/cross_mob/unwire".to_string(),
"mobkit/cross_mob/send".to_string(),
]);
}
crate::mobpack::MobpackRuntimeCatalogState {
loaded_modules,
runtime_methods,
has_contact_directory: self.has_contact_directory(),
has_peer_mob_handles,
has_inproc_contacts: self.has_inproc_contacts(),
runtime_flow_rows: crate::mobpack::runtime_flow_registry_rows_from_definition(
self.mob_handle().definition(),
),
runtime_agent_definition_sources:
crate::mobpack::runtime_agent_definition_sources_from_definition(
self.mob_handle().definition(),
),
runtime_skill_realms: crate::mobpack::runtime_skill_realms_from_definition(
self.mob_handle().definition(),
),
}
}
pub fn session_bridge(&self) -> Option<&Arc<dyn crate::identity_first::bridge::SessionBridge>> {
self.session_bridge.as_ref()
}
pub fn identity_first_context(
&self,
) -> Option<&Arc<crate::identity_first::IdentityFirstRuntimeContext>> {
self.identity_first_context.as_ref()
}
pub fn identity_runtime(&self) -> Option<&Arc<crate::identity_first::IdentityRuntime>> {
self.identity_first_context.as_ref().map(|ctx| &ctx.runtime)
}
pub async fn remember_agent_memory(
&self,
realm: &str,
identity: &crate::identity_first::AgentIdentity,
memory: crate::identity_first::NewAgentMemory,
) -> Result<crate::identity_first::AgentMemoryRecord, crate::identity_first::AgentMemoryError>
{
let runtime = self.identity_runtime().ok_or_else(|| {
crate::identity_first::AgentMemoryError::InvalidConfig(
"identity-first runtime is not configured".to_string(),
)
})?;
runtime.remember_agent_memory(realm, identity, memory).await
}
pub async fn recall_agent_memory(
&self,
request: crate::identity_first::AgentMemoryRecallRequest,
) -> Result<
Vec<crate::identity_first::AgentMemoryRecord>,
crate::identity_first::AgentMemoryError,
> {
let runtime = self.identity_runtime().ok_or_else(|| {
crate::identity_first::AgentMemoryError::InvalidConfig(
"identity-first runtime is not configured".to_string(),
)
})?;
runtime.recall_agent_memory(request).await
}
pub async fn forget_agent_memory(
&self,
realm: &str,
identity: &crate::identity_first::AgentIdentity,
memory_id: &str,
) -> Result<
crate::identity_first::AgentMemoryForgetResult,
crate::identity_first::AgentMemoryError,
> {
let runtime = self.identity_runtime().ok_or_else(|| {
crate::identity_first::AgentMemoryError::InvalidConfig(
"identity-first runtime is not configured".to_string(),
)
})?;
runtime
.forget_agent_memory(realm, identity, memory_id)
.await
}
pub fn attach_identity_first_context(
&mut self,
context: Arc<crate::identity_first::IdentityFirstRuntimeContext>,
) {
self.install_identity_first_context_authority(context);
self.start_identity_first_supervisors();
}
pub async fn install_and_bootstrap_identity_first_context(
&mut self,
context: Arc<crate::identity_first::IdentityFirstRuntimeContext>,
roster: &[crate::identity_first::DurableAgentSpec],
) -> Result<crate::identity_first::RestoreFlowResult, crate::identity_first::IdentityRuntimeError>
{
self.install_identity_first_context_authority(Arc::clone(&context));
match context.bootstrap_roster(roster).await {
Ok(result) => {
self.start_identity_first_supervisors();
Ok(result)
}
Err(error) => {
self.shutdown().await;
Err(error)
}
}
}
fn install_identity_first_context_authority(
&mut self,
context: Arc<crate::identity_first::IdentityFirstRuntimeContext>,
) {
self.install_identity_first_flow_target_provisioner(&context.runtime);
self.mob_runtime
.install_identity_runtime_authority(Arc::clone(&context.runtime));
*self
.implicit_delegate_identity_runtime
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(Arc::clone(&context.runtime));
self.identity_first_context = Some(context);
}
pub(crate) fn install_identity_first_flow_target_provisioner(
&self,
runtime: &Arc<crate::identity_first::IdentityRuntime>,
) {
let identity_runtime = Arc::downgrade(runtime);
self.mob_runtime
.handle()
.install_flow_target_provisioner(Arc::new(move || {
let identity_runtime = identity_runtime.clone();
Box::pin(async move {
let runtime = identity_runtime.upgrade().ok_or_else(|| {
MobError::Internal(
"identity-first flow provisioner is no longer available".to_string(),
)
})?;
runtime
.materialize_all_required_tracked()
.await
.map(|_| ())
.map_err(|error| {
MobError::Internal(format!(
"identity-first flow materialization failed: {error}"
))
})
})
}));
}
fn start_identity_first_supervisors(&mut self) {
let Some(context) = self.identity_first_context.clone() else {
return;
};
let lease_task = context.runtime.clone().spawn_tracked_lease_renewal_task();
if let Some(previous) = self
.identity_lease_renewal_task
.get_mut()
.replace(lease_task)
{
previous.cancel();
tokio::spawn(previous.cancel_and_join());
}
let repair_task = context.spawn_tracked_broken_identity_repair_task(Default::default());
if let Some(previous) = self
.identity_continuity_repair_task
.get_mut()
.replace(repair_task)
{
previous.cancel();
tokio::spawn(previous.cancel_and_join());
}
}
pub async fn refresh_desired_topology(
&self,
) -> Result<
Option<crate::identity_first::RestoreFlowResult>,
crate::identity_first::IdentityRuntimeError,
> {
match self.identity_first_context.as_ref() {
Some(ctx) => ctx.refresh_desired_topology_tracked().await.map(Some),
None => Ok(None),
}
}
pub async fn materialize_identity_first_for_flow(
&self,
) -> Result<
Vec<crate::identity_first::ContinuityRecord>,
crate::identity_first::IdentityRuntimeError,
> {
match self.identity_runtime() {
Some(runtime) => runtime.materialize_all_required_tracked().await,
None => Ok(Vec::new()),
}
}
pub fn metadata_table(&self) -> &Arc<RuntimeMetadataTable> {
&self.metadata_table
}
pub fn set_access_controller(&mut self, controller: crate::access::AccessController) {
self.access_controller = Some(controller);
}
pub fn set_console_identity_roster(
&self,
roster: Arc<crate::identity_first::MutableRosterProvider>,
) {
*self
.console_identity_roster
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(roster);
}
pub fn console_identity_roster(
&self,
) -> Option<Arc<crate::identity_first::MutableRosterProvider>> {
self.console_identity_roster
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn set_memory_panel_store(
&self,
store: Arc<dyn crate::memory::capabilities::MemoryPanelStore>,
) {
*self
.memory_panel_store
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(store);
}
pub fn memory_panel_store(
&self,
) -> Option<Arc<dyn crate::memory::capabilities::MemoryPanelStore>> {
self.memory_panel_store
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn workgraph_service(&self) -> Option<meerkat::WorkGraphService> {
self.workgraph_service.clone()
}
pub fn workgraph_fact_hub(&self) -> Option<crate::workgraph_events::WorkGraphFactHub> {
self.workgraph_fact_hub.clone()
}
pub fn resolved_storage(&self) -> Option<crate::storage_health::ResolvedStorageSummary> {
self.mob_runtime.resolved_storage()
}
pub(crate) fn workgraph_admission(
&self,
) -> std::sync::Arc<crate::workgraph_admission::WorkGraphAdmission> {
self.mob_runtime.workgraph_admission()
}
pub fn set_console_operator_resolver(
&self,
resolver: Arc<crate::memory::coordinator::ConsolePrincipalOperatorResolver>,
) {
*self
.console_operator_resolver
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(resolver);
}
pub fn console_operator_resolver(
&self,
) -> Option<Arc<crate::memory::coordinator::ConsolePrincipalOperatorResolver>> {
self.console_operator_resolver
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn access_controller(&self) -> Option<&crate::access::AccessController> {
self.access_controller.as_ref()
}
pub fn topology_controller(&self) -> &crate::topology_control::TopologyController {
&self.topology_controller
}
pub fn topology_runtime_handle(&self) -> crate::topology_control::TopologyRuntimeHandle {
crate::topology_control::TopologyRuntimeHandle::new(
self.mob_handle(),
self.edge_discovery.clone(),
Arc::clone(&self.managed_dynamic_edges),
self.topology_controller.clone(),
self.identity_first_context.clone(),
)
}
pub fn set_topology_control_policy(
&self,
policy: crate::topology_control::TopologyControlPolicy,
) -> Result<(), crate::topology_control::TopologyControlError> {
self.topology_controller.set_policy(policy)
}
pub fn persistent_metadata(&self) -> &Arc<dyn PersistentMetadataStore> {
&self.persistent_metadata
}
pub async fn set_mob_labels(&self, labels: BTreeMap<String, String>) {
self.metadata_table
.set_labels(MetadataScope::Mob(self.mob_id()), labels)
.await;
}
pub async fn get_mob_labels(&self) -> BTreeMap<String, String> {
self.metadata_table
.get_labels(&MetadataScope::Mob(self.mob_id()))
.await
}
pub async fn delete_mob_labels(&self) {
let _ = self
.metadata_table
.delete_labels(&MetadataScope::Mob(self.mob_id()))
.await;
}
pub async fn set_run_labels(&self, run_id: &str, labels: BTreeMap<String, String>) {
self.metadata_table
.set_labels(
MetadataScope::Run(self.mob_id(), run_id.to_string()),
labels,
)
.await;
}
pub async fn get_run_labels(&self, run_id: &str) -> BTreeMap<String, String> {
self.metadata_table
.get_labels(&MetadataScope::Run(self.mob_id(), run_id.to_string()))
.await
}
pub async fn delete_run_labels(&self, run_id: &str) {
let _ = self
.metadata_table
.delete_labels(&MetadataScope::Run(self.mob_id(), run_id.to_string()))
.await;
}
pub fn event_log_store(&self) -> Option<std::sync::Arc<dyn event_log::EventLogStore>> {
self.event_log
.as_ref()
.map(event_log::EventLogHandle::store)
}
pub fn console_log_store(&self) -> Arc<dyn ConsoleLogStore> {
self.console_log_store.clone()
}
pub fn set_console_log_store(&mut self, store: Arc<dyn ConsoleLogStore>) {
self.console_log_store = store;
}
pub async fn query_mob_events(
&self,
query: &EventQuery,
) -> Result<Vec<MobStructuralEventEnvelope>, mob_events::MobEventsQueryError> {
let events = self.mob_runtime.handle().events();
mob_events::query_ledger_with_filter(&events, &self.mob_events, query).await
}
pub fn subscribe_mob_events(
&self,
) -> tokio::sync::broadcast::Receiver<MobStructuralEventEnvelope> {
self.mob_events.subscribe()
}
pub(crate) fn ingest_event(&self, event: &EventEnvelope<UnifiedEvent>) {
if let Some(ref log) = self.event_log {
log.ingest(event.clone());
}
}
pub(crate) async fn record_console_lifecycle(
&self,
identity: &str,
event_type: &str,
data: serde_json::Value,
) {
self.console_events
.record_lifecycle(identity, event_type, data)
.await;
}
pub async fn reserve_identity_interaction(
&self,
identity: &str,
runtime_member_id: Option<&str>,
interaction_id: &str,
origin: &str,
content: serde_json::Value,
) -> Result<(), &'static str> {
self.console_events
.reserve_interaction_value(identity, runtime_member_id, interaction_id, origin, content)
.await
}
pub(crate) async fn project_console_event_from_unified(
&self,
event: &EventEnvelope<UnifiedEvent>,
) {
self.console_events.project_unified_event(event).await;
}
pub(crate) fn fire_error(&self, event: ErrorEvent) {
fire_error_hook(&self.error_hook, event);
}
fn create_event_ingress(
mob_handle: MobHandle,
agent_mob_mcp_state: Option<Arc<meerkat_mob_mcp::MobMcpState>>,
mob_events: MobEventsStore,
identity_runtime: Arc<
std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>,
>,
) -> MobEventIngress {
let (event_tx, event_rx) = tokio::sync::mpsc::channel(256);
let identity_stream_health_task = tokio::spawn(run_identity_stream_health_monitor(
mob_handle.clone(),
agent_mob_mcp_state.clone(),
identity_runtime,
));
let task = tokio::spawn(run_resilient_mob_agent_event_forwarder(
mob_handle,
agent_mob_mcp_state,
event_tx,
mob_events,
));
MobEventIngress::Forwarder(MobEventForwarder {
event_rx,
task,
identity_stream_health_task,
})
}
#[cfg(test)]
async fn install_test_event_ingress(&self) -> Sender<ForwardedMemberEvent> {
let (event_tx, event_rx) = tokio::sync::mpsc::channel(16);
let replaced = self
.mob_event_ingress
.lock()
.await
.replace(MobEventIngress::Forwarder(MobEventForwarder {
event_rx,
task: tokio::spawn(async {}),
identity_stream_health_task: tokio::spawn(async {}),
}));
if let Some(MobEventIngress::Forwarder(forwarder)) = replaced {
forwarder.task.abort();
forwarder.identity_stream_health_task.abort();
}
event_tx
}
async fn rollback_mob_runtime(
mob_runtime: MobRuntime,
startup_error: UnifiedRuntimeBootstrapError,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
match mob_runtime.handle().stop().await {
Ok(()) => Err(startup_error),
Err(err) => Err(UnifiedRuntimeBootstrapError::ModuleStartupRollbackFailed {
startup_error: Box::new(startup_error),
rollback_error: MobRuntimeError::from(err),
}),
}
}
}
type TaggedAgentEvent = (
AgentRuntimeId,
FenceToken,
ProfileName,
meerkat_core::event::EventEnvelope<AgentEvent>,
Option<Arc<str>>,
);
enum ForwardedAgentEvent {
Event(Box<TaggedAgentEvent>),
Closed(TrackedAgentEventStream),
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
struct TrackedAgentEventStream {
mob_id: String,
durable_identity: Option<String>,
member_identity: AgentIdentity,
runtime_id: AgentRuntimeId,
identity_fencing_token: Option<u64>,
fence_token: FenceToken,
}
type TaggedAgentEventStream = BoxStream<'static, ForwardedAgentEvent>;
struct SubscribeBackoff {
next_attempt: tokio::time::Instant,
consecutive_failures: u32,
}
const SUBSCRIBE_BACKOFF_BASE: Duration = Duration::from_millis(250);
const SUBSCRIBE_BACKOFF_MAX: Duration = Duration::from_secs(30);
const PERMANENT_STREAM_FAILURE_THRESHOLD: u32 = 4;
const RECONCILE_SAFETY_INTERVAL: Duration = Duration::from_secs(30);
struct ReconcileCadence {
machine_watchers: BTreeMap<String, meerkat_mob::MobMachineStateChanges>,
mob_set_changes: Option<tokio::sync::watch::Receiver<u64>>,
next_safety_deadline: tokio::time::Instant,
}
impl ReconcileCadence {
fn new(agent_mob_mcp_state: &Option<Arc<meerkat_mob_mcp::MobMcpState>>) -> Self {
Self {
machine_watchers: BTreeMap::new(),
mob_set_changes: agent_mob_mcp_state
.as_ref()
.map(|state| state.mob_set_changes()),
next_safety_deadline: tokio::time::Instant::now() + RECONCILE_SAFETY_INTERVAL,
}
}
fn rebind(&mut self, handles: &[MobHandle]) {
let mut next = BTreeMap::new();
for handle in handles {
let key = handle.mob_id().to_string();
let watcher = self
.machine_watchers
.remove(&key)
.unwrap_or_else(|| handle.machine_state_changes());
if !watcher.is_closed() {
next.insert(key, watcher);
}
}
self.machine_watchers = next;
self.next_safety_deadline = tokio::time::Instant::now() + RECONCILE_SAFETY_INTERVAL;
}
async fn wait(&mut self, next_backoff_attempt: Option<tokio::time::Instant>) {
let now = tokio::time::Instant::now();
let mut deadline = self.next_safety_deadline;
if let Some(attempt) = next_backoff_attempt {
deadline = deadline.min(attempt.max(now));
}
let Self {
machine_watchers,
mob_set_changes,
..
} = self;
let machine_change = async {
if machine_watchers.is_empty() {
std::future::pending::<()>().await;
return;
}
let keys: Vec<String> = machine_watchers.keys().cloned().collect();
let closed_key = {
let futures: Vec<_> = machine_watchers
.values_mut()
.map(|watcher| Box::pin(watcher.changed()))
.collect();
let (result, index, rest) = futures::future::select_all(futures).await;
drop(rest);
result.is_err().then(|| keys[index].clone())
};
if let Some(key) = closed_key {
machine_watchers.remove(&key);
}
};
let mob_set_change = async {
match mob_set_changes.as_mut() {
Some(rx) => rx.changed().await,
None => std::future::pending().await,
}
};
let mob_set_closed = tokio::select! {
() = machine_change => false,
result = mob_set_change => result.is_err(),
() = tokio::time::sleep_until(deadline) => false,
};
if mob_set_closed {
self.mob_set_changes = None;
}
}
}
fn earliest_backoff_attempt(
subscribe_failures: &HashMap<TrackedAgentEventStream, SubscribeBackoff>,
) -> Option<tokio::time::Instant> {
subscribe_failures
.values()
.map(|backoff| backoff.next_attempt)
.min()
}
fn subscribe_backoff_delay(consecutive_failures: u32) -> Duration {
SUBSCRIBE_BACKOFF_BASE
.saturating_mul(1u32 << consecutive_failures.min(7))
.min(SUBSCRIBE_BACKOFF_MAX)
}
fn forwarder_should_subscribe(status: MobMemberStatus) -> bool {
matches!(status, MobMemberStatus::Active)
}
fn durable_identity_label(labels: &BTreeMap<String, String>) -> Option<String> {
labels.get("agent_identity").cloned()
}
async fn current_identity_fencing_token(
primary_mob_id: &str,
mob_id: &str,
durable_identity: Option<&str>,
identity_runtime: Option<
&Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
>,
) -> Option<u64> {
if mob_id != primary_mob_id {
return None;
}
let durable_identity = durable_identity?;
let identity_runtime = identity_runtime?;
let authority = identity_runtime
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()?;
let identity = crate::identity_first::AgentIdentity::parse(durable_identity).ok()?;
authority
.status(&identity)
.await
.ok()?
.lease
.map(|lease| lease.fencing_token.get())
}
async fn record_identity_turn_completion(
identity_runtime: &Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
durable_identity: Option<&str>,
envelope: &meerkat_core::event::EventEnvelope<AgentEvent>,
) {
if !matches!(envelope.payload, AgentEvent::RunCompleted { .. }) {
return;
}
let Some(durable_identity) = durable_identity else {
return;
};
let authority = identity_runtime
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let Some(authority) = authority else {
return;
};
let identity = match crate::identity_first::AgentIdentity::parse(durable_identity) {
Ok(identity) => identity,
Err(error) => {
tracing::debug!(
identity = %durable_identity,
error = %error,
"mobkit identity health monitor: run completion carried an unparseable durable identity"
);
return;
}
};
authority.record_turn_completed(&identity).await;
}
async fn trigger_identity_stream_repair(
primary_mob_id: &str,
tracked_key: &TrackedAgentEventStream,
identity_runtime: &Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
detail: &str,
) {
if tracked_key.mob_id != primary_mob_id {
return;
}
let Some(durable_identity) = tracked_key.durable_identity.as_deref() else {
return;
};
let Some(identity_fencing_token) = tracked_key.identity_fencing_token else {
return;
};
let runtime_alias =
crate::member_comms_id::runtime_alias_str(tracked_key.runtime_id.identity.as_str())
.into_owned();
let authority = identity_runtime
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let Some(authority) = authority else {
return;
};
let identity = match crate::identity_first::AgentIdentity::parse(durable_identity) {
Ok(identity) => identity,
Err(error) => {
tracing::warn!(
identity = %durable_identity,
error = %error,
"mobkit agent event forwarder: roster identity cannot be mapped to identity authority"
);
return;
}
};
if let Err(error) = authority
.mark_active_runtime_broken(&identity, &runtime_alias, identity_fencing_token, detail)
.await
{
tracing::warn!(
identity = %identity,
runtime_id = %tracked_key.runtime_id,
error = %error,
"mobkit agent event forwarder: failed to trigger identity repair after permanent stream loss"
);
}
}
async fn run_resilient_mob_agent_event_forwarder(
handle: MobHandle,
agent_mob_mcp_state: Option<Arc<meerkat_mob_mcp::MobMcpState>>,
event_tx: Sender<ForwardedMemberEvent>,
mob_events: MobEventsStore,
) {
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
let mut tracked = HashSet::new();
let mut subscribe_failures: HashMap<TrackedAgentEventStream, SubscribeBackoff> = HashMap::new();
let mut cadence = ReconcileCadence::new(&agent_mob_mcp_state);
let handles = Box::pin(reconcile_agent_event_streams(
&handle,
&agent_mob_mcp_state,
&mut tracked,
&mut subscribe_failures,
&mut streams,
None,
))
.await;
cadence.rebind(&handles);
loop {
tokio::select! {
Some(forwarded) = streams.next() => {
match forwarded {
ForwardedAgentEvent::Event(event) => {
let (source, source_fence_token, role, envelope, _durable_identity) = *event;
let attributed_event = AttributedEvent {
source,
source_fence_token,
role,
envelope,
};
let _ = mob_events.project_attributed_event(&attributed_event).await;
if event_tx
.send(forwarded_member_event(attributed_event))
.await
.is_err()
{
break;
}
}
ForwardedAgentEvent::Closed(tracked_key) => {
tracked.remove(&tracked_key);
subscribe_failures.remove(&tracked_key);
let handles = Box::pin(reconcile_agent_event_streams(&handle, &agent_mob_mcp_state, &mut tracked, &mut subscribe_failures, &mut streams, None)).await;
cadence.rebind(&handles);
}
}
}
() = cadence.wait(earliest_backoff_attempt(&subscribe_failures)) => {
let handles = Box::pin(reconcile_agent_event_streams(&handle, &agent_mob_mcp_state, &mut tracked, &mut subscribe_failures, &mut streams, None)).await;
cadence.rebind(&handles);
}
}
}
}
async fn run_identity_stream_health_monitor(
handle: MobHandle,
agent_mob_mcp_state: Option<Arc<meerkat_mob_mcp::MobMcpState>>,
identity_runtime: Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
) {
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
let mut tracked = HashSet::new();
let mut subscribe_failures: HashMap<TrackedAgentEventStream, SubscribeBackoff> = HashMap::new();
let mut cadence = ReconcileCadence::new(&agent_mob_mcp_state);
let handles = Box::pin(reconcile_agent_event_streams(
&handle,
&agent_mob_mcp_state,
&mut tracked,
&mut subscribe_failures,
&mut streams,
Some(&identity_runtime),
))
.await;
cadence.rebind(&handles);
loop {
tokio::select! {
Some(forwarded) = streams.next() => {
match forwarded {
ForwardedAgentEvent::Event(event) => {
let (_, _, _, envelope, durable_identity) = *event;
record_identity_turn_completion(
&identity_runtime,
durable_identity.as_deref(),
&envelope,
).await;
}
ForwardedAgentEvent::Closed(tracked_key) => {
tracked.remove(&tracked_key);
subscribe_failures.remove(&tracked_key);
trigger_identity_stream_repair(
handle.mob_id().as_str(),
&tracked_key,
&identity_runtime,
"live agent event stream closed permanently",
).await;
let handles = Box::pin(reconcile_agent_event_streams(
&handle,
&agent_mob_mcp_state,
&mut tracked,
&mut subscribe_failures,
&mut streams,
Some(&identity_runtime),
)).await;
cadence.rebind(&handles);
}
}
}
() = cadence.wait(earliest_backoff_attempt(&subscribe_failures)) => {
let handles = Box::pin(reconcile_agent_event_streams(
&handle,
&agent_mob_mcp_state,
&mut tracked,
&mut subscribe_failures,
&mut streams,
Some(&identity_runtime),
)).await;
cadence.rebind(&handles);
}
}
}
}
async fn reconcile_agent_event_streams(
handle: &MobHandle,
agent_mob_mcp_state: &Option<Arc<meerkat_mob_mcp::MobMcpState>>,
tracked: &mut HashSet<TrackedAgentEventStream>,
subscribe_failures: &mut HashMap<TrackedAgentEventStream, SubscribeBackoff>,
streams: &mut SelectAll<TaggedAgentEventStream>,
identity_runtime: Option<
&Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
>,
) -> Vec<MobHandle> {
let primary_mob_id = handle.mob_id().to_string();
let mut handles = vec![handle.clone()];
if let Some(state) = agent_mob_mcp_state {
handles.extend(
Box::pin(state.mob_handles_snapshot())
.await
.unwrap_or_default()
.into_iter()
.filter_map(|(mob_id, child_handle)| {
if mob_id.as_str() == primary_mob_id {
None
} else {
Some(child_handle)
}
}),
);
}
let mut current: HashSet<TrackedAgentEventStream> = HashSet::new();
for handle in &handles {
let mob_id = handle.mob_id().to_string();
for entry in handle.list_members_including_retiring().await {
let Some((runtime_id, fence_token)) = entry.binding_atoms() else {
continue;
};
let durable_identity = durable_identity_label(&entry.labels);
let identity_fencing_token = current_identity_fencing_token(
&primary_mob_id,
&mob_id,
durable_identity.as_deref(),
identity_runtime,
)
.await;
if identity_runtime.is_some() && identity_fencing_token.is_none() {
continue;
}
current.insert(TrackedAgentEventStream {
mob_id: mob_id.clone(),
durable_identity,
member_identity: entry.agent_identity.clone(),
runtime_id,
identity_fencing_token,
fence_token,
});
}
}
tracked.retain(|tracked_key| current.contains(tracked_key));
subscribe_failures.retain(|key, _| current.contains(key));
for handle in &handles {
let mob_id = handle.mob_id().to_string();
for entry in handle.list_members_including_retiring().await {
let identity = entry.agent_identity.clone();
let Some((runtime_id, fence_token)) = entry.binding_atoms() else {
continue;
};
let durable_identity = durable_identity_label(&entry.labels);
let identity_fencing_token = current_identity_fencing_token(
&primary_mob_id,
&mob_id,
durable_identity.as_deref(),
identity_runtime,
)
.await;
if identity_runtime.is_some() && identity_fencing_token.is_none() {
continue;
}
let tracked_key = TrackedAgentEventStream {
mob_id: mob_id.clone(),
durable_identity,
member_identity: identity.clone(),
runtime_id: runtime_id.clone(),
identity_fencing_token,
fence_token,
};
if tracked.contains(&tracked_key) {
continue;
}
if !forwarder_should_subscribe(entry.status) {
subscribe_failures.remove(&tracked_key);
continue;
}
let now = tokio::time::Instant::now();
if let Some(backoff) = subscribe_failures.get(&tracked_key)
&& now < backoff.next_attempt
{
continue;
}
let role = entry.role.clone();
match subscribe_agent_events_for_console_forwarder(handle, &tracked_key.member_identity)
.await
{
Ok(stream) => {
let close_key = tracked_key.clone();
let durable_identity: Option<Arc<str>> = tracked_key
.durable_identity
.as_deref()
.map(Arc::<str>::from);
subscribe_failures.remove(&tracked_key);
tracked.insert(tracked_key);
let mapped = stream
.map(move |envelope| {
ForwardedAgentEvent::Event(Box::new((
runtime_id.clone(),
fence_token,
role.clone(),
envelope,
durable_identity.clone(),
)))
})
.chain(futures::stream::once(async move {
ForwardedAgentEvent::Closed(close_key)
}))
.boxed();
streams.push(mapped);
}
Err(error) => {
let repair_key = tracked_key.clone();
let backoff =
subscribe_failures
.entry(tracked_key)
.or_insert(SubscribeBackoff {
next_attempt: now,
consecutive_failures: 0,
});
if identity_runtime.is_none() && backoff.consecutive_failures == 0 {
tracing::warn!(
mob_id = %mob_id,
identity = %identity,
error = %error,
"mobkit agent event forwarder: failed to subscribe; will retry with backoff"
);
} else if identity_runtime.is_none() {
tracing::debug!(
mob_id = %mob_id,
identity = %identity,
error = %error,
consecutive_failures = backoff.consecutive_failures,
"mobkit agent event forwarder: subscribe still failing; backing off"
);
}
backoff.next_attempt =
now + subscribe_backoff_delay(backoff.consecutive_failures);
backoff.consecutive_failures = backoff.consecutive_failures.saturating_add(1);
let stream_is_permanently_lost =
backoff.consecutive_failures >= PERMANENT_STREAM_FAILURE_THRESHOLD;
if stream_is_permanently_lost && let Some(identity_runtime) = identity_runtime {
trigger_identity_stream_repair(
&primary_mob_id,
&repair_key,
identity_runtime,
"live agent event stream remained unavailable after bounded retries",
)
.await;
}
}
}
}
}
handles
}
async fn subscribe_agent_events_for_console_forwarder(
handle: &MobHandle,
identity: &AgentIdentity,
) -> Result<EventStream, meerkat_mob::MobError> {
handle.subscribe_agent_events(identity).await
}
async fn run_mob_events_subscription(
handle: MobHandle,
store: MobEventsStore,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
) {
let mob_id = handle.mob_id().as_str().to_string();
let resume_cursor = match persistent_metadata.get_subscription_cursor(&mob_id).await {
Ok(value) => value,
Err(err) => {
tracing::warn!(
mob_id = %mob_id,
error = %err,
"mob_events subscription: failed to read persisted cursor; resuming from latest"
);
None
}
};
let events = handle.events();
let mut subscription = match resume_cursor {
Some(cursor) => match events.subscribe_after(cursor).await {
Ok(sub) => sub,
Err(MobError::StaleEventCursor {
after_cursor,
latest_cursor,
}) => {
tracing::warn!(
mob_id = %mob_id,
after_cursor,
latest_cursor,
"mob_events subscription: persisted cursor is past ledger frontier; resuming at latest"
);
match events.subscribe().await {
Ok(sub) => sub,
Err(err) => {
tracing::warn!(
mob_id = %mob_id,
error = %err,
"mob_events subscription: failed to subscribe at latest after stale-cursor recovery"
);
return;
}
}
}
Err(err) => {
tracing::warn!(
mob_id = %mob_id,
error = %err,
"mob_events subscription: failed to resume from persisted cursor"
);
return;
}
},
None => match events.subscribe().await {
Ok(sub) => sub,
Err(err) => {
tracing::warn!(
mob_id = %mob_id,
error = %err,
"mob_events subscription: initial subscribe failed"
);
return;
}
},
};
while let Some(event) = subscription.event_rx.recv().await {
let envelope = store.project_mob_event(&event).await;
if let Err(err) = persistent_metadata
.set_subscription_cursor(&mob_id, envelope.cursor)
.await
{
tracing::warn!(
mob_id = %mob_id,
cursor = envelope.cursor,
error = %err,
"mob_events subscription: failed to persist cursor; continuing"
);
}
}
}
const ACTOR_LOOP_PROBE_INTERVAL: Duration = Duration::from_mins(1);
const ACTOR_LOOP_PROBE_BUDGET: Duration = Duration::from_secs(30);
fn actor_loop_probe_interval() -> Duration {
parse_probe_secs(
std::env::var("MOBKIT_ACTOR_LOOP_PROBE_INTERVAL_SECS")
.ok()
.as_deref(),
ACTOR_LOOP_PROBE_INTERVAL,
)
}
fn actor_loop_probe_budget() -> Duration {
parse_probe_secs(
std::env::var("MOBKIT_ACTOR_LOOP_PROBE_BUDGET_SECS")
.ok()
.as_deref(),
ACTOR_LOOP_PROBE_BUDGET,
)
}
fn parse_probe_secs(raw: Option<&str>, default: Duration) -> Duration {
raw.and_then(|value| value.trim().parse::<u64>().ok())
.map(|secs| Duration::from_secs(secs.clamp(1, 3600)))
.unwrap_or(default)
}
async fn run_actor_loop_probe<P, F>(
mut probe: P,
error_hook: SharedErrorHook,
interval: Duration,
budget: Duration,
) where
P: FnMut() -> F,
F: Future<Output = Result<MobState, MobError>>,
{
let mut next_stall_id: u64 = 1;
let mut resolved_stalls: u64 = 0;
loop {
tokio::time::sleep(interval).await;
let started = tokio::time::Instant::now();
let round_trip = probe();
tokio::pin!(round_trip);
let result = match tokio::time::timeout(budget, &mut round_trip).await {
Ok(result) => result,
Err(_) => {
let stall_id = next_stall_id;
next_stall_id += 1;
fire_error_hook(
&error_hook,
ErrorEvent::ActorLoopStalled {
probe_waited_secs: started.elapsed().as_secs(),
detail: format!(
"QueryPhase probe round trip unanswered after {}s; the mob actor \
is one serialized command loop, so every member's dispatch is \
queued behind whatever is blocking it; the probe stays parked on \
this round trip and will not stack another. THIS STALL PAGES \
ONCE: no further stall events will be emitted for it, so silence \
is NOT recovery — hold the incident open until an \
actor_loop_recovered arrives with stall_id {}",
budget.as_secs(),
stall_id
),
stall_id: Some(stall_id),
prior_resolved_stalls: Some(resolved_stalls),
},
);
let result = round_trip.await;
resolved_stalls += 1;
fire_error_hook(
&error_hook,
ErrorEvent::ActorLoopRecovered {
stall_id,
stalled_for_secs: started.elapsed().as_secs(),
},
);
result
}
};
if matches!(
result,
Ok(MobState::Stopped | MobState::Completed | MobState::Destroyed)
) {
break;
}
}
}
fn forwarded_member_event(attributed: AttributedEvent) -> ForwardedMemberEvent {
ForwardedMemberEvent {
alert: compaction_rejection_alert(&attributed),
envelope: attributed_event_to_unified(attributed),
}
}
fn compaction_rejection_alert(attributed: &AttributedEvent) -> Option<ErrorEvent> {
let AgentEvent::CompactionFailed { reason } = &attributed.envelope.payload else {
return None;
};
let session_id = attributed
.envelope
.source
.session_id()
.map(ToString::to_string)
.unwrap_or_default();
let (preserved_history, attempted_entries) = match reason {
meerkat_core::event::CompactionFailureReason::ProjectionHandoffRefused {
preserved_history,
attempted_entries,
..
} => (Some((*preserved_history).into()), Some(*attempted_entries)),
_ => (None, None),
};
Some(ErrorEvent::CompactionPersistenceRejected {
identity: crate::member_comms_id::runtime_event_alias(&attributed.source),
session_id,
error: reason.to_string(),
preserved_history,
attempted_entries,
})
}
fn attributed_event_to_unified(attributed: AttributedEvent) -> EventEnvelope<UnifiedEvent> {
EventEnvelope {
event_id: format!("evt-agent-{}", attributed.envelope.event_id),
source: "agent".to_string(),
timestamp_ms: attributed.envelope.timestamp_ms,
event: UnifiedEvent::Agent {
agent_id: crate::member_comms_id::runtime_event_alias(&attributed.source),
event_type: agent_event_type(&attributed.envelope.payload).to_string(),
payload: Some(crate::mob_handle_runtime::console_agent_event_payload(
&attributed.envelope.payload,
)),
},
}
}
struct ConsoleMemoryEventSink {
store: ConsoleEventStore,
handle: tokio::runtime::Handle,
}
impl crate::memory::events::MemoryEventSink for ConsoleMemoryEventSink {
fn emit(&self, event: crate::memory::events::MemoryTimelineEvent) {
let store = self.store.clone();
let identity = event
.identity()
.map(str::to_string)
.unwrap_or_else(|| crate::console_contracts::SYSTEM_EVENT_IDENTITY.to_string());
let event_type = event.event_type().to_string();
let data = event.data();
self.handle.spawn(async move {
store.append(identity, None, event_type, data).await;
});
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::panic)]
mod tests {
use std::sync::atomic::Ordering;
use super::*;
use meerkat_mob::ids::Generation;
fn attributed_text_delta(member_id: &str, generation: u64) -> AttributedEvent {
AttributedEvent {
source: AgentRuntimeId::new(
AgentIdentity::from(member_id),
Generation::new(generation),
),
source_fence_token: FenceToken::new(1),
role: ProfileName::from("worker"),
envelope: meerkat_core::event::EventEnvelope {
event_id: Default::default(),
source: meerkat_core::event::EventSourceIdentity::runtime("test"),
seq: 0,
mob_id: None,
timestamp_ms: 1,
payload: AgentEvent::TextDelta {
delta: "hello".to_string(),
},
},
}
}
#[test]
fn identity_stream_tracking_uses_trusted_durable_identity_label() {
let labels =
BTreeMap::from([("agent_identity".to_string(), "review:singleton".to_string())]);
assert_eq!(
durable_identity_label(&labels).as_deref(),
Some("review:singleton")
);
assert_eq!(
durable_identity_label(&BTreeMap::new()),
None,
"ordinary mobs must not be guessed into identity authority"
);
}
#[test]
fn attributed_event_ingest_decodes_encoded_roster_member_ids() {
let encoded = crate::member_comms_id::mob_member_id_str("rt:review:singleton:0");
assert!(encoded.starts_with("mk--"), "precondition: alias encodes");
let unified = attributed_event_to_unified(attributed_text_delta(&encoded, 1));
let UnifiedEvent::Agent { agent_id, .. } = unified.event else {
panic!("expected agent event");
};
assert_eq!(agent_id, "rt:review:singleton:0:1");
}
#[test]
fn attributed_event_ingest_passes_plain_member_ids_through() {
let unified = attributed_event_to_unified(attributed_text_delta("worker-one", 0));
let UnifiedEvent::Agent { agent_id, .. } = unified.event else {
panic!("expected agent event");
};
assert_eq!(agent_id, "worker-one:0");
}
#[tokio::test(start_paused = true)]
async fn reconcile_cadence_safety_deadline_survives_recreated_waits() {
let mut cadence = ReconcileCadence::new(&None);
let mut fired = false;
for _ in 0..7 {
tokio::select! {
() = cadence.wait(None) => {
fired = true;
break;
}
() = tokio::time::sleep(Duration::from_secs(5)) => {}
}
}
assert!(
fired,
"safety reconcile starved: recreated wait futures reset the deadline"
);
}
#[tokio::test(start_paused = true)]
async fn reconcile_cadence_rebind_rearms_safety_deadline() {
let mut cadence = ReconcileCadence::new(&None);
tokio::time::sleep(Duration::from_secs(20)).await;
cadence.rebind(&[]);
tokio::select! {
() = cadence.wait(None) => {
panic!("deadline fired 30s after construction despite rebind re-arm")
}
() = tokio::time::sleep(Duration::from_secs(29)) => {}
}
tokio::select! {
() = cadence.wait(None) => {}
() = tokio::time::sleep(Duration::from_secs(2)) => {
panic!("re-armed deadline did not fire 30s after rebind")
}
}
}
#[test]
fn forwarder_only_subscribes_active_members() {
assert!(forwarder_should_subscribe(MobMemberStatus::Active));
assert!(!forwarder_should_subscribe(MobMemberStatus::Retiring));
assert!(!forwarder_should_subscribe(MobMemberStatus::Broken));
assert!(!forwarder_should_subscribe(MobMemberStatus::Completed));
assert!(!forwarder_should_subscribe(MobMemberStatus::Unknown));
}
#[test]
fn subscribe_backoff_grows_and_caps() {
const { assert!(PERMANENT_STREAM_FAILURE_THRESHOLD > 1) };
assert_eq!(subscribe_backoff_delay(0), SUBSCRIBE_BACKOFF_BASE);
assert_eq!(subscribe_backoff_delay(1), SUBSCRIBE_BACKOFF_BASE * 2);
assert_eq!(subscribe_backoff_delay(3), SUBSCRIBE_BACKOFF_BASE * 8);
assert_eq!(subscribe_backoff_delay(7), SUBSCRIBE_BACKOFF_MAX);
assert_eq!(subscribe_backoff_delay(50), SUBSCRIBE_BACKOFF_MAX);
assert!(subscribe_backoff_delay(2) > subscribe_backoff_delay(1));
}
async fn bootstrap_minimal_runtime(temp_dir: &tempfile::TempDir) -> UnifiedRuntime {
let session_path = temp_dir.path().join("sessions");
std::fs::create_dir_all(&session_path).expect("session path");
let factory = meerkat::AgentFactory::new(&session_path);
let session_service: Arc<dyn meerkat_mob::MobSessionService> = Arc::new(
meerkat::build_ephemeral_service(factory, meerkat::Config::default(), 16),
);
let definition = meerkat_mob::MobDefinition::from_toml(
r#"
[mob]
id = "compaction-alert-mob"
[profiles.worker]
model = "gpt-5.5"
"#,
)
.expect("parse mob definition");
let mob_spec = MobBootstrapSpec::new(
definition,
meerkat_mob::MobStorage::in_memory(),
session_service,
)
.with_options(crate::mob_handle_runtime::MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: Some(Arc::new(meerkat_client::TestClient::for_provider(
meerkat_core::Provider::OpenAI,
))),
});
let module_config = MobKitConfig {
modules: vec![],
discovery: crate::types::DiscoverySpec {
namespace: "compaction-alert".to_string(),
modules: vec![],
},
pre_spawn: vec![],
};
UnifiedRuntime::bootstrap(mob_spec, module_config, Duration::from_secs(2))
.await
.expect("bootstrap unified runtime")
}
async fn compaction_alerts_through_drain(
session_id: &meerkat_core::types::SessionId,
reasons: Vec<meerkat_core::event::CompactionFailureReason>,
) -> Vec<ErrorEvent> {
let temp_dir = tempfile::tempdir().expect("temp dir");
let mut runtime = bootstrap_minimal_runtime(&temp_dir).await;
let captured: Arc<tokio::sync::Mutex<Vec<ErrorEvent>>> =
Arc::new(tokio::sync::Mutex::new(Vec::new()));
let hook_captured = captured.clone();
let hook: ErrorHook = Arc::new(move |event| {
let hook_captured = hook_captured.clone();
Box::pin(async move {
hook_captured.lock().await.push(event);
})
});
runtime.set_error_hook(hook);
let event_tx = runtime.install_test_event_ingress().await;
let expected = reasons.len();
for (seq, reason) in reasons.into_iter().enumerate() {
let attributed = AttributedEvent {
source: AgentRuntimeId::new(
AgentIdentity::from("compaction-worker"),
Generation::new(0),
),
source_fence_token: FenceToken::new(1),
role: ProfileName::from("worker"),
envelope: meerkat_core::event::EventEnvelope {
event_id: Default::default(),
source: meerkat_core::event::EventSourceIdentity::session(session_id.clone()),
seq: seq as u64,
mob_id: None,
timestamp_ms: 9,
payload: AgentEvent::CompactionFailed { reason },
},
};
event_tx
.send(forwarded_member_event(attributed))
.await
.expect("send forwarded event");
}
runtime
.drain_mob_agent_events()
.await
.expect("drain member events");
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
loop {
if captured.lock().await.len() >= expected {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"error hook did not receive all {expected} compaction persistence rejections"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
let alerts = captured.lock().await.clone();
runtime.shutdown().await;
alerts
}
#[tokio::test]
async fn compaction_failures_page_error_hook_with_typed_fit_through_drain() {
let session_id = meerkat_core::types::SessionId::new();
let alerts = compaction_alerts_through_drain(
&session_id,
vec![
meerkat_core::event::CompactionFailureReason::TranscriptRewriteFailed {
message: "runtime epoch mismatch".to_string(),
},
meerkat_core::event::CompactionFailureReason::ProjectionHandoffRefused {
refusal: meerkat_core::memory::CompactionHandoffRefusal::RuntimeEpochRotated,
preserved_history:
meerkat_core::event::CompactionPreservedHistoryFit::StillFits,
attempted_entries: 12,
message: "runtime epoch rotated under the coordinator".to_string(),
},
],
)
.await;
let mut untyped = None;
let mut typed = None;
for alert in &alerts {
match alert {
ErrorEvent::CompactionPersistenceRejected {
identity,
session_id: rejected_session,
error,
preserved_history,
attempted_entries,
} => {
assert_eq!(identity, "compaction-worker:0");
assert_eq!(rejected_session, &session_id.to_string());
if preserved_history.is_some() {
typed = Some((error.clone(), *preserved_history, *attempted_entries));
} else {
untyped = Some((error.clone(), *attempted_entries));
}
}
other => panic!("expected CompactionPersistenceRejected, got {other:?}"),
}
}
let (untyped_error, untyped_entries) =
untyped.unwrap_or_else(|| panic!("the non-handoff failure must page too: {alerts:?}"));
assert!(
untyped_error.contains("runtime epoch mismatch"),
"error must carry the rejection detail: {untyped_error}"
);
assert_eq!(
untyped_entries, None,
"a non-handoff compaction failure carries no fit verdict"
);
let (typed_error, typed_fit, typed_entries) = typed.unwrap_or_else(|| {
panic!("the projection-handoff refusal must page with its fit: {alerts:?}")
});
assert_eq!(
typed_fit,
Some(CompactionPreservedHistoryFit::StillFits),
"the wedged/progressing discriminator must cross typed"
);
assert_eq!(typed_entries, Some(12));
assert!(
typed_error.contains("runtime epoch rotated under the coordinator"),
"error must keep the human rendering: {typed_error}"
);
}
#[derive(Clone, Default)]
struct CaptureWriter(Arc<std::sync::Mutex<Vec<u8>>>);
impl CaptureWriter {
fn contents(&self) -> String {
String::from_utf8_lossy(
&self
.0
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
)
.into_owned()
}
}
impl std::io::Write for CaptureWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CaptureWriter {
type Writer = CaptureWriter;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
fn capture_tracing<T>(body: impl FnOnce() -> T) -> (T, String) {
let writer = CaptureWriter::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(writer.clone())
.with_max_level(tracing::Level::INFO)
.with_ansi(false)
.finish();
let value = tracing::subscriber::with_default(subscriber, body);
(value, writer.contents())
}
fn sample_alert() -> ErrorEvent {
ErrorEvent::CompactionPersistenceRejected {
identity: "compaction-worker:0".to_string(),
session_id: "sess-1".to_string(),
error: "runtime refused the durable compaction projection handoff".to_string(),
preserved_history: Some(CompactionPreservedHistoryFit::OverWindow),
attempted_entries: Some(12),
}
}
#[test]
fn error_event_without_hook_still_reaches_the_log() {
let slot: SharedErrorHook = Arc::new(std::sync::RwLock::new(None));
let ((), logged) = capture_tracing(|| fire_error_hook(&slot, sample_alert()));
assert!(
logged.contains("CompactionPersistenceRejected"),
"the typed variant must be named in the record: {logged}"
);
assert!(
logged.contains("OverWindow") && logged.contains("compaction-worker:0"),
"typed fields must ride the record, not just the Display string: {logged}"
);
assert!(
logged.contains("hook_registered=false"),
"the record must say nobody was listening: {logged}"
);
}
#[test]
fn error_event_with_hook_is_logged_as_delivered() {
let hook: ErrorHook = Arc::new(move |_event| Box::pin(async move {}));
let slot: SharedErrorHook = Arc::new(std::sync::RwLock::new(Some(hook)));
let runtime = tokio::runtime::Runtime::new().expect("tokio runtime");
let _guard = runtime.enter();
let ((), logged) = capture_tracing(|| fire_error_hook(&slot, sample_alert()));
assert!(
logged.contains("hook_registered=true"),
"a wired host still gets the record, marked as delivered: {logged}"
);
}
#[tokio::test]
async fn registered_error_hook_fires_exactly_once() {
let captured: Arc<tokio::sync::Mutex<Vec<ErrorEvent>>> =
Arc::new(tokio::sync::Mutex::new(Vec::new()));
let hook_captured = captured.clone();
let hook: ErrorHook = Arc::new(move |event| {
let hook_captured = hook_captured.clone();
Box::pin(async move {
hook_captured.lock().await.push(event);
})
});
let slot: SharedErrorHook = Arc::new(std::sync::RwLock::new(Some(hook)));
fire_error_hook(&slot, sample_alert());
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
loop {
if !captured.lock().await.is_empty() {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"registered hook never received the event"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
captured.lock().await.len(),
1,
"the hook must fire exactly once per event"
);
}
#[test]
fn error_hook_absent_notice_points_at_the_registration_call() {
let ((), logged) = capture_tracing(emit_error_hook_absent_notice);
assert!(
logged.contains("no error hook is registered"),
"the notice must state the condition plainly: {logged}"
);
assert!(
logged.contains("on_error"),
"the notice must point at the registration call: {logged}"
);
}
#[test]
fn compaction_rejection_wire_shape_stays_additive() {
let legacy_json = serde_json::json!({
"category": "compaction_persistence_rejected",
"identity": "compaction-worker:0",
"session_id": "sess-1",
"error": "compaction curator failed: no summary",
});
let alert: ErrorEvent =
serde_json::from_value(legacy_json.clone()).expect("pre-field payload must load");
assert_eq!(
alert,
ErrorEvent::CompactionPersistenceRejected {
identity: "compaction-worker:0".to_string(),
session_id: "sess-1".to_string(),
error: "compaction curator failed: no summary".to_string(),
preserved_history: None,
attempted_entries: None,
}
);
assert_eq!(
serde_json::to_value(&alert).expect("serialize"),
legacy_json,
"an alert with no fit verdict must not emit the new keys"
);
let typed = ErrorEvent::CompactionPersistenceRejected {
identity: "compaction-worker:0".to_string(),
session_id: "sess-1".to_string(),
error: "runtime refused the durable compaction projection handoff".to_string(),
preserved_history: Some(CompactionPreservedHistoryFit::OverWindow),
attempted_entries: Some(12),
};
let encoded = serde_json::to_value(&typed).expect("serialize typed");
assert_eq!(
encoded["preserved_history"],
serde_json::json!("over_window")
);
assert_eq!(encoded["attempted_entries"], serde_json::json!(12));
assert_eq!(
serde_json::from_value::<ErrorEvent>(encoded).expect("round trip"),
typed
);
}
fn capturing_error_hook_slot() -> (SharedErrorHook, Arc<tokio::sync::Mutex<Vec<ErrorEvent>>>) {
let captured: Arc<tokio::sync::Mutex<Vec<ErrorEvent>>> =
Arc::new(tokio::sync::Mutex::new(Vec::new()));
let hook_captured = captured.clone();
let hook: ErrorHook = Arc::new(move |event| {
let hook_captured = hook_captured.clone();
Box::pin(async move {
hook_captured.lock().await.push(event);
})
});
let slot: SharedErrorHook = Arc::new(std::sync::RwLock::new(Some(hook)));
(slot, captured)
}
#[tokio::test(start_paused = true)]
async fn stalled_actor_loop_pages_error_hook() {
let (slot, captured) = capturing_error_hook_slot();
let probe_task = tokio::spawn(run_actor_loop_probe(
std::future::pending::<Result<MobState, MobError>>,
slot,
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_secs(95)).await;
tokio::task::yield_now().await;
let events = captured.lock().await;
assert_eq!(events.len(), 1, "exactly one stall page: {events:?}");
match &events[0] {
ErrorEvent::ActorLoopStalled {
probe_waited_secs,
detail,
stall_id,
prior_resolved_stalls,
} => {
assert_eq!(*probe_waited_secs, 30);
assert!(
detail.contains("QueryPhase"),
"detail must name the probe round trip: {detail}"
);
assert!(
detail.contains("PAGES ONCE") && detail.contains("silence is NOT recovery"),
"detail must say the absence of further pages is not recovery: {detail}"
);
assert!(
detail.contains("actor_loop_recovered") && detail.contains("stall_id 1"),
"detail must name the resolution to wait for, by id: {detail}"
);
assert_eq!(
*stall_id,
Some(1),
"the stall must carry the id its resolution will echo"
);
assert_eq!(*prior_resolved_stalls, Some(0));
}
other => panic!("expected ActorLoopStalled, got {other:?}"),
}
drop(events);
probe_task.abort();
}
#[tokio::test(start_paused = true)]
async fn wedged_actor_loop_pages_once_and_never_resolves() {
let (slot, captured) = capturing_error_hook_slot();
let probe_task = tokio::spawn(run_actor_loop_probe(
std::future::pending::<Result<MobState, MobError>>,
slot,
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_mins(10)).await;
tokio::task::yield_now().await;
let events = captured.lock().await;
assert_eq!(
events.len(),
1,
"a wedged loop pages once and cannot page again: {events:?}"
);
assert!(
!events
.iter()
.any(|event| matches!(event, ErrorEvent::ActorLoopRecovered { .. })),
"a wedged loop must never emit a resolution: {events:?}"
);
drop(events);
probe_task.abort();
}
#[tokio::test(start_paused = true)]
async fn resolved_stall_emits_recovery_correlated_by_stall_id() {
let (slot, captured) = capturing_error_hook_slot();
let probe_task = tokio::spawn(run_actor_loop_probe(
|| async {
tokio::time::sleep(Duration::from_secs(50)).await;
Ok(MobState::Running)
},
slot,
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_mins(2)).await;
tokio::task::yield_now().await;
let events = captured.lock().await;
let stall_id = events
.iter()
.find_map(|event| match event {
ErrorEvent::ActorLoopStalled { stall_id, .. } => Some(*stall_id),
_ => None,
})
.unwrap_or_else(|| panic!("expected a stall page: {events:?}"));
let (recovered_id, stalled_for_secs) = events
.iter()
.find_map(|event| match event {
ErrorEvent::ActorLoopRecovered {
stall_id,
stalled_for_secs,
} => Some((*stall_id, *stalled_for_secs)),
_ => None,
})
.unwrap_or_else(|| panic!("expected a resolution: {events:?}"));
assert_eq!(
Some(recovered_id),
stall_id,
"the resolution must name the stall it closes, or a receiver \
cannot close the incident it opened: {events:?}"
);
assert!(
stalled_for_secs >= 50,
"the resolution must carry how long the loop was stalled, got \
{stalled_for_secs}s"
);
drop(events);
probe_task.abort();
}
#[test]
fn resolved_stall_is_not_logged_as_an_error() {
let recovered = ErrorEvent::ActorLoopRecovered {
stall_id: 7,
stalled_for_secs: 42,
};
let ((), logged) = capture_tracing(|| log_error_event(&recovered, true));
assert!(
logged.contains("INFO"),
"a resolution must log at INFO: {logged}"
);
assert!(
!logged.contains("ERROR"),
"a resolution must not log at ERROR: {logged}"
);
assert!(
logged.contains("actor_loop_recovered") && logged.contains("42"),
"the record must still carry the resolution facts: {logged}"
);
}
#[test]
fn actor_loop_stall_wire_shape_stays_additive() {
let legacy_json = serde_json::json!({
"category": "actor_loop_stalled",
"probe_waited_secs": 30,
"detail": "QueryPhase probe round trip unanswered",
});
let alert: ErrorEvent =
serde_json::from_value(legacy_json.clone()).expect("pre-field payload must load");
assert_eq!(
alert,
ErrorEvent::ActorLoopStalled {
probe_waited_secs: 30,
detail: "QueryPhase probe round trip unanswered".to_string(),
stall_id: None,
prior_resolved_stalls: None,
}
);
assert_eq!(
serde_json::to_value(&alert).expect("serialize"),
legacy_json,
"a stall with no correlation must not emit the new keys"
);
let correlated = ErrorEvent::ActorLoopStalled {
probe_waited_secs: 30,
detail: "stalled".to_string(),
stall_id: Some(3),
prior_resolved_stalls: Some(2),
};
let encoded = serde_json::to_value(&correlated).expect("serialize correlated");
assert_eq!(encoded["stall_id"], serde_json::json!(3));
assert_eq!(
serde_json::from_value::<ErrorEvent>(encoded).expect("round trip"),
correlated
);
}
#[tokio::test(start_paused = true)]
async fn healthy_actor_loop_never_pages() {
let (slot, captured) = capturing_error_hook_slot();
let probes = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let counted = probes.clone();
let probe_task = tokio::spawn(run_actor_loop_probe(
move || {
counted.fetch_add(1, Ordering::SeqCst);
std::future::ready(Ok(MobState::Running))
},
slot,
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_secs(60 * 5 + 5)).await;
tokio::task::yield_now().await;
assert!(
probes.load(Ordering::SeqCst) >= 4,
"probe must keep its cadence on a healthy loop: {}",
probes.load(Ordering::SeqCst)
);
assert!(
captured.lock().await.is_empty(),
"a healthy loop must not page"
);
probe_task.abort();
}
#[tokio::test(start_paused = true)]
async fn stalled_probe_never_stacks_and_resumes_after_recovery() {
let (slot, captured) = capturing_error_hook_slot();
let probes = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let released = Arc::new(std::sync::atomic::AtomicBool::new(false));
let gate = Arc::new(tokio::sync::Notify::new());
let counted = probes.clone();
let released_probe = released.clone();
let gate_probe = gate.clone();
let probe_task = tokio::spawn(run_actor_loop_probe(
move || {
counted.fetch_add(1, Ordering::SeqCst);
let released = released_probe.clone();
let gate = gate_probe.clone();
async move {
if !released.load(Ordering::SeqCst) {
gate.notified().await;
}
Ok(MobState::Running)
}
},
slot,
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_mins(10)).await;
tokio::task::yield_now().await;
assert_eq!(
probes.load(Ordering::SeqCst),
1,
"a stalled probe must not stack another round trip"
);
assert_eq!(
captured.lock().await.len(),
1,
"a persisting stall pages exactly once"
);
released.store(true, Ordering::SeqCst);
gate.notify_one();
tokio::time::sleep(Duration::from_secs(61)).await;
tokio::task::yield_now().await;
assert!(
probes.load(Ordering::SeqCst) >= 2,
"probe must resume after the stalled round trip drains: {}",
probes.load(Ordering::SeqCst)
);
let events = captured.lock().await;
assert_eq!(
events.len(),
2,
"the drained stall must produce its resolution: {events:?}"
);
let stall_id = match &events[0] {
ErrorEvent::ActorLoopStalled { stall_id, .. } => *stall_id,
other => panic!("expected the stall first, got {other:?}"),
};
match &events[1] {
ErrorEvent::ActorLoopRecovered {
stall_id: closed, ..
} => {
assert_eq!(
Some(*closed),
stall_id,
"the resolution must close the stall that opened: {events:?}"
);
}
other => panic!("expected the resolution second, got {other:?}"),
}
drop(events);
probe_task.abort();
}
#[test]
fn probe_env_knob_parses_and_clamps() {
assert_eq!(
parse_probe_secs(None, ACTOR_LOOP_PROBE_INTERVAL),
Duration::from_mins(1)
);
assert_eq!(
parse_probe_secs(Some("120"), ACTOR_LOOP_PROBE_INTERVAL),
Duration::from_mins(2)
);
assert_eq!(
parse_probe_secs(Some(" 15 "), ACTOR_LOOP_PROBE_BUDGET),
Duration::from_secs(15)
);
assert_eq!(
parse_probe_secs(Some("0"), ACTOR_LOOP_PROBE_BUDGET),
Duration::from_secs(1)
);
assert_eq!(
parse_probe_secs(Some("999999"), ACTOR_LOOP_PROBE_BUDGET),
Duration::from_hours(1)
);
assert_eq!(
parse_probe_secs(Some("junk"), ACTOR_LOOP_PROBE_BUDGET),
ACTOR_LOOP_PROBE_BUDGET
);
}
}