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 serde_json::json;
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 live_compose;
pub mod mob_events;
pub mod mob_ops;
pub(crate) mod mob_stop;
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 mob_stop::{MOB_STOP_FLOW_SETTLE_BUDGET, MobStopFlowRunsUnsettled};
pub use types::{
CompactionPreservedHistoryFit, ErrorEvent, IdentityAuthorityReleaseOutcome, MobStopOutcome,
MobTerminalShutdownOutcome, 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
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum UnifiedRuntimeBootstrapPhase {
RegisterPersistedOwners,
PrewarmPersistedAuthority,
Activate,
RestoreRoster,
FailureCleanup,
}
pub type UnifiedRuntimeBootstrapPhaseObserver =
Arc<dyn Fn(UnifiedRuntimeBootstrapPhase) + Send + Sync>;
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>>,
bootstrap_phase_observer: std::sync::OnceLock<UnifiedRuntimeBootstrapPhaseObserver>,
pending_mob_activation:
tokio::sync::Mutex<Option<crate::mob_handle_runtime::PendingMobActivation>>,
event_log: Option<event_log::EventLogHandle>,
console_log_store: Arc<dyn ConsoleLogStore>,
console_projection: std::sync::OnceLock<crate::console_aggregator::MobKitConsoleAggregator>,
live_plan: Option<live_compose::LivePlan>,
live_composition: tokio::sync::OnceCell<live_compose::LiveComposition>,
console_events: ConsoleEventStore,
mob_events: MobEventsStore,
mob_events_subscriber_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
actor_loop_probe_task: tokio::sync::Mutex<Option<JoinHandle<()>>>,
actor_loop_health: Arc<crate::actor_loop_health::ActorLoopHealth>,
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_wake: tokio::sync::watch::Sender<u64>,
identity_completion_ledger: Arc<HealthCompletionLedger>,
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>>,
retired_supervisor_cleanups:
tokio::sync::Mutex<tokio::task::JoinSet<types::RetiredSupervisorKind>>,
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>,
}
#[cfg(test)]
pub(crate) struct TestConsoleEventIngress(Sender<ForwardedMemberEvent>);
#[cfg(test)]
impl TestConsoleEventIngress {
pub(crate) async fn send(
&self,
envelope: EventEnvelope<UnifiedEvent>,
) -> Result<(), &'static str> {
self.0
.send(ForwardedMemberEvent {
envelope,
alert: None,
})
.await
.map_err(|_| "test console ingress is closed")
}
}
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 identity_lease_wake = tokio::sync::watch::channel(0_u64).0;
let identity_completion_ledger = Arc::new(HealthCompletionLedger::default());
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),
Arc::clone(&identity_completion_ledger),
identity_lease_wake.subscribe(),
));
let mob_events_task = Self::spawn_mob_events_subscriber(
mob_runtime.handle(),
mob_events_store.clone(),
persistent_metadata.clone(),
mob_runtime.implicit_delegate_retirement_overrides(),
);
let error_hook: SharedErrorHook = Arc::new(std::sync::RwLock::new(None));
let actor_loop_health = crate::actor_loop_health::ActorLoopHealth::shared();
let actor_loop_probe_task = Self::spawn_actor_loop_probe(
mob_runtime.handle(),
Arc::clone(&error_hook),
Arc::clone(&actor_loop_health),
);
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),
bootstrap_phase_observer: std::sync::OnceLock::new(),
pending_mob_activation: tokio::sync::Mutex::new(None),
event_log: None,
console_log_store: Arc::new(InMemoryConsoleLogStore::new()),
console_projection: std::sync::OnceLock::new(),
live_plan: None,
live_composition: tokio::sync::OnceCell::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),
actor_loop_health,
implicit_delegate_retirement_task: tokio::sync::Mutex::new(None),
implicit_delegate_identity_runtime: identity_runtime_authority,
identity_lease_wake,
identity_completion_ledger,
identity_lease_renewal_task: tokio::sync::Mutex::new(None),
identity_continuity_repair_task: tokio::sync::Mutex::new(None),
retired_supervisor_cleanups: tokio::sync::Mutex::new(tokio::task::JoinSet::new()),
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>,
idle_retire_overrides: Option<
crate::mob_handle_runtime::ImplicitDelegateRetirementOverrides,
>,
) -> Option<JoinHandle<()>> {
let runtime_handle = tokio::runtime::Handle::try_current().ok()?;
Some(runtime_handle.spawn(run_mob_events_subscription(
handle,
store,
persistent_metadata,
idle_retire_overrides,
)))
}
fn spawn_actor_loop_probe(
handle: MobHandle,
error_hook: SharedErrorHook,
health: Arc<crate::actor_loop_health::ActorLoopHealth>,
) -> 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,
health,
actor_loop_probe_interval(),
actor_loop_probe_budget(),
)))
}
pub fn actor_loop_health(&self) -> &Arc<crate::actor_loop_health::ActorLoopHealth> {
&self.actor_loop_health
}
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, pending_mob_activation) = MobRuntime::prepare(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;
let staged = pending_mob_activation.is_some();
*runtime.pending_mob_activation.lock().await = pending_mob_activation;
if staged {
tracing::info!(
"identity-first resume: mob left Stopped pending continuity registration; \
bootstrap edge reconciliation deferred until activation"
);
} else {
runtime.reconcile_bootstrap_edges_if_configured().await;
}
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 activate_without_identity_context(&mut self) -> Result<(), MobRuntimeError> {
if self.identity_first_context.is_some() {
return Err(MobRuntimeError::InvalidConfig(
"cannot activate without identity authority after an identity context is installed"
.to_string(),
));
}
let pending = self.pending_mob_activation.lock().await.take();
if let Some(pending) = pending {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::Activate);
if let Err(error) = pending.activate().await {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::FailureCleanup);
self.shutdown().await;
return Err(error);
}
self.reconcile_bootstrap_edges_if_configured().await;
}
Ok(())
}
pub fn set_bootstrap_phase_observer(
&self,
observer: UnifiedRuntimeBootstrapPhaseObserver,
) -> bool {
self.bootstrap_phase_observer.set(observer).is_ok()
}
fn report_bootstrap_phase(&self, phase: UnifiedRuntimeBootstrapPhase) {
if let Some(observer) = self.bootstrap_phase_observer.get() {
observer(phase);
}
}
async fn reconcile_bootstrap_edges_if_configured(&self) {
if self.edge_discovery.is_some()
|| self.topology_controller.revision().await > 0
|| self.topology_controller.has_pending().await
{
let report = self.reconcile_edges().await;
*self.bootstrap_edges_report.write().await = Some(report);
}
}
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));
let pending = self.pending_mob_activation.lock().await.take();
let reconcile_after_materialization =
pending.is_some() && !context.bootstrap_mode().is_lazy();
if let Some(pending) = pending {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::RegisterPersistedOwners);
let registered_sessions = match context
.runtime
.register_persisted_continuity_owners(roster)
.await
{
Ok(registered) => {
tracing::info!(
registered = registered.len(),
"registered persisted continuity owners before mob activation"
);
registered
}
Err(error) => {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::FailureCleanup);
self.shutdown().await;
return Err(error);
}
};
context
.runtime
.prepublish_customizer_tools(roster, context.customizer.as_deref())
.await;
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::PrewarmPersistedAuthority);
let converged = self
.mob_runtime
.prewarm_persisted_runtime_authority(®istered_sessions)
.await;
tracing::info!(
converged,
requested = registered_sessions.len(),
"converged persisted runtime authority before mob activation"
);
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::Activate);
if let Err(error) = pending.activate().await {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::FailureCleanup);
self.shutdown().await;
return Err(crate::identity_first::IdentityRuntimeError::Internal(
format!(
"activating the prepared mob after registering continuity owners: {error}"
),
));
}
if !reconcile_after_materialization {
self.reconcile_bootstrap_edges_if_configured().await;
}
}
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::RestoreRoster);
match context.bootstrap_roster(roster).await {
Ok(result) => {
if reconcile_after_materialization {
self.reconcile_bootstrap_edges_if_configured().await;
}
self.start_identity_first_supervisors();
Ok(result)
}
Err(error) => {
self.report_bootstrap_phase(UnifiedRuntimeBootstrapPhase::FailureCleanup);
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);
context
.runtime
.install_actor_loop_health(Arc::clone(&self.actor_loop_health));
self.mob_runtime
.install_identity_runtime_authority(Arc::clone(&context.runtime));
if let Some(overrides) = self.mob_runtime.implicit_delegate_retirement_overrides() {
context.runtime.install_session_rotation_observer(Arc::new(
crate::mob_handle_runtime::IdleRetireRotationObserver::new(
overrides,
self.mob_runtime.handle().mob_id().to_string(),
),
));
}
*self
.implicit_delegate_identity_runtime
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(Arc::clone(&context.runtime));
if let Some(session_service) = self.mob_runtime.session_service().cloned() {
context
.runtime
.install_pending_completion_drain(Arc::new(ReplayedCompletionDrain {
session_service,
ledger: Arc::clone(&self.identity_completion_ledger),
}));
}
context
.runtime
.install_lease_observer(self.identity_lease_wake.clone());
self.identity_first_context = Some(context);
if let Some(projection) = self.console_projection.get() {
projection.update_identity_authority("default", self.identity_runtime().cloned());
}
}
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(),
)
})?;
tracing::info!(
cause = "flow_target_provisioning",
"identity-first flow barrier is materializing every registered \
identity before the flow run id is minted; cost is O(fleet), \
not O(flow targets)"
);
let started = std::time::Instant::now();
let outcome = runtime
.materialize_all_required_tracked()
.await
.map(|records| records.len())
.map_err(|error| {
MobError::Internal(format!(
"identity-first flow materialization failed: {error}"
))
});
match &outcome {
Ok(materialized) => tracing::info!(
materialized = *materialized,
elapsed_ms = started.elapsed().as_millis() as u64,
cause = "flow_target_provisioning",
"identity-first flow barrier complete"
),
Err(error) => tracing::warn!(
elapsed_ms = started.elapsed().as_millis() as u64,
cause = "flow_target_provisioning",
%error,
"identity-first flow barrier failed"
),
}
outcome.map(|_| ())
})
}));
}
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();
self.retain_retired_supervisor_cleanup(
types::RetiredSupervisorKind::LeaseRenewal,
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();
self.retain_retired_supervisor_cleanup(
types::RetiredSupervisorKind::ContinuityRepair,
previous.cancel_and_join(),
);
}
}
fn retain_retired_supervisor_cleanup(
&mut self,
kind: types::RetiredSupervisorKind,
cleanup: impl std::future::Future<Output = ()> + Send + 'static,
) {
let retired = self.retired_supervisor_cleanups.get_mut();
while retired.try_join_next().is_some() {}
retired.spawn(async move {
cleanup.await;
kind
});
}
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_realm_migration(
&self,
) -> Option<std::sync::Arc<crate::workgraph_realm::WorkGraphRealmMigrationReport>> {
self.mob_runtime.workgraph_realm_migration()
}
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;
if let Some(projection) = self.console_projection.take() {
projection.unregister_runtime("default");
}
}
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 reserve_caller_identity_interaction(
&self,
identity: &str,
runtime_member_id: Option<&str>,
interaction_id: &str,
origin: &str,
content: serde_json::Value,
) -> Result<(), console_events::InteractionIdInFlight> {
self.console_events
.reserve_caller_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>>>,
>,
completion_ledger: Arc<HealthCompletionLedger>,
identity_lease_changes: tokio::sync::watch::Receiver<u64>,
) -> 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,
completion_ledger,
identity_lease_changes,
));
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();
let _ = forwarder.task.await;
let _ = forwarder.identity_stream_health_task.await;
}
event_tx
}
#[cfg(test)]
pub(crate) async fn install_test_console_event_ingress(&self) -> TestConsoleEventIngress {
TestConsoleEventIngress(self.install_test_event_ingress().await)
}
async fn rollback_mob_runtime(
mob_runtime: MobRuntime,
startup_error: UnifiedRuntimeBootstrapError,
) -> Result<Self, UnifiedRuntimeBootstrapError> {
match mob_runtime.handle().stop().await {
Ok(_report) => 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>>,
Option<meerkat_core::comms::SessionEventEpoch>,
);
enum ForwardedAgentEvent {
Event(Box<TaggedAgentEvent>),
Closed(TrackedAgentEventStream, u64),
Abandoned {
key: TrackedAgentEventStream,
session_id: meerkat_core::types::SessionId,
generation: u64,
},
}
const PREDECESSOR_DRAIN_DEADLINE: Duration = Duration::from_secs(10);
#[derive(Debug, serde::Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
enum ForwarderStreamGap {
PredecessorStreamAbandoned { drain_deadline_ms: u64 },
}
fn predecessor_stream_gap_event(
key: &TrackedAgentEventStream,
session_id: &meerkat_core::types::SessionId,
generation: u64,
) -> ForwardedMemberEvent {
let timestamp_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|elapsed| u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX))
.unwrap_or_default();
let gap = ForwarderStreamGap::PredecessorStreamAbandoned {
drain_deadline_ms: u64::try_from(PREDECESSOR_DRAIN_DEADLINE.as_millis())
.unwrap_or(u64::MAX),
};
ForwardedMemberEvent {
envelope: EventEnvelope {
event_id: format!("evt-agent-gap-{session_id}-{timestamp_ms}-{generation}"),
source: "agent".to_string(),
timestamp_ms,
event: UnifiedEvent::Agent {
agent_id: crate::member_comms_id::runtime_event_alias(&key.runtime_id),
event_type: "stream_truncated".to_string(),
payload: Some(json!({
"reason": gap,
"session_id": session_id,
})),
},
},
alert: None,
}
}
#[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>;
#[derive(Default)]
struct AttachedStreams {
current: HashMap<TrackedAgentEventStream, StreamAttachment>,
departed: HashMap<u64, DepartedStream>,
forwarded: HashMap<MemberStreamOwner, Arc<std::sync::Mutex<ForwardedEventIds>>>,
}
const FORWARDED_EVENT_ID_WINDOW: usize = 2048;
#[derive(Default)]
struct ForwardedEventIds {
order: std::collections::VecDeque<uuid::Uuid>,
seen: HashSet<uuid::Uuid>,
}
impl ForwardedEventIds {
fn first_sighting(&mut self, event_id: uuid::Uuid) -> bool {
if !self.seen.insert(event_id) {
return false;
}
self.order.push_back(event_id);
while self.order.len() > FORWARDED_EVENT_ID_WINDOW {
if let Some(evicted) = self.order.pop_front() {
self.seen.remove(&evicted);
}
}
true
}
}
struct DepartedStream {
owner: MemberStreamOwner,
actor: meerkat_session::LiveSessionActorWitness,
attachment: StreamAttachment,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
struct MemberStreamOwner {
mob_id: String,
member_identity: AgentIdentity,
}
impl MemberStreamOwner {
fn of(key: &TrackedAgentEventStream) -> Self {
Self {
mob_id: key.mob_id.clone(),
member_identity: key.member_identity.clone(),
}
}
}
impl AttachedStreams {
fn close(&mut self, key: &TrackedAgentEventStream, generation: u64) -> bool {
self.departed.remove(&generation);
if self
.current
.get(key)
.is_some_and(|attachment| attachment.generation == generation)
{
self.current.remove(key);
true
} else {
false
}
}
fn depart(&mut self, key: &TrackedAgentEventStream, cut_unknown_actor: bool) {
let Some(attachment) = self.current.remove(key) else {
return;
};
if attachment.actor.is_none() && cut_unknown_actor {
attachment.abort.abort();
return;
}
if let Some(actor) = attachment.actor.clone() {
self.departed.insert(
attachment.generation,
DepartedStream {
owner: MemberStreamOwner::of(key),
actor,
attachment,
},
);
}
}
fn predecessor_draining(
&mut self,
owner: &MemberStreamOwner,
bound_session: Option<&meerkat_core::types::SessionId>,
now: tokio::time::Instant,
) -> bool {
let mut draining = false;
for departed in self.departed.values_mut() {
if departed.owner != *owner {
continue;
}
let predecessor = !departed.actor.is_live()
|| bound_session
.is_some_and(|session_id| departed.actor.session_id() != session_id);
if predecessor {
departed.attachment.hold(now);
draining = true;
}
}
draining
}
fn live_departed_for(
&self,
owner: &MemberStreamOwner,
session_id: &meerkat_core::types::SessionId,
) -> Option<u64> {
self.departed
.iter()
.find(|(_, departed)| {
departed.owner == *owner
&& departed.actor.is_live()
&& departed.actor.session_id() == session_id
})
.map(|(generation, _)| *generation)
}
fn forwarded_for(
&mut self,
owner: &MemberStreamOwner,
) -> Arc<std::sync::Mutex<ForwardedEventIds>> {
Arc::clone(self.forwarded.entry(owner.clone()).or_default())
}
fn retain_forwarded_for(&mut self, roster: &HashSet<TrackedAgentEventStream>) {
let Self {
current,
departed,
forwarded,
} = self;
forwarded.retain(|owner, _| {
roster
.iter()
.any(|key| MemberStreamOwner::of(key) == *owner)
|| current
.keys()
.any(|key| MemberStreamOwner::of(key) == *owner)
|| departed.values().any(|departed| departed.owner == *owner)
});
}
fn holds_departed(&self, owner: &MemberStreamOwner) -> bool {
self.departed
.values()
.any(|departed| departed.owner == *owner)
}
fn earliest_drain_deadline(&self) -> Option<tokio::time::Instant> {
self.current
.values()
.chain(self.departed.values().map(|departed| &departed.attachment))
.filter_map(StreamAttachment::drain_deadline)
.min()
}
}
fn earliest_reconcile_deadline(
subscribe_failures: &HashMap<TrackedAgentEventStream, SubscribeBackoff>,
tracked: &AttachedStreams,
) -> Option<tokio::time::Instant> {
earliest_backoff_attempt(subscribe_failures)
.into_iter()
.chain(tracked.earliest_drain_deadline())
.min()
}
struct StreamAttachment {
generation: u64,
attribution: Arc<std::sync::Mutex<StreamAttribution>>,
actor: Option<meerkat_session::LiveSessionActorWitness>,
abort: futures::stream::AbortHandle,
held_since: Option<tokio::time::Instant>,
}
impl StreamAttachment {
fn hold(&mut self, now: tokio::time::Instant) {
let since = *self.held_since.get_or_insert(now);
if now >= since + PREDECESSOR_DRAIN_DEADLINE && !self.abort.is_aborted() {
tracing::warn!(
identity = %lock_attribution(&self.attribution).key.member_identity,
deadline = ?PREDECESSOR_DRAIN_DEADLINE,
"mobkit agent event forwarder: a revoked predecessor's stream stayed open past \
its drain deadline; cutting it off with an explicit gap so the successor attaches"
);
self.abort.abort();
}
}
fn drain_deadline(&self) -> Option<tokio::time::Instant> {
self.held_since
.filter(|_| !self.abort.is_aborted())
.map(|since| since + PREDECESSOR_DRAIN_DEADLINE)
}
}
struct StreamAttribution {
key: TrackedAgentEventStream,
role: ProfileName,
durable_identity: Option<Arc<str>>,
}
impl StreamAttribution {
fn new(key: TrackedAgentEventStream, role: ProfileName) -> Self {
let durable_identity = key.durable_identity.as_deref().map(Arc::<str>::from);
Self {
key,
role,
durable_identity,
}
}
}
fn lock_attribution(
attribution: &std::sync::Mutex<StreamAttribution>,
) -> std::sync::MutexGuard<'_, StreamAttribution> {
attribution
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn attach_member_event_stream(
streams: &mut SelectAll<TaggedAgentEventStream>,
key: TrackedAgentEventStream,
role: ProfileName,
stream: EventStream,
actor: Option<meerkat_session::LiveSessionActorWitness>,
epoch: Option<meerkat_core::comms::SessionEventEpoch>,
forwarded: Option<Arc<std::sync::Mutex<ForwardedEventIds>>>,
) -> StreamAttachment {
static NEXT_GENERATION: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let generation = NEXT_GENERATION.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let attribution = Arc::new(std::sync::Mutex::new(StreamAttribution::new(key, role)));
let events = Arc::clone(&attribution);
let closing = Arc::clone(&attribution);
let stream = match forwarded {
Some(forwarded) => stream
.filter(move |envelope| {
let first = forwarded
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.first_sighting(envelope.event_id);
futures::future::ready(first)
})
.boxed(),
None => stream,
};
let (stream, abort) = futures::stream::abortable(stream);
let cut_off = abort.clone();
let session_id = actor.as_ref().map(|actor| actor.session_id().clone());
let mapped = stream
.map(move |envelope| {
let attribution = lock_attribution(&events);
ForwardedAgentEvent::Event(Box::new((
attribution.key.runtime_id.clone(),
attribution.key.fence_token,
attribution.role.clone(),
envelope,
attribution.durable_identity.clone(),
epoch,
)))
})
.chain(
futures::stream::once(async move {
let key = lock_attribution(&closing).key.clone();
let abandoned = match session_id {
Some(session_id) if cut_off.is_aborted() => {
Some(ForwardedAgentEvent::Abandoned {
key: key.clone(),
session_id,
generation,
})
}
_ => None,
};
futures::stream::iter(
abandoned
.into_iter()
.chain([ForwardedAgentEvent::Closed(key, generation)]),
)
})
.flatten(),
)
.boxed();
streams.push(mapped);
StreamAttachment {
generation,
attribution,
actor,
abort,
held_since: None,
}
}
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>>,
lease_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()),
lease_changes: None,
next_safety_deadline: tokio::time::Instant::now() + RECONCILE_SAFETY_INTERVAL,
}
}
fn with_lease_changes(mut self, lease_changes: tokio::sync::watch::Receiver<u64>) -> Self {
self.lease_changes = Some(lease_changes);
self
}
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,
lease_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 lease_change = async {
match lease_changes.as_mut() {
Some(rx) => rx.changed().await,
None => std::future::pending().await,
}
};
let (mob_set_closed, lease_closed) = tokio::select! {
() = machine_change => (false, false),
result = mob_set_change => (result.is_err(), false),
result = lease_change => (false, result.is_err()),
() = tokio::time::sleep_until(deadline) => (false, false),
};
if lease_closed {
self.lease_changes = None;
}
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())
}
#[derive(Default)]
struct HealthCompletionLedger {
sessions: std::sync::Mutex<HashMap<meerkat_core::types::SessionId, SessionCreditLedger>>,
}
#[derive(Default)]
struct SessionCreditLedger {
current: ResumeSpace,
high_water: HashMap<Option<meerkat_core::comms::SessionEventEpoch>, u64>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
enum ResumeSpace {
#[default]
Unclaimed,
Unnamed,
Named(meerkat_core::comms::SessionEventEpoch),
}
impl HealthCompletionLedger {
fn lock(
&self,
) -> std::sync::MutexGuard<'_, HashMap<meerkat_core::types::SessionId, SessionCreditLedger>>
{
self.sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn claim(
&self,
session_id: &meerkat_core::types::SessionId,
epoch: Option<meerkat_core::comms::SessionEventEpoch>,
seq: u64,
) -> bool {
if seq == 0 {
return false;
}
let mut sessions = self.lock();
let session = sessions.entry(session_id.clone()).or_default();
if !session.high_water.contains_key(&epoch) {
session.current = match epoch {
Some(epoch) => ResumeSpace::Named(epoch),
None => ResumeSpace::Unnamed,
};
}
let high_water = session.high_water.entry(epoch).or_default();
if seq <= *high_water {
return false;
}
*high_water = seq;
true
}
fn resume_cursor(
&self,
session_id: &meerkat_core::types::SessionId,
) -> meerkat_core::comms::SessionEventCursor {
let sessions = self.lock();
let Some(session) = sessions.get(session_id) else {
return meerkat_core::comms::SessionEventCursor::Earliest;
};
match session.current {
ResumeSpace::Named(epoch) => meerkat_core::comms::SessionEventCursor::After {
epoch,
seq: session.high_water.get(&Some(epoch)).copied().unwrap_or(0),
},
_ => meerkat_core::comms::SessionEventCursor::Earliest,
}
}
}
struct ReplayedCompletionDrain {
session_service: Arc<dyn meerkat_mob::MobSessionService>,
ledger: Arc<HealthCompletionLedger>,
}
#[async_trait::async_trait]
impl crate::identity_first::runtime::PendingCompletionDrain for ReplayedCompletionDrain {
async fn drain(
&self,
session_id: &meerkat_core::types::SessionId,
) -> crate::identity_first::runtime::PendingRunTerminals {
use futures::FutureExt as _;
use meerkat_core::comms::{SessionEventCursorRejection, StreamError};
let mut terminals = crate::identity_first::runtime::PendingRunTerminals::default();
let mut cursor = self.ledger.resume_cursor(session_id);
let subscription = loop {
match self
.session_service
.subscribe_agent_session_events_from(session_id, cursor)
.await
{
Ok(subscription) => break subscription,
Err(StreamError::CursorRejected {
reason: SessionEventCursorRejection::EpochMismatch { .. },
..
}) if cursor != meerkat_core::comms::SessionEventCursor::Earliest => {
cursor = meerkat_core::comms::SessionEventCursor::Earliest;
}
Err(_) => return terminals,
}
};
let epoch = subscription.epoch;
let mut stream = subscription.stream;
while let Some(Some(envelope)) = stream.next().now_or_never() {
let claimed = self.ledger.claim(session_id, epoch, envelope.seq);
match envelope.payload {
AgentEvent::RunCompleted { .. } if claimed => terminals.completed += 1,
AgentEvent::RunFailed { .. } if claimed => terminals.failed += 1,
_ => {}
}
}
terminals
}
}
async fn record_identity_turn_completion(
identity_runtime: &Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
ledger: &HealthCompletionLedger,
durable_identity: Option<&str>,
epoch: Option<meerkat_core::comms::SessionEventEpoch>,
envelope: &meerkat_core::event::EventEnvelope<AgentEvent>,
) -> bool {
let claim = || {
envelope
.source_session_id()
.is_some_and(|session_id| ledger.claim(session_id, epoch, envelope.seq))
};
let failed = match envelope.payload {
AgentEvent::RunCompleted { .. } => false,
AgentEvent::RunFailed { .. } => true,
_ => {
claim();
return false;
}
};
let authority = identity_runtime
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let identity = durable_identity.and_then(|label| {
crate::identity_first::AgentIdentity::parse(label)
.inspect_err(|error| {
tracing::debug!(
identity = %label,
error = %error,
"mobkit identity health monitor: run completion carried an unparseable durable identity"
);
})
.ok()
});
let (Some(authority), Some(identity)) = (authority, identity) else {
claim();
return false;
};
authority
.credit_observed_terminal(&identity, failed, claim)
.await
}
async fn trigger_identity_stream_repair(
handle: &MobHandle,
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;
}
};
match crate::identity_first::bridge::identity_actuation_custody(handle, durable_identity, None)
.await
{
crate::identity_first::bridge::CollisionCustody::IdentityFirstOwns
| crate::identity_first::bridge::CollisionCustody::LegacyRepairMayProveCustody => {}
crate::identity_first::bridge::CollisionCustody::MobMachineOwns => {
tracing::info!(
identity = %identity,
runtime_id = %tracked_key.runtime_id,
detail,
"mobkit agent event forwarder: durable intent wants this identity present, so \
MobMachine owns its materialization; not marking identity-first broken and \
leaving reattachment to the reconcile loop"
);
return;
}
crate::identity_first::bridge::CollisionCustody::Indeterminate => {
tracing::warn!(
identity = %identity,
runtime_id = %tracked_key.runtime_id,
detail,
"mobkit agent event forwarder: identity intent is unreadable or does not name \
this identity, so repair ownership cannot be established; declining to mark \
identity-first broken (degraded, retried on the next closure)"
);
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 = AttachedStreams::default();
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, _epoch) = *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::Abandoned { key, session_id, generation } => {
if event_tx
.send(predecessor_stream_gap_event(&key, &session_id, generation))
.await
.is_err()
{
break;
}
}
ForwardedAgentEvent::Closed(tracked_key, generation) => {
if tracked.close(&tracked_key, generation) {
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_reconcile_deadline(&subscribe_failures, &tracked)) => {
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>>>>,
completion_ledger: Arc<HealthCompletionLedger>,
lease_changes: tokio::sync::watch::Receiver<u64>,
) {
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
let mut tracked = AttachedStreams::default();
let mut subscribe_failures: HashMap<TrackedAgentEventStream, SubscribeBackoff> = HashMap::new();
let mut cadence = ReconcileCadence::new(&agent_mob_mcp_state).with_lease_changes(lease_changes);
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, epoch) = *event;
let _credited = record_identity_turn_completion(
&identity_runtime,
&completion_ledger,
durable_identity.as_deref(),
epoch,
&envelope,
).await;
}
ForwardedAgentEvent::Abandoned { .. } => {}
ForwardedAgentEvent::Closed(tracked_key, generation) => {
if tracked.close(&tracked_key, generation) {
subscribe_failures.remove(&tracked_key);
}
trigger_identity_stream_repair(
&handle,
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_reconcile_deadline(&subscribe_failures, &tracked)) => {
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 AttachedStreams,
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,
});
}
}
let departed: Vec<TrackedAgentEventStream> = tracked
.current
.keys()
.filter(|tracked_key| !current.contains(*tracked_key))
.cloned()
.collect();
for tracked_key in departed {
tracked.depart(&tracked_key, identity_runtime.is_some());
}
subscribe_failures.retain(|key, _| current.contains(key));
tracked.retain_forwarded_for(¤t);
let console = identity_runtime.is_none();
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,
identity_fencing_token,
fence_token,
};
if let Some(attachment) = tracked.current.get_mut(&tracked_key) {
if attachment
.actor
.as_ref()
.is_some_and(|actor| !actor.is_live())
{
attachment.hold(tokio::time::Instant::now());
}
continue;
}
if !forwarder_should_subscribe(entry.status) {
subscribe_failures.remove(&tracked_key);
continue;
}
let owner = MemberStreamOwner::of(&tracked_key);
if tracked.holds_departed(&owner) {
let current_session = handle.resolve_bridge_session_id(&identity).await;
if tracked.predecessor_draining(
&owner,
current_session.as_ref(),
tokio::time::Instant::now(),
) {
continue;
}
}
let bound_session = if tracked.holds_departed(&owner) {
match handle.resolve_bridge_session_id(&identity).await {
Some(candidate) if tracked.live_departed_for(&owner, &candidate).is_some() => {
match observe_bound_session(handle, &tracked_key).await {
Some(session_id) => Some(session_id),
None => continue,
}
}
_ => None,
}
} else {
None
};
if let Some(generation) = bound_session
.as_ref()
.and_then(|session_id| tracked.live_departed_for(&owner, session_id))
&& let Some(DepartedStream { attachment, .. }) =
tracked.departed.remove(&generation)
{
*lock_attribution(&attachment.attribution) =
StreamAttribution::new(tracked_key.clone(), entry.role.clone());
subscribe_failures.remove(&tracked_key);
tracked.current.insert(tracked_key, attachment);
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();
let replay = console
|| identity_runtime.is_some_and(|slot| {
slot.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
.is_some_and(|runtime| runtime.has_pending_completion_drain())
});
let subscription =
subscribe_agent_events_for_console_forwarder(handle, &identity, console, replay)
.await;
match subscription {
Ok((stream, actor, epoch)) => {
subscribe_failures.remove(&tracked_key);
let forwarded = console.then(|| tracked.forwarded_for(&owner));
let attachment = attach_member_event_stream(
streams,
tracked_key.clone(),
role,
stream,
actor,
epoch,
forwarded,
);
tracked.current.insert(tracked_key, attachment);
}
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(
handle,
&primary_mob_id,
&repair_key,
identity_runtime,
"live agent event stream remained unavailable after bounded retries",
)
.await;
}
}
}
}
}
handles
}
async fn observe_bound_session(
handle: &MobHandle,
key: &TrackedAgentEventStream,
) -> Option<meerkat_core::types::SessionId> {
let before = handle
.resolve_bridge_session_id(&key.member_identity)
.await?;
let bound = handle
.list_members_including_retiring()
.await
.into_iter()
.any(|entry| {
entry.agent_identity == key.member_identity
&& forwarder_should_subscribe(entry.status)
&& entry
.binding_atoms()
.is_some_and(|(runtime_id, fence_token)| {
runtime_id == key.runtime_id && fence_token == key.fence_token
})
});
let after = handle.resolve_bridge_session_id(&key.member_identity).await;
(bound && after.as_ref() == Some(&before)).then_some(before)
}
async fn subscribe_agent_events_for_console_forwarder(
handle: &MobHandle,
identity: &AgentIdentity,
console: bool,
replay: bool,
) -> Result<
(
EventStream,
Option<meerkat_session::LiveSessionActorWitness>,
Option<meerkat_core::comms::SessionEventEpoch>,
),
meerkat_mob::MobError,
> {
let cursor = if replay {
meerkat_core::comms::SessionEventCursor::Earliest
} else {
meerkat_core::comms::SessionEventCursor::Live
};
let subscription = handle.subscribe_agent_events_from(identity, cursor).await?;
let actor = if console { subscription.actor } else { None };
Ok((subscription.stream, actor, subscription.epoch))
}
async fn run_mob_events_subscription(
handle: MobHandle,
store: MobEventsStore,
persistent_metadata: Arc<dyn PersistentMetadataStore>,
idle_retire_overrides: Option<crate::mob_handle_runtime::ImplicitDelegateRetirementOverrides>,
) {
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 Some(overrides) = idle_retire_overrides.as_ref() {
overrides.observe_mob_event(&mob_id, &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)
}
fn actor_terminated_detail(result: &Result<MobState, MobError>) -> Option<String> {
let error = match result {
Ok(_) => return None,
Err(error) => error,
};
match error {
MobError::ActorCommandChannelClosed | MobError::ActorReplyChannelClosed => {
Some(error.to_string())
}
MobError::Internal(text)
if text.contains(&MobError::ActorCommandChannelClosed.to_string())
|| text.contains(&MobError::ActorReplyChannelClosed.to_string()) =>
{
Some(text.clone())
}
_ => None,
}
}
async fn run_actor_loop_probe<P, F>(
mut probe: P,
error_hook: SharedErrorHook,
health: Arc<crate::actor_loop_health::ActorLoopHealth>,
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;
health.mark_stalled(stall_id);
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;
if let Some(detail) = actor_terminated_detail(&result) {
health.mark_terminated(Some(stall_id), detail.clone());
tracing::error!(
stall_id,
stalled_for_secs = started.elapsed().as_secs(),
detail = %detail,
"mob actor loop TERMINATED while stall {stall_id} was open: the actor \
command channel closed, so the loop did not recover and cannot \
recover in this process; every delivery will now fail fast; the \
process must restart"
);
fire_error_hook(
&error_hook,
ErrorEvent::ActorLoopTerminated {
stall_id: Some(stall_id),
detail,
},
);
break;
}
resolved_stalls += 1;
health.mark_recovered(stall_id);
fire_error_hook(
&error_hook,
ErrorEvent::ActorLoopRecovered {
stall_id,
stalled_for_secs: started.elapsed().as_secs(),
},
);
result
}
};
if let Some(detail) = actor_terminated_detail(&result) {
health.mark_terminated(None, detail.clone());
tracing::error!(
detail = %detail,
"mob actor loop TERMINATED: the actor command channel closed; every delivery \
will now fail fast; the process must restart"
);
fire_error_hook(
&error_hook,
ErrorEvent::ActorLoopTerminated {
stall_id: None,
detail,
},
);
break;
}
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,
})
}
static SOURCE_EVENT_PROCESS_EPOCH: std::sync::LazyLock<String> =
std::sync::LazyLock::new(|| uuid::Uuid::new_v4().simple().to_string());
fn source_event_epoch(attributed: &AttributedEvent) -> String {
format!(
"{}.{}.{}",
*SOURCE_EVENT_PROCESS_EPOCH,
attributed.source.generation.get(),
attributed.source_fence_token.get()
)
}
fn attributed_event_to_unified(attributed: AttributedEvent) -> EventEnvelope<UnifiedEvent> {
let mut payload =
crate::mob_handle_runtime::console_agent_event_payload(&attributed.envelope.payload);
if let Some(object) = payload.as_object_mut() {
if let Some(session_id) = attributed.envelope.source.session_id() {
object.insert("session_id".to_string(), json!(session_id));
object.insert(
"source_sequence".to_string(),
json!(attributed.envelope.seq),
);
object.insert(
"source_epoch".to_string(),
json!(source_event_epoch(&attributed)),
);
}
}
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(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 {
assistant_message_id: None,
delta: "hello".to_string(),
},
},
}
}
#[test]
fn attributed_projection_preserves_exact_source_session_and_sequence() {
let session_id = meerkat_core::SessionId::new();
let mut attributed = attributed_text_delta("router", 1);
attributed.envelope.source =
meerkat_core::event::EventSourceIdentity::session(session_id.clone());
attributed.envelope.seq = 41;
let projected = attributed_event_to_unified(attributed);
let UnifiedEvent::Agent {
payload: Some(payload),
..
} = projected.event
else {
panic!("expected agent payload");
};
assert_eq!(payload["session_id"], json!(session_id));
assert_eq!(payload["source_sequence"], json!(41));
assert_eq!(payload["delta"], "hello");
let epoch = payload["source_epoch"].as_str().expect("source epoch");
assert!(epoch.ends_with(".1.1"), "{epoch}");
let mut respawned = attributed_text_delta("router", 2);
respawned.envelope.source = meerkat_core::event::EventSourceIdentity::session(session_id);
let UnifiedEvent::Agent {
payload: Some(respawned),
..
} = attributed_event_to_unified(respawned).event
else {
panic!("expected respawned agent payload");
};
assert_ne!(respawned["source_epoch"], payload["source_epoch"]);
let legacy = attributed_event_to_unified(attributed_text_delta("router", 1));
let UnifiedEvent::Agent {
payload: Some(payload),
..
} = legacy.event
else {
panic!("expected legacy agent payload");
};
assert!(payload.get("session_id").is_none());
assert!(payload.get("source_sequence").is_none());
assert!(payload.get("source_epoch").is_none());
}
#[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,
crate::actor_loop_health::ActorLoopHealth::shared(),
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,
crate::actor_loop_health::ActorLoopHealth::shared(),
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,
crate::actor_loop_health::ActorLoopHealth::shared(),
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,
crate::actor_loop_health::ActorLoopHealth::shared(),
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,
crate::actor_loop_health::ActorLoopHealth::shared(),
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();
}
#[tokio::test(start_paused = true)]
async fn channel_closed_on_parked_probe_is_termination_not_recovery() {
let (slot, captured) = capturing_error_hook_slot();
let health = crate::actor_loop_health::ActorLoopHealth::shared();
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);
async {
tokio::time::sleep(Duration::from_secs(50)).await;
Err(MobError::ActorCommandChannelClosed)
}
},
slot,
Arc::clone(&health),
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;
assert_eq!(
events.len(),
2,
"one stall page then one termination: {events:?}"
);
assert!(
matches!(
&events[0],
ErrorEvent::ActorLoopStalled {
stall_id: Some(1),
..
}
),
"the stall opens first: {events:?}"
);
match &events[1] {
ErrorEvent::ActorLoopTerminated { stall_id, detail } => {
assert_eq!(*stall_id, Some(1), "termination names the stall it closed");
assert!(
detail.contains("command channel closed"),
"detail carries the channel-closed evidence: {detail}"
);
}
other => panic!("expected ActorLoopTerminated, got {other:?}"),
}
assert!(
!events
.iter()
.any(|event| matches!(event, ErrorEvent::ActorLoopRecovered { .. })),
"a dead actor must NEVER be reported as recovered: {events:?}"
);
drop(events);
assert!(
matches!(
health.snapshot(),
crate::actor_loop_health::ActorLoopHealthState::Terminated {
stall_id: Some(1),
..
}
),
"shared health must read terminated: {:?}",
health.snapshot()
);
tokio::time::sleep(Duration::from_mins(10)).await;
assert_eq!(
probes.load(Ordering::SeqCst),
1,
"a terminated loop has nothing left to probe"
);
assert!(
probe_task.is_finished(),
"the probe task must end on termination"
);
}
#[test]
fn laundered_channel_closed_text_classifies_as_termination() {
let laundered = MobError::Internal(format!(
"{}; last actor-published phase is Running, so the phase watch is not terminal \
lifecycle authority",
MobError::ActorCommandChannelClosed
));
assert!(actor_terminated_detail(&Err(laundered)).is_some());
assert!(actor_terminated_detail(&Err(MobError::ActorReplyChannelClosed)).is_some());
assert!(actor_terminated_detail(&Err(MobError::ActorCommandChannelClosed)).is_some());
assert!(
actor_terminated_detail(&Err(MobError::Internal(
"scope denied for this command".to_string()
)))
.is_none()
);
assert!(actor_terminated_detail(&Ok(MobState::Running)).is_none());
}
#[tokio::test(start_paused = true)]
async fn channel_closed_on_fresh_probe_terminates_without_a_stall_id() {
let (slot, captured) = capturing_error_hook_slot();
let health = crate::actor_loop_health::ActorLoopHealth::shared();
let probe_task = tokio::spawn(run_actor_loop_probe(
|| std::future::ready(Err(MobError::ActorReplyChannelClosed)),
slot,
Arc::clone(&health),
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_secs(65)).await;
tokio::task::yield_now().await;
let events = captured.lock().await;
assert_eq!(events.len(), 1, "{events:?}");
assert!(
matches!(
&events[0],
ErrorEvent::ActorLoopTerminated { stall_id: None, .. }
),
"{events:?}"
);
drop(events);
assert!(health.snapshot().refuses_admission());
assert!(probe_task.is_finished());
}
#[tokio::test(start_paused = true)]
async fn health_reads_stalled_while_parked_and_live_after_recovery() {
let (slot, _captured) = capturing_error_hook_slot();
let health = crate::actor_loop_health::ActorLoopHealth::shared();
let probe_task = tokio::spawn(run_actor_loop_probe(
|| async {
tokio::time::sleep(Duration::from_secs(50)).await;
Ok(MobState::Running)
},
slot,
Arc::clone(&health),
Duration::from_mins(1),
Duration::from_secs(30),
));
tokio::time::sleep(Duration::from_secs(95)).await;
tokio::task::yield_now().await;
assert_eq!(
health.snapshot().open_stall_id(),
Some(1),
"the delivery path must see the open stall: {:?}",
health.snapshot()
);
tokio::time::sleep(Duration::from_secs(20)).await;
tokio::task::yield_now().await;
assert_eq!(
health.snapshot(),
crate::actor_loop_health::ActorLoopHealthState::Live,
"recovery must reopen the delivery path"
);
probe_task.abort();
}
#[test]
fn termination_is_logged_as_an_error_naming_the_restart() {
let terminated = ErrorEvent::ActorLoopTerminated {
stall_id: Some(4),
detail: "mob actor command channel closed".to_string(),
};
let ((), logged) = capture_tracing(|| log_error_event(&terminated, true));
assert!(logged.contains("ERROR"), "{logged}");
assert!(
logged.contains("process must restart") && logged.contains("NOT a recovery"),
"the operator must be told the loop is gone, not recovered: {logged}"
);
}
#[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
);
}
const IDLE_WORKER: &str = "worker";
async fn mob_with_idle_worker(
mob_id: &str,
temp_dir: &tempfile::TempDir,
) -> (MobRuntime, meerkat_core::types::SessionId) {
let definition = meerkat_mob::MobDefinition::from_toml(&format!(
"[mob]\nid = \"{mob_id}\"\n\n[profiles.worker]\nmodel = \"gpt-5.5\"\n\
external_addressable = true\n\n[profiles.worker.tools]\ncomms = true\n"
))
.expect("mob definition");
let spec = MobBootstrapSpec::ephemeral_runtime_backed_inner(
definition,
meerkat_mob::MobStorage::in_memory(),
temp_dir.path().to_path_buf(),
16,
None,
"test session store",
None,
None,
None,
None,
crate::mob_handle_runtime::CapabilityFlags::default(),
None,
None,
)
.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::default())),
});
let mob_runtime = MobRuntime::bootstrap(spec)
.await
.expect("bootstrap mob runtime");
let handle = mob_runtime.handle();
let mut member = SpawnMemberSpec::new(
ProfileName::from("worker"),
AgentIdentity::from(IDLE_WORKER),
);
member.runtime_mode = Some(meerkat_mob::MobRuntimeMode::TurnDriven);
handle.ensure_member(member).await.expect("seat worker");
let session_id = handle
.resolve_bridge_session_id(&AgentIdentity::from(IDLE_WORKER))
.await
.expect("worker session binding");
(mob_runtime, session_id)
}
async fn run_worker_turn(
handle: &MobHandle,
content: &str,
) -> Vec<meerkat_core::event::EventEnvelope<AgentEvent>> {
let mut probe = handle
.subscribe_agent_events(&AgentIdentity::from(IDLE_WORKER))
.await
.expect("probe subscription");
crate::mob_handle_runtime::send_message_on_mob(handle, IDLE_WORKER, content)
.await
.expect("send to worker");
let mut events = Vec::new();
tokio::time::timeout(crate::test_wait::STRUCTURAL_BACKSTOP, async {
while let Some(event) = probe.next().await {
let terminal = matches!(
event.payload,
AgentEvent::RunCompleted { .. } | AgentEvent::RunFailed { .. }
);
events.push(event);
if terminal {
break;
}
}
})
.await
.expect("worker turn reaches a terminal");
assert!(
matches!(
events.last().map(|event| &event.payload),
Some(AgentEvent::RunCompleted { .. })
),
"worker turn completes: {events:?}"
);
events
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn console_forwarder_delivers_runs_that_started_before_it_attached() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let (mob_runtime, session_id) =
mob_with_idle_worker("console-replay-first-run", &temp_dir).await;
let handle = mob_runtime.handle();
let first_run = run_worker_turn(&handle, "probe-1").await;
let first_terminal_seq = first_run.last().expect("first run events").seq;
let module_runtime = std::thread::spawn(|| {
start_mobkit_runtime_with_options(
MobKitConfig {
modules: vec![],
discovery: crate::types::DiscoverySpec {
namespace: "console-replay-first-run".to_string(),
modules: vec![],
},
pre_spawn: vec![],
},
Vec::new(),
Duration::from_secs(2),
RuntimeOptions::default(),
)
})
.join()
.expect("module runtime thread")
.expect("module runtime");
let runtime = UnifiedRuntime::from_parts(
mob_runtime,
module_runtime,
Arc::new(InMemoryMetadataStore::new()),
)
.await;
let second_run = run_worker_turn(&handle, "probe-2").await;
let second_terminal_seq = second_run.last().expect("second run events").seq;
let worker_frames =
async || -> Vec<crate::console_contracts::ConsoleIdentityEventEnvelope> {
runtime
.drain_mob_agent_events()
.await
.expect("drain member events");
runtime
.console_events()
.replay_all(None)
.await
.expect("console replay")
.into_iter()
.filter(|frame| frame.identity == IDLE_WORKER)
.collect()
};
crate::test_wait::poll_until(
"probe-2's terminal reaches the console timeline",
crate::test_wait::STRUCTURAL_BACKSTOP,
async || {
worker_frames().await.iter().any(|frame| {
frame.event_type == "interaction_complete"
&& frame.data.get("session_id") == Some(&json!(session_id))
&& frame.data.get("source_sequence") == Some(&json!(second_terminal_seq))
})
},
)
.await;
let frames = worker_frames().await;
let run_started_prompts: Vec<Option<String>> = frames
.iter()
.filter(|frame| frame.event_type == "run_started")
.map(|frame| {
frame
.data
.get("input")
.cloned()
.and_then(|input| {
serde_json::from_value::<meerkat_core::types::RunInput>(input).ok()
})
.and_then(|input| input.prompt_text())
})
.collect();
assert_eq!(
run_started_prompts,
vec![Some("probe-1".to_string()), Some("probe-2".to_string())],
"both runs start on the console timeline, in order"
);
let mut first_run_sequences: Vec<u64> = frames
.iter()
.filter(|frame| frame.data.get("session_id") == Some(&json!(session_id)))
.filter_map(|frame| {
frame
.data
.get("source_sequence")
.and_then(serde_json::Value::as_u64)
})
.filter(|seq| *seq <= first_terminal_seq)
.collect();
first_run_sequences.sort_unstable();
assert_eq!(
first_run_sequences,
(1..=first_terminal_seq).collect::<Vec<_>>(),
"the pre-attach run arrives complete from meerkat's first sequence"
);
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn reconcile_attaches_from_the_actors_first_event_after_a_finished_run() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let (mob_runtime, _session_id) =
mob_with_idle_worker("console-replay-backoff", &temp_dir).await;
let handle = mob_runtime.handle();
run_worker_turn(&handle, "probe-1").await;
let entry = handle
.list_members_including_retiring()
.await
.into_iter()
.find(|entry| entry.agent_identity == IDLE_WORKER)
.expect("worker roster entry");
let (runtime_id, fence_token) = entry.binding_atoms().expect("worker binding atoms");
let key = TrackedAgentEventStream {
mob_id: handle.mob_id().to_string(),
durable_identity: durable_identity_label(&entry.labels),
member_identity: entry.agent_identity.clone(),
runtime_id,
identity_fencing_token: None,
fence_token,
};
let mut tracked = AttachedStreams::default();
let mut subscribe_failures = HashMap::from([(
key.clone(),
SubscribeBackoff {
next_attempt: tokio::time::Instant::now(),
consecutive_failures: 1,
},
)]);
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut subscribe_failures,
&mut streams,
None,
))
.await;
assert!(
tracked
.current
.get(&key)
.is_some_and(|attachment| attachment.actor.is_some()),
"the subscription names the member's actor"
);
assert!(
subscribe_failures.is_empty(),
"the subscription clears the lapsed backoff"
);
let first = tokio::time::timeout(Duration::from_secs(5), streams.next())
.await
.expect("attached stream yields")
.expect("attached stream open");
let ForwardedAgentEvent::Event(event) = first else {
panic!("attached stream closed before yielding the first run");
};
let (_, _, _, envelope, _, _) = *event;
assert!(
matches!(envelope.payload, AgentEvent::RunStarted { .. }),
"the stream starts at the finished run's start, got {:?}",
envelope.payload
);
assert_eq!(envelope.seq, 1);
drop(streams);
let _ = handle.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn reconcile_rekeys_a_stream_onto_a_same_actor_rebinding() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let (mob_runtime, _session_id) =
mob_with_idle_worker("console-replay-rekey", &temp_dir).await;
let handle = mob_runtime.handle();
let mut tracked = AttachedStreams::default();
let mut subscribe_failures = HashMap::new();
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut subscribe_failures,
&mut streams,
None,
))
.await;
let (key, attachment) = tracked
.current
.drain()
.next()
.expect("the worker's stream attached");
assert!(attachment.actor.is_some());
let generation = attachment.generation;
let mut previous = key.clone();
previous.fence_token = FenceToken::new(key.fence_token.get() + 1_000);
{
let mut attribution = lock_attribution(&attachment.attribution);
let role = attribution.role.clone();
*attribution = StreamAttribution::new(previous.clone(), role);
}
tracked.current.insert(previous, attachment);
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut subscribe_failures,
&mut streams,
None,
))
.await;
assert_eq!(tracked.current.len(), 1);
let rekeyed = tracked
.current
.get(&key)
.expect("the live binding tracks the stream");
assert_eq!(
rekeyed.generation, generation,
"the same stream was re-keyed, not a second subscription opened"
);
assert_eq!(streams.len(), 1, "one stream to the actor");
run_worker_turn(&handle, "probe-1").await;
let first = tokio::time::timeout(crate::test_wait::STRUCTURAL_BACKSTOP, streams.next())
.await
.expect("re-keyed stream yields")
.expect("re-keyed stream open");
let ForwardedAgentEvent::Event(event) = first else {
panic!("re-keyed stream closed");
};
let (runtime_id, fence_token, _, envelope, _, _) = *event;
assert!(matches!(envelope.payload, AgentEvent::RunStarted { .. }));
assert_eq!(
(runtime_id, fence_token),
(key.runtime_id.clone(), key.fence_token),
"events carry the new binding"
);
drop(streams);
let _ = handle.shutdown().await;
}
const RESTART_MEMBER: &str = "lead-1";
async fn boot_identity_first_persistent(
mob_id: &str,
mob_path: &std::path::Path,
state_root: &std::path::Path,
) -> (UnifiedRuntime, Arc<crate::identity_first::IdentityRuntime>) {
let boot = boot_persistent_before_activation(mob_id, mob_path, state_root).await;
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
(runtime, identity_runtime)
}
struct PersistentBootBeforeActivation {
runtime: UnifiedRuntime,
identity_runtime: Arc<crate::identity_first::IdentityRuntime>,
context: Arc<crate::identity_first::IdentityFirstRuntimeContext>,
roster: Vec<crate::identity_first::DurableAgentSpec>,
}
impl PersistentBootBeforeActivation {
async fn activate(self) -> UnifiedRuntime {
let Self {
mut runtime,
context,
roster,
..
} = self;
runtime
.install_and_bootstrap_identity_first_context(context, &roster)
.await
.expect("activate and restore identity-first runtime");
runtime
}
}
async fn boot_persistent_before_activation(
mob_id: &str,
mob_path: &std::path::Path,
state_root: &std::path::Path,
) -> PersistentBootBeforeActivation {
use crate::identity_first::{
AgentAddressability, AgentRuntimeServices, ContinuityStore, DurabilityPolicy,
DurableAgentSpec, IdentityFirstRuntimeContext, IdentityRuntime, IdentityRuntimeConfig,
LocalContinuityStore, LocalLeaseProvider, MobSessionBridge, MutableRosterProvider,
};
std::fs::create_dir_all(state_root).expect("state root");
let definition = meerkat_mob::MobDefinition::from_toml(&format!(
"[mob]\nid = \"{mob_id}\"\n\n[profiles.lead]\nmodel = \"gpt-5.5\"\n\
external_addressable = true\nruntime_mode = \"turn_driven\"\n\n\
[profiles.lead.tools]\ncomms = true\n"
))
.expect("mob definition");
let session_store = Arc::new(
meerkat_store::SqliteSessionStore::open(state_root.join("sessions.sqlite3"))
.expect("open session store"),
);
let (storage, provenance) =
crate::mob_composition_manifest::persistent_mob_storage(mob_path.to_path_buf())
.expect("open persistent mob storage");
let spec = MobBootstrapSpec::persistent(
definition,
storage,
state_root.to_path_buf(),
16,
session_store,
)
.expect("compose persistent MobKit stores")
.with_mob_storage_provenance(provenance)
.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::default())),
});
let runtime = UnifiedRuntime::bootstrap(
spec,
MobKitConfig {
modules: Vec::new(),
discovery: crate::types::DiscoverySpec {
namespace: mob_id.to_string(),
modules: Vec::new(),
},
pre_spawn: Vec::new(),
},
Duration::from_secs(2),
)
.await
.expect("bootstrap unified runtime");
let roster = vec![DurableAgentSpec {
identity: crate::identity_first::AgentIdentity::parse(RESTART_MEMBER)
.expect("identity"),
profile: ProfileName::from("lead"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
placement: None,
}];
let continuity_store = Arc::new(
LocalContinuityStore::open(state_root.join("identity-continuity.sqlite3"))
.expect("open identity continuity store"),
);
let identity_runtime = Arc::new(
IdentityRuntime::new(IdentityRuntimeConfig {
continuity_store: continuity_store as Arc<dyn ContinuityStore>,
lease_provider: Arc::new(LocalLeaseProvider::new()),
runtime_instance_id: format!("{mob_id}-instance"),
has_runtime_store: true,
durability_policy: DurabilityPolicy::SyncWriteThrough,
bridge: Some(Arc::new(MobSessionBridge::with_session_service(
runtime.mob_handle(),
runtime
.mob_runtime()
.session_service()
.cloned()
.expect("persistent runtime has a session service"),
))),
default_timeout: None,
})
.with_runtime_services(AgentRuntimeServices::new(runtime.mob_handle())),
);
let context = Arc::new(IdentityFirstRuntimeContext::new(
identity_runtime.clone(),
Arc::new(MutableRosterProvider::new(roster.clone())),
None,
None,
Some(runtime.mob_handle().definition().clone()),
));
PersistentBootBeforeActivation {
runtime,
identity_runtime,
context,
roster,
}
}
async fn commit_restart_member_turn(
identity_runtime: &crate::identity_first::IdentityRuntime,
prompt: &str,
) -> meerkat_core::types::SessionId {
let identity =
crate::identity_first::AgentIdentity::parse(RESTART_MEMBER).expect("identity");
identity_runtime
.send_awaiting_commit(
&identity,
&meerkat_core::ContentInput::Text(prompt.to_string()),
)
.await
.expect("complete member turn");
identity_runtime
.status(&identity)
.await
.expect("member status")
.session_id
.expect("member session")
}
#[derive(Clone, Copy)]
enum RestartShape {
Live,
Retired,
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn console_timeline_carries_a_revived_members_first_run_after_restart() {
console_timeline_after_restart("console-replay-restart", RestartShape::Live).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn console_timeline_carries_an_archived_members_first_run_after_restart() {
console_timeline_after_restart("console-replay-restart-archived", RestartShape::Retired)
.await;
}
async fn health_restart(
mob_id: &str,
temp: &tempfile::TempDir,
) -> (
UnifiedRuntime,
Arc<crate::identity_first::IdentityRuntime>,
Arc<HealthCompletionLedger>,
crate::identity_first::AgentIdentity,
) {
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
commit_restart_member_turn(&first_identity, "before restart").await;
first.shutdown().await;
let boot = boot_persistent_before_activation(mob_id, &mob_path, &state_root).await;
let _ingress = boot.runtime.install_test_event_ingress().await;
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
let ledger = Arc::clone(&runtime.identity_completion_ledger);
let identity =
crate::identity_first::AgentIdentity::parse(RESTART_MEMBER).expect("identity");
assert!(
identity_runtime
.status(&identity)
.await
.expect("status")
.lease
.is_some(),
"the lease landed"
);
(runtime, identity_runtime, ledger, identity)
}
fn drive_health_streams(
mut streams: SelectAll<TaggedAgentEventStream>,
slot: Arc<std::sync::RwLock<Option<Arc<crate::identity_first::IdentityRuntime>>>>,
ledger: Arc<HealthCompletionLedger>,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
while let Some(forwarded) = streams.next().await {
if let ForwardedAgentEvent::Event(event) = forwarded {
let (_, _, _, envelope, durable_identity, epoch) = *event;
let _credited = record_identity_turn_completion(
&slot,
&ledger,
durable_identity.as_deref(),
epoch,
&envelope,
)
.await;
}
}
})
}
struct NoPendingCompletions;
#[async_trait::async_trait]
impl crate::identity_first::runtime::PendingCompletionDrain for NoPendingCompletions {
async fn drain(
&self,
_session_id: &meerkat_core::types::SessionId,
) -> crate::identity_first::runtime::PendingRunTerminals {
crate::identity_first::runtime::PendingRunTerminals::default()
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn the_completion_drain_is_installed_before_the_lease_observer() {
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, _ledger, _identity) =
health_restart("identity-health-drain-order", &temp).await;
assert!(identity_runtime.has_pending_completion_drain());
assert!(
identity_runtime.lease_observer_installed_after_drain(),
"the lease observer saw the drain installed"
);
runtime.shutdown().await;
}
struct FailingOnceDrain {
armed: std::sync::atomic::AtomicBool,
calls: tokio::sync::watch::Sender<u64>,
}
#[async_trait::async_trait]
impl crate::identity_first::runtime::PendingCompletionDrain for FailingOnceDrain {
async fn drain(
&self,
_session_id: &meerkat_core::types::SessionId,
) -> crate::identity_first::runtime::PendingRunTerminals {
self.calls.send_modify(|calls| *calls += 1);
crate::identity_first::runtime::PendingRunTerminals {
completed: 0,
failed: u64::from(self.armed.swap(false, std::sync::atomic::Ordering::AcqRel)),
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_drained_failed_run_ends_a_completion_wait_typed() {
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, _ledger, identity) =
health_restart("identity-health-drained-failure", &temp).await;
let (calls, mut observed) = tokio::sync::watch::channel(0_u64);
let drain = Arc::new(FailingOnceDrain {
armed: std::sync::atomic::AtomicBool::new(false),
calls,
});
identity_runtime.install_pending_completion_drain(drain.clone());
let baseline = identity_runtime.completion_cursor(&identity).await;
let calls_before_wait = *observed.borrow_and_update();
let wait = {
let identity_runtime = identity_runtime.clone();
let identity = identity.clone();
tokio::spawn(async move {
identity_runtime
.await_completion(&identity, Some(baseline), Duration::from_secs(10))
.await
})
};
observed
.wait_for(|calls| *calls >= calls_before_wait + 2)
.await
.expect("drain calls observed");
drain
.armed
.store(true, std::sync::atomic::Ordering::Release);
identity_runtime.completion_cursor(&identity).await;
let outcome = tokio::time::timeout(Duration::from_secs(5), wait)
.await
.expect("the wait ends promptly")
.expect("wait task");
assert!(
matches!(outcome, crate::identity_first::CompletionWait::RunFailed(_)),
"{outcome:?}"
);
runtime.shutdown().await;
}
#[test]
fn health_completion_ledger_claims_each_sequence_once_per_space() {
use meerkat_core::comms::{SessionEventCursor, SessionEventEpoch};
let ledger = HealthCompletionLedger::default();
let session = meerkat_core::types::SessionId::new();
let space_epoch = SessionEventEpoch::new();
let space = Some(space_epoch);
assert_eq!(ledger.resume_cursor(&session), SessionEventCursor::Earliest);
assert!(
!ledger.claim(&session, space, 0),
"gap markers are never claimed"
);
assert!(ledger.claim(&session, space, 1));
assert!(ledger.claim(&session, space, 3));
assert!(!ledger.claim(&session, space, 2), "below the high-water");
assert!(!ledger.claim(&session, space, 3), "claimed once");
assert_eq!(
ledger.resume_cursor(&session),
SessionEventCursor::After {
epoch: space_epoch,
seq: 3,
}
);
let restarted_epoch = SessionEventEpoch::new();
let restarted = Some(restarted_epoch);
assert!(ledger.claim(&session, restarted, 1));
assert!(!ledger.claim(&session, restarted, 1));
assert_eq!(
ledger.resume_cursor(&session),
SessionEventCursor::After {
epoch: restarted_epoch,
seq: 1,
}
);
assert!(
!ledger.claim(&session, space, 3),
"credited in its own space"
);
assert!(!ledger.claim(&session, space, 1));
assert!(
ledger.claim(&session, space, 4),
"only past its own high-water"
);
assert_eq!(
ledger.resume_cursor(&session),
SessionEventCursor::After {
epoch: restarted_epoch,
seq: 1,
}
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn identity_health_drain_and_monitor_race_credits_each_turn_once() {
const TURNS: u64 = 4;
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, ledger, identity) =
health_restart("identity-health-race", &temp).await;
let slot = Arc::new(std::sync::RwLock::new(Some(identity_runtime.clone())));
let handle = runtime.mob_handle();
let mut tracked = AttachedStreams::default();
let mut failures = HashMap::new();
let mut streams = SelectAll::new();
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut failures,
&mut streams,
Some(&slot),
))
.await;
assert!(!tracked.current.is_empty(), "the monitor attached");
let before = identity_runtime.completion_cursor(&identity).await;
let monitor = drive_health_streams(streams, Arc::clone(&slot), Arc::clone(&ledger));
let committed = Arc::new(std::sync::atomic::AtomicU64::new(0));
let reader = {
let identity_runtime = identity_runtime.clone();
let identity = identity.clone();
let committed = Arc::clone(&committed);
tokio::spawn(async move {
let mut last = before;
loop {
let cursor = identity_runtime.completion_cursor(&identity).await;
assert!(cursor >= last, "the cursor never rewinds");
let done = committed.load(std::sync::atomic::Ordering::Acquire);
assert!(
cursor <= (0..=done).fold(before, |cursor, _| cursor.advanced()),
"a read never runs ahead of completed turns"
);
last = cursor;
tokio::task::yield_now().await;
}
})
};
for turn in 0..TURNS {
commit_restart_member_turn(&identity_runtime, &format!("race turn {turn}")).await;
committed.fetch_add(1, std::sync::atomic::Ordering::Release);
}
let expected = (0..TURNS).fold(before, |cursor, _| cursor.advanced());
identity_runtime
.wait_for_completion(
&identity,
(0..TURNS - 1).fold(before, |cursor, _| cursor.advanced()),
Duration::from_secs(10),
)
.await
.expect("the last turn is credited");
assert_eq!(
identity_runtime.completion_cursor(&identity).await,
expected,
"each turn credited exactly once"
);
reader.abort();
assert!(
reader.await.err().is_some_and(|error| error.is_cancelled()),
"the reader's invariants held"
);
monitor.abort();
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn identity_health_lease_rotation_replay_credits_nothing_twice() {
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, ledger, identity) =
health_restart("identity-health-rotation", &temp).await;
commit_restart_member_turn(&identity_runtime, "before rotation").await;
let slot = Arc::new(std::sync::RwLock::new(Some(identity_runtime.clone())));
let handle = runtime.mob_handle();
let mut tracked = AttachedStreams::default();
let mut failures = HashMap::new();
let mut streams = SelectAll::new();
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut failures,
&mut streams,
Some(&slot),
))
.await;
let (key, attachment) = tracked
.current
.drain()
.next()
.expect("the monitor attached");
let mut previous = key.clone();
previous.identity_fencing_token = key.identity_fencing_token.map(|token| token + 1_000);
tracked.current.insert(previous, attachment);
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut failures,
&mut streams,
Some(&slot),
))
.await;
assert!(
tracked.current.contains_key(&key),
"a new stream serves the rotated lease"
);
let baseline = identity_runtime.completion_cursor(&identity).await;
let monitor = drive_health_streams(streams, Arc::clone(&slot), Arc::clone(&ledger));
assert!(
identity_runtime
.wait_for_completion(&identity, baseline, Duration::from_millis(300))
.await
.is_err(),
"neither the cut-off stream nor the new stream's replay credits a turn again"
);
commit_restart_member_turn(&identity_runtime, "after rotation").await;
let completed = identity_runtime
.wait_for_completion(&identity, baseline, Duration::from_secs(5))
.await
.expect("the next turn completes the wait");
assert_eq!(completed, baseline.advanced(), "counted exactly once");
monitor.abort();
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn identity_health_credits_a_window_turn_before_any_later_baseline() {
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, ledger, identity) =
health_restart("identity-health-window", &temp).await;
let before = identity_runtime.completion_cursor(&identity).await;
commit_restart_member_turn(&identity_runtime, "in the window").await;
let baseline = identity_runtime.completion_cursor(&identity).await;
assert_eq!(
baseline,
before.advanced(),
"the read credited the window turn first, exactly once"
);
let slot = Arc::new(std::sync::RwLock::new(Some(identity_runtime.clone())));
let handle = runtime.mob_handle();
let mut tracked = AttachedStreams::default();
let mut failures = HashMap::new();
let mut streams = SelectAll::new();
Box::pin(reconcile_agent_event_streams(
&handle,
&None,
&mut tracked,
&mut failures,
&mut streams,
Some(&slot),
))
.await;
assert!(!tracked.current.is_empty(), "the monitor attached");
let monitor = drive_health_streams(streams, Arc::clone(&slot), Arc::clone(&ledger));
assert!(
identity_runtime
.wait_for_completion(&identity, baseline, Duration::from_millis(300))
.await
.is_err(),
"the replayed window turn is not credited twice, and never satisfies a newer baseline"
);
commit_restart_member_turn(&identity_runtime, "after the attach").await;
let completed = identity_runtime
.wait_for_completion(&identity, baseline, Duration::from_secs(5))
.await
.expect("the newer turn completes the wait");
assert_eq!(completed, baseline.advanced(), "counted exactly once");
assert_eq!(
identity_runtime.completion_cursor(&identity).await,
baseline.advanced(),
"the monitor's replay and the drain never both credit it"
);
monitor.abort();
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn the_health_monitor_attaches_on_the_lease_wake() {
let temp = tempfile::tempdir().expect("temp dir");
let (runtime, identity_runtime, ledger, identity) =
health_restart("identity-health-wake", &temp).await;
identity_runtime.install_pending_completion_drain(Arc::new(NoPendingCompletions));
let slot = Arc::new(std::sync::RwLock::new(None));
let (lease_wake, lease_changes) = tokio::sync::watch::channel(0_u64);
let monitor = tokio::spawn(run_identity_stream_health_monitor(
runtime.mob_handle(),
None,
Arc::clone(&slot),
ledger,
lease_changes,
));
let baseline = identity_runtime.completion_cursor(&identity).await;
tokio::task::yield_now().await;
*slot
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(identity_runtime.clone());
identity_runtime.install_lease_observer(lease_wake);
commit_restart_member_turn(&identity_runtime, "after the wake").await;
identity_runtime
.wait_for_completion(&identity, baseline, Duration::from_secs(10))
.await
.expect("the lease wake attached the monitor, well inside the safety tick");
monitor.abort();
runtime.shutdown().await;
}
async fn console_timeline_after_restart(mob_id: &str, shape: RestartShape) {
let temp = tempfile::tempdir().expect("temp dir");
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let session_id = commit_restart_member_turn(&first_identity, "before restart").await;
if let RestartShape::Retired = shape {
let handle = first.mob_handle();
let members = handle.list_members().await;
assert_eq!(members.len(), 1, "one durable member");
let member = members[0].agent_identity.clone();
handle
.retire(member.clone())
.await
.expect("mob-plane retire archives the session");
crate::test_wait::poll_until(
"the mob-plane retire finalizes",
crate::test_wait::STRUCTURAL_BACKSTOP,
async || {
!handle
.list_members()
.await
.iter()
.any(|entry| entry.agent_identity == member && !entry.is_final)
},
)
.await;
}
first.shutdown().await;
let (second, second_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let revived_session = commit_restart_member_turn(&second_identity, "after restart").await;
assert_eq!(
revived_session, session_id,
"the restore resumed the same session"
);
let session_frames =
async || -> Vec<crate::console_contracts::ConsoleIdentityEventEnvelope> {
second
.drain_mob_agent_events()
.await
.expect("drain member events");
second
.console_events()
.replay_all(None)
.await
.expect("console replay")
.into_iter()
.filter(|frame| frame.data.get("session_id") == Some(&json!(session_id)))
.collect()
};
let count_at = |frames: &[crate::console_contracts::ConsoleIdentityEventEnvelope],
event_type: &str,
seq: u64| {
frames
.iter()
.filter(|frame| {
frame.event_type == event_type
&& frame.data.get("source_sequence") == Some(&json!(seq))
})
.count()
};
crate::test_wait::poll_until(
"the revived member's first run starts on the console timeline",
crate::test_wait::STRUCTURAL_BACKSTOP,
async || {
let frames = session_frames().await;
count_at(&frames, "run_started", 1) > 0 && count_at(&frames, "turn_started", 2) > 0
},
)
.await;
let frames = session_frames().await;
assert_eq!(count_at(&frames, "run_started", 1), 1, "{frames:#?}");
assert_eq!(count_at(&frames, "turn_started", 2), 1, "{frames:#?}");
second.shutdown().await;
}
struct DrivenForwarder {
tracked: AttachedStreams,
subscribe_failures: HashMap<TrackedAgentEventStream, SubscribeBackoff>,
streams: SelectAll<TaggedAgentEventStream>,
}
impl DrivenForwarder {
fn new() -> Self {
Self {
tracked: AttachedStreams::default(),
subscribe_failures: HashMap::new(),
streams: SelectAll::new(),
}
}
async fn reconcile(
&mut self,
handle: &MobHandle,
agent_mob_mcp_state: &Option<Arc<meerkat_mob_mcp::MobMcpState>>,
) {
Box::pin(reconcile_agent_event_streams(
handle,
agent_mob_mcp_state,
&mut self.tracked,
&mut self.subscribe_failures,
&mut self.streams,
None,
))
.await;
}
async fn next(&mut self) -> ForwardedAgentEvent {
tokio::time::timeout(crate::test_wait::STRUCTURAL_BACKSTOP, self.streams.next())
.await
.expect("the forwarder's streams yield")
.expect("the forwarder holds a stream")
}
async fn forward_terminals(
&mut self,
runtime: &UnifiedRuntime,
ingress: &Sender<ForwardedMemberEvent>,
terminals: usize,
) -> Vec<ForwarderStep> {
let mut seen = 0;
self.forward_until(runtime, ingress, |payload| {
if matches!(
payload,
AgentEvent::RunCompleted { .. } | AgentEvent::RunFailed { .. }
) {
seen += 1;
}
seen == terminals
})
.await
}
async fn forward_until(
&mut self,
runtime: &UnifiedRuntime,
ingress: &Sender<ForwardedMemberEvent>,
mut done: impl FnMut(&AgentEvent) -> bool,
) -> Vec<ForwarderStep> {
let handle = runtime.mob_handle();
let mut steps = Vec::new();
let mut finished = false;
while !finished {
match self.next().await {
ForwardedAgentEvent::Event(event) => {
let (source, source_fence_token, role, envelope, _, _) = *event;
finished = done(&envelope.payload);
steps.push(ForwarderStep::Event(Box::new(envelope.clone())));
ingress
.send(forwarded_member_event(AttributedEvent {
source,
source_fence_token,
role,
envelope,
}))
.await
.expect("console ingress open");
runtime
.drain_mob_agent_events()
.await
.expect("drain member events");
}
ForwardedAgentEvent::Abandoned {
key,
session_id,
generation,
} => {
steps.push(ForwarderStep::Abandoned { generation });
ingress
.send(predecessor_stream_gap_event(&key, &session_id, generation))
.await
.expect("console ingress open");
runtime
.drain_mob_agent_events()
.await
.expect("drain member events");
}
ForwardedAgentEvent::Closed(key, generation) => {
let served = self.tracked.close(&key, generation);
if served {
self.subscribe_failures.remove(&key);
}
steps.push(ForwarderStep::Closed { generation, served });
self.reconcile(&handle, &None).await;
}
}
}
steps
}
}
#[derive(Debug)]
enum ForwarderStep {
Event(Box<meerkat_core::event::EventEnvelope<AgentEvent>>),
Closed { generation: u64, served: bool },
Abandoned { generation: u64 },
}
fn restart_member_key(
entry: &meerkat_mob::runtime::MobMemberListEntry,
mob_id: &str,
) -> TrackedAgentEventStream {
let (runtime_id, fence_token) = entry.binding_atoms().expect("member binding atoms");
TrackedAgentEventStream {
mob_id: mob_id.to_string(),
durable_identity: durable_identity_label(&entry.labels),
member_identity: entry.agent_identity.clone(),
runtime_id,
identity_fencing_token: None,
fence_token,
}
}
async fn restart_member_entry(handle: &MobHandle) -> meerkat_mob::runtime::MobMemberListEntry {
let mut entries = handle.list_members_including_retiring().await;
assert_eq!(entries.len(), 1, "one durable member");
entries.remove(0)
}
async fn session_timeline(
runtime: &UnifiedRuntime,
session_id: &meerkat_core::types::SessionId,
) -> Vec<crate::console_contracts::ConsoleIdentityEventEnvelope> {
runtime
.console_events()
.replay_all(None)
.await
.expect("console replay")
.into_iter()
.filter(|frame| frame.data.get("session_id") == Some(&json!(session_id)))
.collect()
}
fn run_started_prompt(
frame: &crate::console_contracts::ConsoleIdentityEventEnvelope,
) -> Option<String> {
frame
.data
.get("input")
.cloned()
.and_then(|input| serde_json::from_value::<meerkat_core::types::RunInput>(input).ok())
.and_then(|input| input.prompt_text())
}
fn assert_one_attributed_run(
frames: &[crate::console_contracts::ConsoleIdentityEventEnvelope],
prompt: &str,
) {
let types: Vec<&str> = frames
.iter()
.map(|frame| frame.event_type.as_str())
.collect();
assert_eq!(types.first(), Some(&"run_started"), "timeline: {types:?}");
assert_eq!(types.get(1), Some(&"turn_started"), "timeline: {types:?}");
assert_eq!(
types.last(),
Some(&"interaction_complete"),
"timeline: {types:?}"
);
assert_eq!(
types.iter().filter(|kind| **kind == "run_started").count(),
1,
"one run: {types:?}"
);
assert_eq!(run_started_prompt(&frames[0]).as_deref(), Some(prompt));
let sequences: Vec<u64> = frames
.iter()
.map(|frame| {
frame
.data
.get("source_sequence")
.and_then(serde_json::Value::as_u64)
.expect("source sequence")
})
.collect();
let first = sequences.first().copied().expect("frames carry sequences");
assert_eq!(
sequences,
(first..first + frames.len() as u64).collect::<Vec<_>>(),
"meerkat's own contiguous sequence for the run, none missing or repeated"
);
let run_id = frames[0]
.data
.get("run_id")
.filter(|run_id| !run_id.is_null())
.expect("run_started carries its typed run id")
.clone();
for frame in frames {
if let Some(frame_run) = frame.data.get("run_id") {
assert_eq!(
frame_run, &run_id,
"{} is attributed to the run it belongs to",
frame.event_type
);
}
}
for kind in ["turn_started", "interaction_complete"] {
let frame = frames
.iter()
.find(|frame| frame.event_type == kind)
.expect("frame present");
assert_eq!(frame.data.get("run_id"), Some(&run_id), "{kind} lineage");
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn restored_members_first_run_reaches_the_console_after_stream_not_found_backoff() {
let mob_id = "console-replay-restore-backoff";
let temp = tempfile::tempdir().expect("temp dir");
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let session_id = commit_restart_member_turn(&first_identity, "before restart").await;
first.shutdown().await;
let boot = boot_persistent_before_activation(mob_id, &mob_path, &state_root).await;
let ingress = boot.runtime.install_test_event_ingress().await;
let handle = boot.runtime.mob_handle();
let mut forwarder = DrivenForwarder::new();
let restored = restart_member_entry(&handle).await;
assert_eq!(restored.status, MobMemberStatus::Active);
let key = restart_member_key(&restored, mob_id);
let not_live = match handle
.subscribe_agent_events(&restored.agent_identity)
.await
{
Ok(_) => panic!("no actor serves the restored member's session yet"),
Err(error) => error,
};
assert!(
matches!(
¬_live,
meerkat_mob::MobError::MemberSessionNotLive { session_id: missing, .. }
if missing == &session_id
),
"{not_live:?}"
);
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"no actor to attach to yet"
);
let backoff = forwarder
.subscribe_failures
.get_mut(&key)
.expect("the failed subscription parked the member in backoff");
assert_eq!(backoff.consecutive_failures, 1);
backoff.next_attempt = tokio::time::Instant::now() + Duration::from_hours(1);
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
assert_eq!(
restart_member_key(&restart_member_entry(&handle).await, mob_id),
key,
"the Resume route kept the restored binding"
);
let revived = commit_restart_member_turn(&identity_runtime, "after restart").await;
assert_eq!(revived, session_id);
assert!(
forwarder.tracked.current.is_empty(),
"the run finished unobserved"
);
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"the pending backoff still parks the member"
);
forwarder
.subscribe_failures
.get_mut(&key)
.expect("parked member")
.next_attempt = tokio::time::Instant::now();
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder
.tracked
.current
.get(&key)
.is_some_and(|attachment| attachment.actor.is_some()),
"the late subscription attached to the revived actor"
);
assert!(forwarder.subscribe_failures.is_empty());
forwarder.forward_terminals(&runtime, &ingress, 1).await;
assert_one_attributed_run(
&session_timeline(&runtime, &session_id).await,
"after restart",
);
runtime.shutdown().await;
}
async fn inject_activation_outcomes(
runtime: &UnifiedRuntime,
max_consecutive_stalls: u32,
outcomes: Vec<meerkat_mob::MobError>,
) {
let mut slot = runtime.pending_mob_activation.lock().await;
let pending = slot
.as_mut()
.expect("restarting a Stopped persistent mob stages an activation");
pending.set_stall_policy_for_test(
crate::mob_activation_retry::ActivationStallPolicy::with_max_consecutive_stalls(
max_consecutive_stalls,
),
);
pending.inject_resume_outcomes_for_test(outcomes);
}
fn restart_member_activation_stall() -> meerkat_mob::MobError {
meerkat_mob::MobError::LifecycleOperationProgressStalled {
intent: "explicit_resume".to_string(),
member_id: Some(meerkat_mob::AgentIdentity::from(RESTART_MEMBER)),
stage: "resume_member",
}
}
async fn restart_before_activation(
mob_id: &str,
temp: &tempfile::TempDir,
) -> (
PersistentBootBeforeActivation,
meerkat_core::types::SessionId,
) {
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let session_id = commit_restart_member_turn(&first_identity, "before restart").await;
first.shutdown().await;
(
boot_persistent_before_activation(mob_id, &mob_path, &state_root).await,
session_id,
)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_stalled_activation_that_later_resumes_completes_bootstrap() {
let temp = tempfile::tempdir().expect("temp dir");
let (boot, session_id) =
restart_before_activation("activation-stall-recovers", &temp).await;
inject_activation_outcomes(
&boot.runtime,
3,
vec![
restart_member_activation_stall(),
restart_member_activation_stall(),
],
)
.await;
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
assert_eq!(
runtime.mob_handle().status().await.expect("mob status"),
meerkat_mob::MobState::Running,
"the re-joined resume lifted the mob"
);
let revived = commit_restart_member_turn(&identity_runtime, "after restart").await;
assert_eq!(revived, session_id, "the restore resumed the same session");
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn consecutive_activation_stalls_at_one_point_fail_bootstrap() {
let temp = tempfile::tempdir().expect("temp dir");
let (boot, _) = restart_before_activation("activation-stall-exhausts", &temp).await;
inject_activation_outcomes(
&boot.runtime,
2,
vec![
restart_member_activation_stall(),
restart_member_activation_stall(),
],
)
.await;
let PersistentBootBeforeActivation {
mut runtime,
context,
roster,
..
} = boot;
let error = runtime
.install_and_bootstrap_identity_first_context(context, &roster)
.await
.expect_err("K consecutive stalls exhaust the activation");
assert!(
matches!(
&error,
crate::identity_first::IdentityRuntimeError::Internal(message)
if message.starts_with("activating the prepared mob")
),
"{error:?}"
);
assert_ne!(
runtime.mob_handle().status_observation_snapshot(),
meerkat_mob::MobState::Running,
"the failed activation never lifted the mob"
);
assert!(
runtime.pending_mob_activation.lock().await.is_none(),
"the obligation was consumed, not re-staged"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_non_stall_activation_error_fails_bootstrap_immediately() {
let temp = tempfile::tempdir().expect("temp dir");
let (boot, _) = restart_before_activation("activation-error-fatal", &temp).await;
inject_activation_outcomes(
&boot.runtime,
10,
vec![
meerkat_mob::MobError::Internal("injected activation failure".to_string()),
restart_member_activation_stall(),
],
)
.await;
let PersistentBootBeforeActivation {
mut runtime,
context,
roster,
..
} = boot;
let error = runtime
.install_and_bootstrap_identity_first_context(context, &roster)
.await
.expect_err("a non-stall activation error is fatal");
assert!(
matches!(
&error,
crate::identity_first::IdentityRuntimeError::Internal(message)
if message.contains("injected activation failure")
),
"{error:?}"
);
}
async fn discard_restart_member_actor(
runtime: &UnifiedRuntime,
session_id: &meerkat_core::types::SessionId,
) {
let service = runtime
.mob_runtime()
.session_service()
.cloned()
.expect("persistent runtime has a session service");
meerkat_mob::MobSessionService::discard_live_session(service.as_ref(), session_id)
.await
.expect("discard the live actor");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn same_binding_successor_attaches_only_after_its_predecessor_drained() {
let mob_id = "console-replay-successor";
let temp = tempfile::tempdir().expect("temp dir");
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let session_id = commit_restart_member_turn(&first_identity, "before restart").await;
first.shutdown().await;
let boot = boot_persistent_before_activation(mob_id, &mob_path, &state_root).await;
let ingress = boot.runtime.install_test_event_ingress().await;
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
let handle = runtime.mob_handle();
let key = restart_member_key(&restart_member_entry(&handle).await, mob_id);
let mut forwarder = DrivenForwarder::new();
forwarder.reconcile(&handle, &None).await;
let predecessor = forwarder
.tracked
.current
.get(&key)
.map(|attachment| attachment.generation)
.expect("the revived actor's stream attached");
commit_restart_member_turn(&identity_runtime, "predecessor run").await;
discard_restart_member_actor(&runtime, &session_id).await;
commit_restart_member_turn(&identity_runtime, "successor run").await;
assert_eq!(
restart_member_key(&restart_member_entry(&handle).await, mob_id),
key,
"the successor kept the binding atoms"
);
forwarder.reconcile(&handle, &None).await;
assert_eq!(
forwarder
.tracked
.current
.get(&key)
.map(|attachment| attachment.generation),
Some(predecessor),
"the predecessor's stream still serves the binding; the successor waits"
);
assert_eq!(forwarder.streams.len(), 1);
let steps = forwarder.forward_terminals(&runtime, &ingress, 2).await;
let close = steps
.iter()
.position(|step| matches!(step, ForwarderStep::Closed { .. }))
.expect("the predecessor's stream closed");
assert!(matches!(
steps[close],
ForwarderStep::Closed { generation, served: true } if generation == predecessor
));
let events = |steps: &[ForwarderStep]| -> Vec<(u64, &'static str)> {
steps
.iter()
.filter_map(|step| match step {
ForwarderStep::Event(envelope) => Some((
envelope.seq,
meerkat_core::event::agent_event_type(&envelope.payload),
)),
ForwarderStep::Closed { .. } | ForwarderStep::Abandoned { .. } => None,
})
.collect()
};
let (before, after) = (events(&steps[..close]), events(&steps[close + 1..]));
assert_eq!(before.first(), Some(&(1, "run_started")), "{before:?}");
assert_eq!(
before.last().map(|event| event.1),
Some("run_completed"),
"{before:?}"
);
let predecessor_tail = before.last().map_or(0, |event| event.0);
assert_eq!(
after.first().map(|event| event.1),
Some("run_started"),
"{after:?}"
);
assert!(
after
.first()
.is_some_and(|event| event.0 > predecessor_tail),
"the successor's sequence continues the predecessor's: {after:?}"
);
assert_eq!(
after.last().map(|event| event.1),
Some("run_completed"),
"{after:?}"
);
assert!(
forwarder
.tracked
.current
.get(&key)
.is_some_and(
|attachment| attachment.generation != predecessor && attachment.actor.is_some()
),
"the successor's stream serves the binding now"
);
let timeline = session_timeline(&runtime, &session_id).await;
let second_start = timeline
.iter()
.rposition(|frame| frame.event_type == "run_started")
.expect("successor run_started");
assert_one_attributed_run(&timeline[..second_start], "predecessor run");
assert_one_attributed_run(&timeline[second_start..], "successor run");
assert_ne!(
timeline[0].data.get("run_id"),
timeline[second_start].data.get("run_id")
);
runtime.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn discarded_actor_is_never_attributed_and_its_successor_attaches_from_its_first_event() {
let mob_id = "console-replay-discarded";
let temp = tempfile::tempdir().expect("temp dir");
let mob_path = temp.path().join("mob.sqlite3");
let state_root = temp.path().join("state");
let (first, first_identity) =
boot_identity_first_persistent(mob_id, &mob_path, &state_root).await;
let session_id = commit_restart_member_turn(&first_identity, "before restart").await;
first.shutdown().await;
let boot = boot_persistent_before_activation(mob_id, &mob_path, &state_root).await;
let ingress = boot.runtime.install_test_event_ingress().await;
let identity_runtime = boot.identity_runtime.clone();
let runtime = boot.activate().await;
let handle = runtime.mob_handle();
let key = restart_member_key(&restart_member_entry(&handle).await, mob_id);
commit_restart_member_turn(&identity_runtime, "revoked run").await;
discard_restart_member_actor(&runtime, &session_id).await;
let mut forwarder = DrivenForwarder::new();
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"nothing is attached from a revoked actor, and no actor serves the session"
);
assert!(forwarder.subscribe_failures.contains_key(&key));
commit_restart_member_turn(&identity_runtime, "successor run").await;
forwarder
.subscribe_failures
.get_mut(&key)
.expect("parked member")
.next_attempt = tokio::time::Instant::now();
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder
.tracked
.current
.get(&key)
.is_some_and(|attachment| attachment.actor.is_some()),
"the successor's stream attached"
);
forwarder.forward_terminals(&runtime, &ingress, 1).await;
assert_one_attributed_run(
&session_timeline(&runtime, &session_id).await,
"successor run",
);
runtime.shutdown().await;
}
#[tokio::test]
async fn stale_close_never_untracks_the_stream_serving_its_key() {
let key = TrackedAgentEventStream {
mob_id: "stale-close".to_string(),
durable_identity: None,
member_identity: AgentIdentity::from("worker"),
runtime_id: AgentRuntimeId::new(AgentIdentity::from("worker"), Generation::new(0)),
identity_fencing_token: None,
fence_token: FenceToken::new(1),
};
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
let superseded = attach_member_event_stream(
&mut streams,
key.clone(),
ProfileName::from("worker"),
Box::pin(futures::stream::empty()),
None,
None,
None,
);
let superseded_generation = superseded.generation;
let serving = attach_member_event_stream(
&mut streams,
key.clone(),
ProfileName::from("worker"),
Box::pin(futures::stream::pending()),
None,
None,
None,
);
let serving_generation = serving.generation;
let mut tracked = AttachedStreams::default();
tracked.current.insert(key.clone(), serving);
let Some(ForwardedAgentEvent::Closed(closed_key, generation)) = streams.next().await else {
panic!("the superseded stream closes");
};
assert_eq!(
(closed_key.clone(), generation),
(key.clone(), superseded_generation)
);
assert!(!tracked.close(&closed_key, generation));
assert_eq!(
tracked
.current
.get(&key)
.map(|attachment| attachment.generation),
Some(serving_generation)
);
assert!(tracked.close(&key, serving_generation));
assert!(tracked.current.is_empty());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn child_mob_members_first_run_is_replayed_under_the_child_mob() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let (mob_runtime, _) = mob_with_idle_worker("console-replay-parent", &temp_dir).await;
let state = mob_runtime
.agent_mob_mcp_state()
.expect("stock constructor installs agent mob tools");
let child_definition = meerkat_mob::MobDefinition::from_toml(
"[mob]\nid = \"console-replay-child\"\n\n[profiles.worker]\nmodel = \"gpt-5.5\"\n\
external_addressable = true\n\n[profiles.worker.tools]\ncomms = true\n",
)
.expect("child mob definition");
state
.mob_create_definition(child_definition)
.await
.expect("create child mob");
let child = Box::pin(state.mob_handles_snapshot())
.await
.expect("managed mobs")
.into_iter()
.find(|(mob_id, _)| mob_id.as_str() == "console-replay-child")
.map(|(_, handle)| handle)
.expect("child mob handle");
let mut member = SpawnMemberSpec::new(
ProfileName::from("worker"),
AgentIdentity::from(IDLE_WORKER),
);
member.runtime_mode = Some(meerkat_mob::MobRuntimeMode::TurnDriven);
child
.ensure_member(member)
.await
.expect("seat child worker");
run_worker_turn(&child, "child probe").await;
let mut forwarder = DrivenForwarder::new();
forwarder
.reconcile(&mob_runtime.handle(), &Some(state))
.await;
let ForwardedAgentEvent::Event(event) = forwarder.next().await else {
panic!("the child stream yields its first run");
};
let (runtime_id, _, _, envelope, _, _) = *event;
assert!(matches!(envelope.payload, AgentEvent::RunStarted { .. }));
assert_eq!(envelope.seq, 1);
assert!(
forwarder
.tracked
.current
.keys()
.any(|key| key.mob_id == "console-replay-child"
&& key.runtime_id == runtime_id
&& forwarder.tracked.current[key].actor.is_some()),
"attributed to the child mob's member, with its actor"
);
drop(forwarder);
let _ = mob_runtime.handle().shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn stuck_predecessor_is_cut_off_at_its_drain_deadline() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let (mob_runtime, _) =
mob_with_idle_worker("console-replay-stuck-predecessor", &temp_dir).await;
let handle = mob_runtime.handle();
let module_runtime = std::thread::spawn(|| {
start_mobkit_runtime_with_options(
MobKitConfig {
modules: vec![],
discovery: crate::types::DiscoverySpec {
namespace: "console-replay-stuck-predecessor".to_string(),
modules: vec![],
},
pre_spawn: vec![],
},
Vec::new(),
Duration::from_secs(2),
RuntimeOptions::default(),
)
})
.join()
.expect("module runtime thread")
.expect("module runtime");
let runtime = UnifiedRuntime::from_parts(
mob_runtime,
module_runtime,
Arc::new(InMemoryMetadataStore::new()),
)
.await;
let ingress = runtime.install_test_event_ingress().await;
handle
.respawn(AgentIdentity::from(IDLE_WORKER), None)
.await
.expect("respawn onto a fresh actor");
let predecessor_session = handle
.resolve_bridge_session_id(&AgentIdentity::from(IDLE_WORKER))
.await
.expect("predecessor session");
let entry = handle
.list_members_including_retiring()
.await
.into_iter()
.find(|entry| entry.agent_identity == IDLE_WORKER)
.expect("worker roster entry");
let predecessor_key = restart_member_key(&entry, handle.mob_id().as_str());
let subscription = handle
.subscribe_agent_events_from(
&AgentIdentity::from(IDLE_WORKER),
meerkat_core::comms::SessionEventCursor::Earliest,
)
.await
.expect("subscribe to the fresh actor");
assert!(subscription.actor.is_some(), "the fresh actor is named");
let mut forwarder = DrivenForwarder::new();
let attachment = attach_member_event_stream(
&mut forwarder.streams,
predecessor_key.clone(),
entry.role.clone(),
Box::pin(subscription.stream.chain(futures::stream::pending())),
subscription.actor,
subscription.epoch,
None,
);
let predecessor = attachment.generation;
forwarder
.tracked
.current
.insert(predecessor_key.clone(), attachment);
run_worker_turn(&handle, "predecessor run").await;
forwarder.forward_terminals(&runtime, &ingress, 1).await;
handle
.respawn(AgentIdentity::from(IDLE_WORKER), None)
.await
.expect("respawn the member");
let successor_session = handle
.resolve_bridge_session_id(&AgentIdentity::from(IDLE_WORKER))
.await
.expect("successor session");
assert_ne!(successor_session, predecessor_session);
run_worker_turn(&handle, "successor run").await;
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"the successor's new binding waits for the predecessor"
);
let held_since = forwarder
.tracked
.departed
.get(&predecessor)
.filter(|departed| !departed.actor.is_live())
.and_then(|departed| departed.attachment.held_since)
.expect("the revoked predecessor is draining, and its hold is recorded");
assert_eq!(
earliest_reconcile_deadline(&forwarder.subscribe_failures, &forwarder.tracked),
Some(held_since + PREDECESSOR_DRAIN_DEADLINE),
"the forwarder wakes itself at the drain deadline"
);
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"the successor still waits before the deadline"
);
{
use futures::FutureExt;
assert!(
forwarder.streams.next().now_or_never().is_none(),
"the predecessor's stream is still open"
);
}
forwarder
.tracked
.departed
.get_mut(&predecessor)
.expect("draining predecessor")
.attachment
.held_since = Some(held_since - PREDECESSOR_DRAIN_DEADLINE);
forwarder.reconcile(&handle, &None).await;
assert!(
forwarder.tracked.current.is_empty(),
"the successor attaches only after the gap marker went out"
);
let steps = forwarder.forward_terminals(&runtime, &ingress, 1).await;
assert!(
matches!(
steps.as_slice(),
[
ForwarderStep::Abandoned { generation: abandoned },
ForwarderStep::Closed { generation: closed, served: false },
..
] if *abandoned == predecessor && *closed == predecessor
),
"the predecessor is cut off first: {steps:?}"
);
let successor_events: Vec<(u64, &'static str)> = steps[2..]
.iter()
.filter_map(|step| match step {
ForwarderStep::Event(envelope) => Some((
envelope.seq,
meerkat_core::event::agent_event_type(&envelope.payload),
)),
_ => None,
})
.collect();
assert_eq!(successor_events.first(), Some(&(1, "run_started")));
assert_eq!(
successor_events.last().map(|event| event.1),
Some("run_completed")
);
let successor_key = forwarder
.tracked
.current
.keys()
.next()
.expect("the successor's stream serves its binding")
.clone();
assert_ne!(successor_key, predecessor_key);
let replay = runtime
.console_events()
.replay_all(None)
.await
.expect("console replay");
let gap = replay
.iter()
.position(|frame| frame.event_type == "stream_truncated")
.expect("the gap is on the console timeline");
assert_eq!(
replay[gap].data["reason"]["kind"],
json!("predecessor_stream_abandoned")
);
assert_eq!(replay[gap].data["session_id"], json!(predecessor_session));
let session_frames = |session: &meerkat_core::types::SessionId| {
replay
.iter()
.enumerate()
.filter(|(_, frame)| frame.data.get("session_id") == Some(&json!(session)))
.filter(|(_, frame)| frame.event_type != "stream_truncated")
.map(|(index, frame)| (index, frame.clone()))
.collect::<Vec<_>>()
};
let predecessor_frames = session_frames(&predecessor_session);
let successor_frames = session_frames(&successor_session);
assert!(predecessor_frames.iter().all(|(index, _)| *index < gap));
assert!(
successor_frames.iter().all(|(index, _)| *index > gap),
"every successor frame follows the gap"
);
let frames = |indexed: Vec<(
usize,
crate::console_contracts::ConsoleIdentityEventEnvelope,
)>| {
indexed
.into_iter()
.map(|(_, frame)| frame)
.collect::<Vec<_>>()
};
assert_one_attributed_run(&frames(predecessor_frames), "predecessor run");
assert_one_attributed_run(&frames(successor_frames), "successor run");
runtime.shutdown().await;
}
#[tokio::test]
async fn departed_live_actor_on_another_session_drains_until_its_deadline() {
let dir = tempfile::tempdir().expect("temp dir");
let raw = Arc::new(meerkat::build_ephemeral_service(
meerkat::AgentFactory::new(dir.path()),
meerkat::Config::default(),
16,
)) as Arc<dyn meerkat_mob::MobSessionService>;
let request = meerkat_core::service::CreateSessionRequest {
model: "gpt-5.5".to_string(),
prompt: meerkat_core::ContentInput::Text("probe".to_string()),
system_prompt: meerkat_core::config::SystemPromptOverride::Inherit,
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::Discard,
build: Some(meerkat_core::service::SessionBuildOptions {
llm_client_override: Some(meerkat::encode_llm_client_override_for_service(
Arc::new(meerkat_client::TestClient::default()),
)),
..Default::default()
}),
labels: None,
injected_context: Vec::new(),
};
let old_session = raw
.create_session(request)
.await
.expect("deferred create")
.session_id;
let actor = raw
.subscribe_agent_session_events_from(
&old_session,
meerkat_core::comms::SessionEventCursor::Live,
)
.await
.expect("subscribe")
.actor
.expect("the ephemeral service names its actor");
let new_session = meerkat_core::types::SessionId::new();
let key = TrackedAgentEventStream {
mob_id: "departed-other-session".to_string(),
durable_identity: None,
member_identity: AgentIdentity::from("worker"),
runtime_id: AgentRuntimeId::new(AgentIdentity::from("worker"), Generation::new(0)),
identity_fencing_token: None,
fence_token: FenceToken::new(1),
};
let owner = MemberStreamOwner::of(&key);
let mut streams: SelectAll<TaggedAgentEventStream> = SelectAll::new();
let attachment = attach_member_event_stream(
&mut streams,
key.clone(),
ProfileName::from("worker"),
Box::pin(futures::stream::pending()),
Some(actor),
None,
None,
);
let generation = attachment.generation;
let mut tracked = AttachedStreams::default();
tracked.current.insert(key.clone(), attachment);
tracked.depart(&key, false);
let now = tokio::time::Instant::now();
assert!(!tracked.predecessor_draining(&owner, Some(&old_session), now));
assert_eq!(
tracked.live_departed_for(&owner, &old_session),
Some(generation)
);
assert!(tracked.predecessor_draining(&owner, Some(&new_session), now));
assert_eq!(
tracked.earliest_drain_deadline(),
Some(now + PREDECESSOR_DRAIN_DEADLINE)
);
assert!(tracked.predecessor_draining(
&owner,
Some(&new_session),
now + PREDECESSOR_DRAIN_DEADLINE
));
assert_eq!(tracked.earliest_drain_deadline(), None, "cut off");
let Some(ForwardedAgentEvent::Abandoned {
session_id,
generation: abandoned,
..
}) = streams.next().await
else {
panic!("the cut-off stream reports its gap first");
};
assert_eq!((session_id, abandoned), (old_session, generation));
let Some(ForwardedAgentEvent::Closed(closed_key, closed)) = streams.next().await else {
panic!("then closes");
};
assert!(!tracked.close(&closed_key, closed));
assert!(!tracked.predecessor_draining(&owner, Some(&new_session), now));
}
}