use std::collections::{BTreeMap, HashMap};
use std::future::Future;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use meerkat_core::lifecycle::run_primitive::{
OpenAiProviderTag, ProviderParamsOverride, ProviderTag,
};
use meerkat_core::types::HandlingMode;
use meerkat_mob::ids::AgentIdentity as MobAgentIdentity;
use meerkat_mob::launch::MemberLaunchMode;
use meerkat_mob::{
MobHandle, MobSessionService, ResumeOverrideField, SpawnMemberSpec, SpawnSystemPromptOverride,
WorkOrigin, WorkRef, WorkSpec,
};
use crate::mob_handle_runtime::{
content_input_has_images, is_previous_member_cleanup_ambiguous_error,
is_recoverable_lifecycle_cleanup_error, is_recoverable_session_owned_retire_cleanup_error,
model_capabilities_for_member, topology_restore_failed_peer_ids,
};
use super::adapters::{ContinuitySessionStoreAdapter, SessionRuntimeState};
use super::types::{
AgentBuildDraft, AgentIdentity, AgentRuntimeId, CheckpointVersion, ContinuityGeneration,
DurableAgentSpec, FencingToken, SessionSnapshot,
};
fn is_missing_event_injector_error(error: &str) -> bool {
error.contains("missing event injector capability")
|| (error.contains("missing required capability")
&& error.contains("interaction_event_injector"))
}
fn is_missing_bridge_session_snapshot_error(error: &str) -> bool {
error.contains("missing bridge session snapshot")
}
fn durable_snapshot_is_typed_absent(error: &meerkat_mob::MobError) -> bool {
matches!(
error,
meerkat_mob::MobError::SessionUnavailableForResume {
reason: meerkat_mob::error::SessionResumeUnavailableReason::Absent,
..
}
)
}
fn is_repairable_bridge_delivery_error(error: &str) -> bool {
is_missing_event_injector_error(error)
|| is_missing_bridge_session_snapshot_error(error)
|| is_previous_member_cleanup_ambiguous_error(error)
}
fn is_recoverable_bridge_respawn_cleanup_error(error: &str) -> bool {
is_recoverable_lifecycle_cleanup_error(error)
}
fn is_member_already_exists_error(error: &meerkat_mob::MobError) -> bool {
matches!(error, meerkat_mob::MobError::MemberAlreadyExists(_))
}
async fn abandon_then_retire_reset_superseded<Abandon, AbandonFuture, Retire, RetireFuture>(
member_id: &MobAgentIdentity,
session_id: &meerkat_core::types::SessionId,
abandon: Abandon,
retire: Retire,
) -> Result<(), BridgeError>
where
Abandon: FnOnce() -> AbandonFuture,
AbandonFuture: Future<Output = Result<(), meerkat_store::SessionStoreError>>,
Retire: FnOnce() -> RetireFuture,
RetireFuture: Future<Output = Result<(), meerkat_mob::MobError>>,
{
abandon().await.map_err(|error| {
BridgeError::Mob(format!(
"reset retire could not abandon superseded session {session_id}: {error}"
))
})?;
match retire().await {
Ok(()) | Err(meerkat_mob::MobError::MemberNotFound(_)) => Ok(()),
Err(error) => Err(BridgeError::Mob(format!(
"reset retire cleanup failed for {member_id} after superseded session \
{session_id} was abandoned: {error}"
))),
}
}
#[derive(Debug, PartialEq, Eq)]
enum MemberRepairRespawnFailure {
DegradedTopologyRestore { failed_peer_ids: Vec<String> },
RecoverableCleanup,
Fatal(String),
}
#[derive(Debug)]
enum RepairResumeFailure {
DurableSnapshotMissing { detail: String },
Rejected(BridgeError),
}
fn classify_member_repair_respawn_failure(
error: &meerkat_mob::MobRespawnError,
) -> MemberRepairRespawnFailure {
if let Some(failed_peer_ids) = topology_restore_failed_peer_ids(error) {
return MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids };
}
if is_recoverable_bridge_respawn_cleanup_error(&error.to_string()) {
return MemberRepairRespawnFailure::RecoverableCleanup;
}
MemberRepairRespawnFailure::Fatal(error.to_string())
}
fn fresh_member_spec_from_pre_delivery_entry(
member_id: &MobAgentIdentity,
role: meerkat_mob::ProfileName,
labels: BTreeMap<String, String>,
) -> SpawnMemberSpec {
let mut spec = SpawnMemberSpec::new(role, member_id.clone());
if !labels.is_empty() {
spec = spec.with_labels(labels);
}
spec
}
#[derive(Debug)]
pub enum BridgeError {
Mob(String),
InvalidInput(String),
ResumeRejected {
kind: ResumeRejectionKind,
detail: String,
},
ActorAdmissionTimeout {
operation: &'static str,
identity: MobAgentIdentity,
waited: Duration,
},
}
impl std::fmt::Display for BridgeError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Mob(msg) => write!(f, "session bridge mob error: {msg}"),
Self::InvalidInput(msg) => write!(f, "session bridge invalid input: {msg}"),
Self::ResumeRejected { kind, detail } => write!(
f,
"session bridge resume rejected ({kind:?}): {detail}; durable session preserved, \
identity degraded pending retry"
),
Self::ActorAdmissionTimeout {
operation,
identity,
waited,
} => write!(
f,
"session bridge actor call `{operation}` for member {identity} exceeded the \
admission budget after {waited:?}; the mob actor command loop is not draining"
),
}
}
}
impl std::error::Error for BridgeError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ResumeRejectionKind {
MemberRestoreFailed,
TranscriptContinuity,
ArchivedNotRevivable,
Other,
}
pub(crate) fn archived_not_revivable_park_reason(
session_id: &meerkat_core::types::SessionId,
detail: &str,
) -> String {
format!(
"durable session {session_id} is archived and its runtime lifecycle refuses \
revival (typed ArchivedNotRevivable): a stable verdict retries cannot change, \
so continuity repair parks instead of heal-looping. The transcript is intact \
and preserved. Operator path: upstream archived-session revive lands in \
meerkat 0.8.15; until then reset the identity via `mobkit/reset` (deliberate \
fresh start) or restart the gateway after an upstream fix (the park is \
process-local). Refusal: {detail}"
)
}
fn resume_rejected(
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
error: &meerkat_mob::MobError,
step: &str,
) -> BridgeError {
let kind = classify_resume_error(error);
tracing::error!(
identity = %identity,
session_id = %session_id,
kind = ?kind,
step,
error = %error,
"resume rejected; durable session preserved, identity degraded pending reconcile retry \
(refusing fresh-spawn fallback)"
);
BridgeError::ResumeRejected {
kind,
detail: format!("{step}: {error}"),
}
}
fn classify_resume_error(error: &meerkat_mob::MobError) -> ResumeRejectionKind {
if matches!(
error,
meerkat_mob::MobError::SessionUnavailableForResume {
reason: meerkat_mob::error::SessionResumeUnavailableReason::ArchivedNotRevivable,
..
}
) {
return ResumeRejectionKind::ArchivedNotRevivable;
}
if matches!(error, meerkat_mob::MobError::MemberRestoreFailed { .. }) {
return ResumeRejectionKind::MemberRestoreFailed;
}
let text = error.to_string();
if text.contains("not a continuation of persisted revision")
|| text.contains("TranscriptContinuityViolation")
|| text.contains("continuity preflight")
{
return ResumeRejectionKind::TranscriptContinuity;
}
ResumeRejectionKind::Other
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResumeFallbackReason {
RuntimeIdentityIncompatible { detail: String },
NeverPersisted { detail: String },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResumeSessionOutcome {
Resumed {
session_id: meerkat_core::types::SessionId,
},
FreshSpawned {
session_id: meerkat_core::types::SessionId,
reason: ResumeFallbackReason,
},
}
impl ResumeSessionOutcome {
#[must_use]
pub fn session_id(&self) -> &meerkat_core::types::SessionId {
match self {
Self::Resumed { session_id } | Self::FreshSpawned { session_id, .. } => session_id,
}
}
#[must_use]
pub fn fallback_reason(&self) -> Option<&ResumeFallbackReason> {
match self {
Self::Resumed { .. } => None,
Self::FreshSpawned { reason, .. } => Some(reason),
}
}
}
const BRIDGE_ACTOR_ADMISSION_BUDGET: Duration = Duration::from_mins(10);
fn bridge_actor_admission_budget() -> Duration {
parse_bridge_actor_admission_budget(
std::env::var("MOBKIT_BRIDGE_ACTOR_ADMISSION_SECS")
.ok()
.as_deref(),
)
}
fn parse_bridge_actor_admission_budget(raw: Option<&str>) -> Duration {
raw.and_then(|value| value.trim().parse::<u64>().ok())
.map(|secs| Duration::from_secs(secs.clamp(1, 3600)))
.unwrap_or(BRIDGE_ACTOR_ADMISSION_BUDGET)
}
struct ActorAdmissionDeadline {
started: tokio::time::Instant,
deadline: tokio::time::Instant,
}
impl ActorAdmissionDeadline {
fn new(budget: Duration) -> Self {
let started = tokio::time::Instant::now();
Self {
started,
deadline: started + budget,
}
}
async fn bound<T, F>(
&self,
operation: &'static str,
identity: &MobAgentIdentity,
call: F,
) -> Result<T, BridgeError>
where
F: Future<Output = T>,
{
match tokio::time::timeout_at(self.deadline, call).await {
Ok(value) => Ok(value),
Err(_) => {
let waited = self.started.elapsed();
tracing::warn!(
operation,
identity = %identity,
waited_ms = waited.as_millis(),
"mob actor did not answer within the delivery admission budget; the actor \
command loop is not draining (head-of-line block) — delivery abandoned"
);
Err(BridgeError::ActorAdmissionTimeout {
operation,
identity: identity.clone(),
waited,
})
}
}
}
}
fn internal_bridge_work_spec(
content: &meerkat_core::ContentInput,
system_prompt: Option<&str>,
injected_context: &[meerkat_core::ContentInput],
interaction_id: Option<&str>,
) -> WorkSpec {
let mut spec = WorkSpec::new(content.clone(), WorkOrigin::Internal);
if let Some(system_prompt) = system_prompt {
spec = spec.with_system_prompt(system_prompt);
}
if !injected_context.is_empty() {
spec = spec.with_injected_context(injected_context.to_vec());
}
if let Some(raw) = interaction_id {
match raw.parse::<uuid::Uuid>() {
Ok(id) => {
spec = spec.with_interaction_id(meerkat_core::interaction::InteractionId(id));
}
Err(_) => {
tracing::debug!(
interaction_id = %raw,
"non-UUID interaction id not threaded into runtime admission"
);
}
}
}
spec
}
struct InternalBridgeWork<'a> {
content: &'a meerkat_core::ContentInput,
system_prompt: Option<&'a str>,
injected_context: &'a [meerkat_core::ContentInput],
interaction_id: Option<&'a str>,
delivery_identity: Option<&'a meerkat_mob::MobDeliveryIdentity>,
}
async fn submit_internal_bridge_work(
handle: &MobHandle,
member_id: &MobAgentIdentity,
work: InternalBridgeWork<'_>,
handling_mode: HandlingMode,
deadline: &ActorAdmissionDeadline,
) -> Result<(), BridgeError> {
let entry = deadline
.bound(
"deliver.get_member",
member_id,
handle.get_member(member_id),
)
.await?
.map_err(|err| BridgeError::Mob(err.to_string()))?
.ok_or_else(|| BridgeError::Mob(format!("member not found: {member_id}")))?;
let spec = internal_bridge_work_spec(
work.content,
work.system_prompt,
work.injected_context,
work.interaction_id,
);
match work.delivery_identity {
Some(delivery_identity) => deadline
.bound(
"deliver.submit_work",
member_id,
handle.submit_work_with_mode_and_delivery_identity(
entry.agent_runtime_id.clone(),
entry.fence_token,
spec,
handling_mode,
delivery_identity.clone(),
),
)
.await?
.map(|_| ())
.map_err(|err| BridgeError::Mob(err.to_string())),
None => deadline
.bound(
"deliver.submit_work",
member_id,
handle.submit_work_with_mode(
entry.agent_runtime_id.clone(),
entry.fence_token,
WorkRef::new(),
spec,
handling_mode,
),
)
.await?
.map(|_| ())
.map_err(|err| BridgeError::Mob(err.to_string())),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CommittedBoundaryRepair {
AlreadyCommitted,
Recovered,
Unprovable { reason: String },
Unsupported,
}
#[async_trait]
pub trait CommittedBoundaryRecoverer: Send + Sync {
async fn recover_committed_boundary(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<CommittedBoundaryRepair, BridgeError>;
}
#[async_trait]
impl<B> CommittedBoundaryRecoverer for meerkat_session::PersistentSessionService<B>
where
B: meerkat_session::SessionAgentBuilder + 'static,
{
async fn recover_committed_boundary(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<CommittedBoundaryRepair, BridgeError> {
match meerkat_session::PersistentSessionService::recover_committed_boundary(
self, session_id,
)
.await
{
Ok(meerkat_session::CommittedBoundaryRecovery::AlreadyCommitted) => {
Ok(CommittedBoundaryRepair::AlreadyCommitted)
}
Ok(meerkat_session::CommittedBoundaryRecovery::Recovered { message_count }) => {
tracing::info!(
%session_id,
message_count,
"machine-authorized recovery persisted a committed durable head"
);
Ok(CommittedBoundaryRepair::Recovered)
}
Ok(meerkat_session::CommittedBoundaryRecovery::Unprovable { reason }) => {
Ok(CommittedBoundaryRepair::Unprovable { reason })
}
Err(error) => map_committed_boundary_recovery_error(error),
}
}
}
fn map_committed_boundary_recovery_error(
error: meerkat_core::SessionError,
) -> Result<CommittedBoundaryRepair, BridgeError> {
match error {
error @ (meerkat_core::SessionError::DurableTailRecoveryRefused { .. }
| meerkat_core::SessionError::DurableEvidenceQuarantined { .. }) => {
Ok(CommittedBoundaryRepair::Unprovable {
reason: error.to_string(),
})
}
error => Err(BridgeError::Mob(format!(
"committed-boundary recovery: {error}"
))),
}
}
#[derive(Debug, Clone)]
pub struct BridgeDelivery {
pub content: meerkat_core::ContentInput,
pub handling_mode: HandlingMode,
pub system_prompt: Option<String>,
pub injected_context: Vec<meerkat_core::ContentInput>,
pub interaction_id: Option<String>,
pub delivery_identity: Option<meerkat_mob::MobDeliveryIdentity>,
}
impl BridgeDelivery {
pub fn new(content: meerkat_core::ContentInput, handling_mode: HandlingMode) -> Self {
Self {
content,
handling_mode,
system_prompt: None,
injected_context: Vec::new(),
interaction_id: None,
delivery_identity: None,
}
}
}
#[async_trait]
pub trait SessionBridge: Send + Sync {
async fn raw_member_alias_exists(&self, _alias: &str) -> Result<bool, BridgeError> {
Ok(false)
}
async fn create_session(
&self,
identity: &AgentIdentity,
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
) -> Result<meerkat_core::types::SessionId, BridgeError>;
fn requires_resume_snapshot(&self) -> bool {
true
}
async fn resume_session(
&self,
identity: &AgentIdentity,
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
snapshot: &SessionSnapshot,
) -> Result<ResumeSessionOutcome, BridgeError>;
async fn deliver_admitted(
&self,
runtime_id: &AgentRuntimeId,
delivery: BridgeDelivery,
) -> Result<meerkat_core::types::SessionId, BridgeError>;
async fn deliver(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
self.deliver_admitted(
runtime_id,
BridgeDelivery::new(content.clone(), HandlingMode::Queue),
)
.await
}
async fn deliver_with_mode(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
handling_mode: HandlingMode,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
self.deliver_admitted(
runtime_id,
BridgeDelivery::new(content.clone(), handling_mode),
)
.await
}
async fn deliver_with_mode_and_context(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
injected_context: &[meerkat_core::ContentInput],
handling_mode: HandlingMode,
interaction_id: Option<&str>,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let mut delivery = BridgeDelivery::new(content.clone(), handling_mode);
delivery.injected_context = injected_context.to_vec();
delivery.interaction_id = interaction_id.map(ToString::to_string);
self.deliver_admitted(runtime_id, delivery).await
}
async fn deliver_with_mode_context_and_system_prompt(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
system_prompt: Option<&str>,
injected_context: &[meerkat_core::ContentInput],
handling_mode: HandlingMode,
interaction_id: Option<&str>,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let mut delivery = BridgeDelivery::new(content.clone(), handling_mode);
delivery.system_prompt = system_prompt.map(ToString::to_string);
delivery.injected_context = injected_context.to_vec();
delivery.interaction_id = interaction_id.map(ToString::to_string);
self.deliver_admitted(runtime_id, delivery).await
}
async fn checkpoint_session(
&self,
runtime_id: &AgentRuntimeId,
session_id: &meerkat_core::types::SessionId,
) -> Result<SessionSnapshot, BridgeError>;
async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError>;
async fn retire_reset_superseded_member(
&self,
runtime_id: &AgentRuntimeId,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
self.retire_member(runtime_id).await?;
self.unregister_session_runtime_state(session_id).await
}
async fn wire_peer(&self, _a: &AgentRuntimeId, _b: &AgentRuntimeId) -> Result<(), BridgeError> {
Err(BridgeError::Mob("peer wiring not supported".to_string()))
}
async fn wire_peers_batch(
&self,
edges: &[(AgentRuntimeId, AgentRuntimeId)],
) -> Result<(), BridgeError> {
for (a, b) in edges {
self.wire_peer(a, b).await?;
}
Ok(())
}
async fn current_member_wires(
&self,
) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
Ok(Vec::new())
}
async fn current_member_wires_any_half(
&self,
) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
self.current_member_wires().await
}
async fn unwire_peer(
&self,
_a: &AgentRuntimeId,
_b: &AgentRuntimeId,
) -> Result<(), BridgeError> {
Err(BridgeError::Mob("peer unwiring not supported".to_string()))
}
async fn inspect_member(
&self,
_runtime_id: &AgentRuntimeId,
) -> Result<MemberInspection, BridgeError> {
Err(BridgeError::Mob("inspect not supported".to_string()))
}
async fn register_session_runtime_state(
&self,
_session_id: &meerkat_core::types::SessionId,
_identity: &AgentIdentity,
_generation: ContinuityGeneration,
checkpoint_version: CheckpointVersion,
_fencing_token: FencingToken,
) -> Result<CheckpointVersion, BridgeError> {
Ok(checkpoint_version)
}
async fn suspend_session_runtime_state(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
self.unregister_session_runtime_state(session_id).await
}
async fn unregister_session_runtime_state(
&self,
_session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
Ok(())
}
async fn recover_committed_boundary(
&self,
_session_id: &meerkat_core::types::SessionId,
) -> Result<CommittedBoundaryRepair, BridgeError> {
Ok(CommittedBoundaryRepair::Unsupported)
}
}
#[derive(Debug, Clone)]
pub struct MemberInspection {
pub output_preview: Option<String>,
pub is_final: bool,
pub peer_reachable_count: usize,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct UnmaskedResumeDivergence {
model: bool,
provider: bool,
}
fn unmasked_resume_divergence(
mask: &meerkat_core::service::ResumeOverrideMask,
declared_model: &str,
declared_provider: Option<meerkat_core::Provider>,
restored_model: &str,
restored_provider: meerkat_core::Provider,
) -> UnmaskedResumeDivergence {
UnmaskedResumeDivergence {
model: !mask.model && restored_model != declared_model,
provider: !mask.provider
&& declared_provider.is_some_and(|declared| restored_provider != declared),
}
}
pub struct MobSessionBridge {
handle: MobHandle,
session_store: Option<Arc<dyn meerkat::SessionStore>>,
session_service: Option<Arc<dyn MobSessionService>>,
continuity_session_store: Option<Arc<ContinuitySessionStoreAdapter>>,
runtime_members: Arc<tokio::sync::RwLock<HashMap<String, String>>>,
runtime_sessions: Arc<tokio::sync::RwLock<HashMap<String, meerkat_core::types::SessionId>>>,
generated_external_owner_session: std::sync::OnceLock<meerkat_core::types::SessionId>,
committed_boundary_recoverer: Option<Arc<dyn CommittedBoundaryRecoverer>>,
resume_divergence_logged: std::sync::Mutex<std::collections::HashSet<String>>,
actor_admission_budget: Duration,
runtime_ingress_authority: Option<Arc<dyn meerkat_runtime::SessionServiceRuntimeExt>>,
}
struct CarriedMemberInput {
original_input_id: meerkat_core::lifecycle::InputId,
admission_sequence: Option<u64>,
input: meerkat_runtime::Input,
}
struct PendingIngressCapture {
carryable: Vec<CarriedMemberInput>,
uncarryable: Vec<(meerkat_core::lifecycle::InputId, &'static str, String)>,
}
impl PendingIngressCapture {
fn empty() -> Self {
Self {
carryable: Vec::new(),
uncarryable: Vec::new(),
}
}
fn is_empty(&self) -> bool {
self.carryable.is_empty() && self.uncarryable.is_empty()
}
}
fn remint_carried_input_identity(
mut input: meerkat_runtime::Input,
) -> Option<meerkat_runtime::Input> {
let header = match &mut input {
meerkat_runtime::Input::Prompt(i) => &mut i.header,
meerkat_runtime::Input::Peer(i) => &mut i.header,
meerkat_runtime::Input::FlowStep(i) => &mut i.header,
meerkat_runtime::Input::ExternalEvent(i) => &mut i.header,
meerkat_runtime::Input::Continuation(i) => &mut i.header,
meerkat_runtime::Input::Operation(i) => &mut i.header,
_ => return None,
};
header.id = meerkat_core::lifecycle::InputId::new();
if let Some(key) = header.idempotency_key.take() {
tracing::debug!(
idempotency_key = %key,
readmitted_input_id = %header.id,
"carried input drops its idempotency key: the original admission was \
durably terminalized by the repair disposal"
);
}
Some(input)
}
impl MobSessionBridge {
pub fn new(handle: MobHandle) -> Self {
Self {
handle,
session_store: None,
session_service: None,
continuity_session_store: None,
runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
generated_external_owner_session: std::sync::OnceLock::new(),
committed_boundary_recoverer: None,
resume_divergence_logged: std::sync::Mutex::new(std::collections::HashSet::new()),
actor_admission_budget: bridge_actor_admission_budget(),
runtime_ingress_authority: None,
}
}
pub fn with_session_service(
handle: MobHandle,
session_service: Arc<dyn MobSessionService>,
) -> Self {
Self {
handle,
session_store: None,
session_service: Some(session_service),
continuity_session_store: None,
runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
generated_external_owner_session: std::sync::OnceLock::new(),
committed_boundary_recoverer: None,
resume_divergence_logged: std::sync::Mutex::new(std::collections::HashSet::new()),
actor_admission_budget: bridge_actor_admission_budget(),
runtime_ingress_authority: None,
}
}
pub fn with_session_store(
handle: MobHandle,
session_store: Arc<dyn meerkat::SessionStore>,
) -> Self {
Self {
handle,
session_store: Some(session_store),
session_service: None,
continuity_session_store: None,
runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
generated_external_owner_session: std::sync::OnceLock::new(),
committed_boundary_recoverer: None,
resume_divergence_logged: std::sync::Mutex::new(std::collections::HashSet::new()),
actor_admission_budget: bridge_actor_admission_budget(),
runtime_ingress_authority: None,
}
}
pub fn with_session_store_and_service(
handle: MobHandle,
session_store: Arc<dyn meerkat::SessionStore>,
session_service: Arc<dyn MobSessionService>,
) -> Self {
Self {
handle,
session_store: Some(session_store),
session_service: Some(session_service),
continuity_session_store: None,
runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
generated_external_owner_session: std::sync::OnceLock::new(),
committed_boundary_recoverer: None,
resume_divergence_logged: std::sync::Mutex::new(std::collections::HashSet::new()),
actor_admission_budget: bridge_actor_admission_budget(),
runtime_ingress_authority: None,
}
}
pub fn with_continuity_session_store(
handle: MobHandle,
session_store: Arc<ContinuitySessionStoreAdapter>,
session_service: Option<Arc<dyn MobSessionService>>,
) -> Self {
Self {
handle,
session_store: Some(session_store.clone()),
session_service,
continuity_session_store: Some(session_store),
runtime_members: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_sessions: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
generated_external_owner_session: std::sync::OnceLock::new(),
committed_boundary_recoverer: None,
resume_divergence_logged: std::sync::Mutex::new(std::collections::HashSet::new()),
actor_admission_budget: bridge_actor_admission_budget(),
runtime_ingress_authority: None,
}
}
#[must_use]
pub fn with_actor_admission_budget(mut self, budget: Duration) -> Self {
self.actor_admission_budget = budget;
self
}
#[must_use]
pub fn with_committed_boundary_recoverer(
mut self,
recoverer: Arc<dyn CommittedBoundaryRecoverer>,
) -> Self {
self.committed_boundary_recoverer = Some(recoverer);
self
}
#[must_use]
pub fn with_runtime_ingress_authority(
mut self,
authority: Arc<dyn meerkat_runtime::SessionServiceRuntimeExt>,
) -> Self {
self.runtime_ingress_authority = Some(authority);
self
}
#[must_use]
pub fn actor_admission_budget(&self) -> Duration {
self.actor_admission_budget
}
async fn remember_runtime_member(
&self,
runtime_id: &AgentRuntimeId,
member_id: &MobAgentIdentity,
) {
self.runtime_members.write().await.insert(
runtime_id.as_str().to_string(),
member_id.as_str().to_string(),
);
}
async fn remember_runtime_session(
&self,
runtime_id: &AgentRuntimeId,
session_id: &meerkat_core::types::SessionId,
) {
self.runtime_sessions
.write()
.await
.insert(runtime_id.as_str().to_string(), session_id.clone());
}
async fn forget_runtime_member(&self, runtime_id: &AgentRuntimeId) {
self.runtime_members
.write()
.await
.remove(runtime_id.as_str());
self.runtime_sessions
.write()
.await
.remove(runtime_id.as_str());
}
async fn member_id_for_runtime_id(&self, runtime_id: &AgentRuntimeId) -> MobAgentIdentity {
let members = self.runtime_members.read().await;
members
.get(runtime_id.as_str())
.map(|member| MobAgentIdentity::from(member.as_str()))
.unwrap_or_else(|| crate::member_comms_id::mob_member_id(runtime_id.as_str()))
}
async fn retire_session_owned_member_to_absence(
&self,
member_id: &MobAgentIdentity,
) -> Result<(), meerkat_mob::MobError> {
let retained_cleanup_error = match self.handle.retire(member_id.clone()).await {
Ok(()) | Err(meerkat_mob::MobError::MemberNotFound(_)) => None,
Err(error) if is_recoverable_session_owned_retire_cleanup_error(&error.to_string()) => {
Some(error)
}
Err(error) => return Err(error),
};
if let Some(initial_error) = retained_cleanup_error {
tracing::warn!(
member_id = %member_id,
error = %initial_error,
"session-owned retire retained a cleanup anchor; retrying exact incarnation"
);
match self.handle.retire(member_id.clone()).await {
Ok(()) | Err(meerkat_mob::MobError::MemberNotFound(_)) => {}
Err(retry_error) => {
return Err(meerkat_mob::MobError::Internal(format!(
"session-owned retire cleanup retry failed for {member_id}: initial: \
{initial_error}; retry: {retry_error}"
)));
}
}
}
if self
.handle
.list_all_members()
.await
.iter()
.any(|entry| entry.agent_identity == *member_id)
{
return Err(meerkat_mob::MobError::Internal(format!(
"session-owned retire reported success but retained roster anchor {member_id}"
)));
}
Ok(())
}
fn resolved_runtime_ingress_authority(
&self,
) -> Option<Arc<dyn meerkat_runtime::SessionServiceRuntimeExt>> {
if let Some(explicit) = self.runtime_ingress_authority.as_ref() {
return Some(Arc::clone(explicit));
}
self.session_service
.as_ref()?
.runtime_adapter()
.map(|machine| machine as Arc<dyn meerkat_runtime::SessionServiceRuntimeExt>)
}
async fn capture_pending_member_ingress(
&self,
session_id: &meerkat_core::types::SessionId,
) -> PendingIngressCapture {
let Some(authority) = self.resolved_runtime_ingress_authority() else {
tracing::debug!(
session_id = %session_id,
"no runtime ingress authority; repair cannot observe or carry queued member inputs"
);
return PendingIngressCapture::empty();
};
let active = match authority.list_active_inputs(session_id).await {
Ok(ids) => ids,
Err(error) => {
tracing::debug!(
session_id = %session_id,
error = %error,
"pending-ingress probe unavailable before repair disposal"
);
return PendingIngressCapture::empty();
}
};
let mut capture = PendingIngressCapture::empty();
for input_id in active {
let stored = match authority.input_state(session_id, &input_id).await {
Ok(Some(stored)) => stored,
Ok(None) => continue,
Err(error) => {
capture
.uncarryable
.push((input_id, "state-unreadable", error.to_string()));
continue;
}
};
if stored.seed.terminal_outcome.is_some() {
continue;
}
match stored.seed.phase {
meerkat_runtime::InputLifecycleState::Accepted
| meerkat_runtime::InputLifecycleState::Queued => {}
meerkat_runtime::InputLifecycleState::Staged
| meerkat_runtime::InputLifecycleState::Applied
| meerkat_runtime::InputLifecycleState::AppliedPendingConsumption => {
capture.uncarryable.push((
input_id,
"mid-run",
format!(
"input was {:?} at repair disposal; the disposal cancel \
terminalizes it and the sender observes that terminal",
stored.seed.phase
),
));
continue;
}
meerkat_runtime::InputLifecycleState::Consumed
| meerkat_runtime::InputLifecycleState::Superseded
| meerkat_runtime::InputLifecycleState::Coalesced
| meerkat_runtime::InputLifecycleState::Abandoned => continue,
other => {
capture.uncarryable.push((
input_id,
"unrecognized-phase",
format!("input was in unrecognized lifecycle phase {other:?}"),
));
continue;
}
}
match stored.state.persisted_input {
Some(
input @ (meerkat_runtime::Input::Prompt(_)
| meerkat_runtime::Input::Peer(_)
| meerkat_runtime::Input::ExternalEvent(_)),
) => {
capture.carryable.push(CarriedMemberInput {
original_input_id: input_id,
admission_sequence: stored.seed.admission_sequence,
input,
});
}
Some(meerkat_runtime::Input::FlowStep(_)) => {
capture.uncarryable.push((
input_id,
"flow-step",
"flow-step correlation is owned by the flow engine and cannot be \
re-admitted raw; the loss is bounded to the interrupted flow, \
which observes its step's terminal and owns the retry"
.to_string(),
));
}
Some(
meerkat_runtime::Input::Continuation(_) | meerkat_runtime::Input::Operation(_),
) => {
capture.uncarryable.push((
input_id,
"runtime-internal",
"continuation/operation inputs are machine-internal and cannot be \
re-admitted raw; the successor runtime re-derives its own; the \
loss is bounded to the disposed runtime's in-flight bookkeeping"
.to_string(),
));
}
Some(_) => {
capture.uncarryable.push((
input_id,
"unrecognized-class",
"the pending input's class is unknown to this build; no carry \
lane exists for it"
.to_string(),
));
}
None => {
capture.uncarryable.push((
input_id,
"payload-unavailable",
"the runtime retained no payload for this pending input".to_string(),
));
}
}
}
capture
.carryable
.sort_by_key(|entry| entry.admission_sequence.unwrap_or(u64::MAX));
capture
}
fn log_pending_ingress_before_repair_disposal(
&self,
member_id: &MobAgentIdentity,
session_id: &meerkat_core::types::SessionId,
capture: &PendingIngressCapture,
) {
if capture.is_empty() {
return;
}
tracing::warn!(
member_id = %member_id,
session_id = %session_id,
carryable = capture.carryable.len(),
destroyed = capture.uncarryable.len(),
"repair disposal found pending queued inputs on the member: carrying \
the carryable set to the healed successor; anything listed below is \
destroyed with the member"
);
for (input_id, class, reason) in &capture.uncarryable {
tracing::warn!(
member_id = %member_id,
session_id = %session_id,
input_id = %input_id,
class,
reason = %reason,
"repair disposal DESTROYS a pending member input it cannot carry"
);
}
}
async fn readmit_carried_inputs(
&self,
member_id: &MobAgentIdentity,
session_id: &meerkat_core::types::SessionId,
capture: PendingIngressCapture,
) {
if capture.carryable.is_empty() {
return;
}
let Some(authority) = self.resolved_runtime_ingress_authority() else {
tracing::error!(
member_id = %member_id,
session_id = %session_id,
lost = capture.carryable.len(),
"runtime ingress authority disappeared between capture and carry; \
captured queued inputs are lost"
);
return;
};
let total = capture.carryable.len();
let mut carried = 0usize;
for entry in capture.carryable {
let CarriedMemberInput {
original_input_id,
input,
..
} = entry;
let Some(input) = remint_carried_input_identity(input) else {
tracing::error!(
member_id = %member_id,
session_id = %session_id,
original_input_id = %original_input_id,
"carried input has no re-identifiable header in this build; the \
input is lost"
);
continue;
};
let readmitted_input_id = input.id().clone();
match authority.accept_input(session_id, input).await {
Ok(meerkat_runtime::AcceptOutcome::Accepted { .. }) => {
carried += 1;
tracing::info!(
member_id = %member_id,
session_id = %session_id,
original_input_id = %original_input_id,
readmitted_input_id = %readmitted_input_id,
"carried a queued member input into the healed successor session"
);
}
Ok(other) => {
let outcome = match &other {
meerkat_runtime::AcceptOutcome::Deduplicated { .. } => "deduplicated",
meerkat_runtime::AcceptOutcome::Rejected { .. } => "rejected",
_ => "unrecognized",
};
tracing::error!(
member_id = %member_id,
session_id = %session_id,
original_input_id = %original_input_id,
outcome,
"successor admission did not accept a carried queued input; \
the input is lost"
);
}
Err(error) => {
tracing::error!(
member_id = %member_id,
session_id = %session_id,
original_input_id = %original_input_id,
error = %error,
"failed to re-admit a carried queued input into the healed \
successor; the input is lost"
);
}
}
}
tracing::warn!(
member_id = %member_id,
session_id = %session_id,
carried,
total,
"repair carried queued member inputs into the healed successor session"
);
}
async fn member_wires(
&self,
require_reciprocal: bool,
) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
let members = self.handle.list_members_including_retiring().await;
let runtime_members = self.runtime_members.read().await;
let member_runtimes = runtime_members
.iter()
.map(|(runtime, member)| (member.clone(), runtime.clone()))
.collect::<HashMap<_, _>>();
let active_ids = members
.iter()
.map(|member| member.agent_identity.to_string())
.collect::<std::collections::BTreeSet<_>>();
let mut directed = std::collections::BTreeSet::new();
for member in &members {
let a = member.agent_identity.to_string();
for peer in &member.wired_to {
let b = peer.to_string();
if active_ids.contains(&b) {
directed.insert((a.clone(), b));
}
}
}
let edges = directed
.iter()
.filter(|(a, b)| !require_reciprocal || directed.contains(&(b.clone(), a.clone())))
.map(|(a, b)| {
if a <= b {
(a.clone(), b.clone())
} else {
(b.clone(), a.clone())
}
})
.collect::<std::collections::BTreeSet<_>>();
Ok(edges
.into_iter()
.filter_map(|(a, b)| {
let a = member_runtimes
.get(&a)
.cloned()
.unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&a).into_owned());
let b = member_runtimes
.get(&b)
.cloned()
.unwrap_or_else(|| crate::member_comms_id::runtime_alias_str(&b).into_owned());
Some((
AgentRuntimeId::parse(&a).ok()?,
AgentRuntimeId::parse(&b).ok()?,
))
})
.collect())
}
async fn runtime_session_id(
&self,
runtime_id: &AgentRuntimeId,
) -> Option<meerkat_core::types::SessionId> {
self.runtime_sessions
.read()
.await
.get(runtime_id.as_str())
.cloned()
}
async fn durable_session_row_is_absent(
&self,
session_id: &meerkat_core::types::SessionId,
) -> bool {
if let Some(store) = self.session_store.as_ref() {
return matches!(store.load_meta(session_id).await, Ok(None));
}
if let Some(store) = self.continuity_session_store.as_ref() {
use meerkat::SessionStore as _;
return matches!(store.load_meta(session_id).await, Ok(None));
}
if let Some(service) = self.session_service.as_ref() {
return matches!(
service.read(session_id).await,
Err(meerkat_core::SessionError::NotFound { .. })
);
}
false
}
async fn resume_source_confirmed_absent(
&self,
session_id: &meerkat_core::types::SessionId,
) -> bool {
let mut consulted_authority = false;
if let Some(store) = self.continuity_session_store.as_ref() {
use meerkat::SessionStore as _;
match store.load_meta(session_id).await {
Ok(Some(_)) => return false,
Ok(None) => consulted_authority = true,
Err(_) => return false,
}
}
if let Some(service) = self.session_service.as_ref() {
match service.session_known_to_archive_authority(session_id).await {
Ok(true) => return false,
Ok(false) => consulted_authority = true,
Err(_) => return false,
}
}
consulted_authority
}
async fn verify_durable_session_after_rejected_resume(
&self,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
) {
let Some(store) = self.session_store.as_ref() else {
return;
};
match store.load_meta(session_id).await {
Ok(Some(_)) => {}
Ok(None) => {
tracing::error!(
identity = %identity,
session_id = %session_id,
"durable session is GONE after the rejected resume (Bug I class); this \
should be impossible on meerkat >=0.7.29 (ask 31: non-destructive resume \
rollback + auto-revival) — report upstream; recovery needs a pre-damage \
store restore"
);
}
Err(error) => {
tracing::warn!(
identity = %identity,
session_id = %session_id,
%error,
"could not verify durable session presence after the rejected resume"
);
}
}
}
fn base_profile_for_spec(&self, spec: &DurableAgentSpec) -> Option<meerkat_mob::Profile> {
self.handle
.definition()
.resolve_inline_profile(&spec.profile)
.cloned()
}
async fn log_unmasked_resume_divergence(
&self,
identity: &AgentIdentity,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
session_id: &meerkat_core::types::SessionId,
) {
let Some(profile) = base_profile else {
return;
};
let mask = profile.resume_override_mask();
if mask.model && mask.provider {
return;
}
{
let logged = self
.resume_divergence_logged
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if logged.contains(identity.as_str()) {
return;
}
}
let Some(service) = self.session_service.as_ref() else {
return;
};
let metadata = match service.load_persisted_session_metadata(session_id).await {
Ok(Some(view)) => view.session_metadata,
Ok(None) => None,
Err(error) => {
tracing::debug!(
identity = %identity,
session_id = %session_id,
%error,
"resume-divergence check skipped: durable metadata read failed"
);
None
}
};
let Some(metadata) = metadata else {
return;
};
let declared_model = draft.model.as_ref().unwrap_or(&profile.model);
let divergence = unmasked_resume_divergence(
&mask,
declared_model,
profile.provider,
&metadata.model,
metadata.provider,
);
if divergence.model || divergence.provider {
tracing::info!(
identity = %identity,
profile = %spec.profile.as_str(),
session_id = %session_id,
restored_model = %metadata.model,
restored_provider = %metadata.provider.as_str(),
profile_model = %declared_model,
profile_provider = ?profile.provider.map(|provider| provider.as_str()),
model_unmasked_divergent = divergence.model,
provider_unmasked_divergent = divergence.provider,
"resume restored an LLM identity (model, provider) that differs from the \
profile declaration; durable metadata wins for the unmasked fields (no \
resume_overrides mask covers them)"
);
self.resume_divergence_logged
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(identity.as_str().to_string());
}
}
async fn resolve_runtime_session_id(
&self,
runtime_id: &AgentRuntimeId,
member_id: &MobAgentIdentity,
missing_message: &'static str,
deadline: &ActorAdmissionDeadline,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
if let Some(session_id) = self.handle.resolve_bridge_session_id(member_id).await {
self.remember_runtime_session(runtime_id, &session_id).await;
return Ok(session_id);
}
if member_id.as_str() != runtime_id.as_str()
&& deadline
.bound(
"resolve_session.get_member",
member_id,
self.handle.get_member(member_id),
)
.await?
.map_err(|err| BridgeError::Mob(err.to_string()))?
.is_some()
&& let Some(session_id) = self.runtime_session_id(runtime_id).await
{
return Ok(session_id);
}
Err(BridgeError::Mob(missing_message.to_string()))
}
async fn repair_member_for_delivery(
&self,
runtime_id: &AgentRuntimeId,
member_id: &MobAgentIdentity,
member_entry_before_delivery: Option<(meerkat_mob::ProfileName, BTreeMap<String, String>)>,
) -> Result<(), BridgeError> {
if let (Some(session_id), Some((role, labels))) = (
self.runtime_session_id(runtime_id).await,
member_entry_before_delivery.clone(),
) {
let capture = self.capture_pending_member_ingress(&session_id).await;
self.log_pending_ingress_before_repair_disposal(member_id, &session_id, &capture);
match self
.resume_repair_member(
runtime_id,
member_id,
role.clone(),
labels.clone(),
&session_id,
)
.await
{
Ok(()) => {
self.readmit_carried_inputs(member_id, &session_id, capture)
.await;
return Ok(());
}
Err(RepairResumeFailure::DurableSnapshotMissing { detail }) => {
tracing::warn!(
runtime_id = %runtime_id,
member_id = %member_id,
session_id = %session_id,
detail = %detail,
"durable session snapshot is gone; repairing with a fresh spawn \
under the same identity (no transcript left to preserve)"
);
self.handle
.ensure_member(fresh_member_spec_from_pre_delivery_entry(
member_id, role, labels,
))
.await
.map_err(|e| BridgeError::Mob(e.to_string()))?;
self.remember_runtime_member(runtime_id, member_id).await;
if let Some(fresh_session_id) =
self.handle.resolve_bridge_session_id(member_id).await
{
self.readmit_carried_inputs(member_id, &fresh_session_id, capture)
.await;
} else if !capture.carryable.is_empty() {
tracing::error!(
runtime_id = %runtime_id,
member_id = %member_id,
lost = capture.carryable.len(),
"fresh repair spawn has no resolvable session id; the \
captured queued inputs are lost"
);
}
return Ok(());
}
Err(RepairResumeFailure::Rejected(err)) => {
if !capture.carryable.is_empty() {
tracing::error!(
runtime_id = %runtime_id,
member_id = %member_id,
session_id = %session_id,
lost = capture.carryable.len(),
"delivery repair failed after its retire; the captured \
queued inputs are lost with the disposed member"
);
}
return Err(err);
}
}
}
match self.handle.respawn(member_id.clone(), None).await {
Ok(_) => Ok(()),
Err(respawn_err) => match classify_member_repair_respawn_failure(&respawn_err) {
MemberRepairRespawnFailure::DegradedTopologyRestore { failed_peer_ids } => {
tracing::warn!(
runtime_id = %runtime_id,
member_id = %member_id,
failed_peer_count = failed_peer_ids.len(),
failed_peer_ids = ?failed_peer_ids,
"identity bridge respawn restored member with isolated peer edges; continuing delivery"
);
Ok(())
}
MemberRepairRespawnFailure::RecoverableCleanup => {
if self
.handle
.get_member(member_id)
.await
.map_err(|err| BridgeError::Mob(err.to_string()))?
.is_none()
&& let Some((role, labels)) = member_entry_before_delivery
{
self.handle
.ensure_member(fresh_member_spec_from_pre_delivery_entry(
member_id, role, labels,
))
.await
.map_err(|e| BridgeError::Mob(e.to_string()))?;
}
Ok(())
}
MemberRepairRespawnFailure::Fatal(message) => Err(BridgeError::Mob(message)),
},
}
}
async fn resume_repair_member(
&self,
runtime_id: &AgentRuntimeId,
member_id: &MobAgentIdentity,
role: meerkat_mob::ProfileName,
labels: BTreeMap<String, String>,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), RepairResumeFailure> {
if let Err(err) = self.retire_session_owned_member_to_absence(member_id).await {
return Err(RepairResumeFailure::Rejected(BridgeError::Mob(format!(
"repair retire before resume: {err}"
))));
}
self.forget_runtime_member(runtime_id).await;
let mut spec = SpawnMemberSpec::new(role, member_id.clone());
if !labels.is_empty() {
spec = spec.with_labels(labels);
}
spec.launch_mode = MemberLaunchMode::Resume {
bridge_session_id: session_id.clone(),
};
match self.spawn_member_spec(spec).await {
Ok(()) => {
self.remember_runtime_member(runtime_id, member_id).await;
self.remember_runtime_session(runtime_id, session_id).await;
tracing::info!(
runtime_id = %runtime_id,
member_id = %member_id,
session_id = %session_id,
"delivery repair resumed the member onto its durable session \
(transcript preserved, no session rotation)"
);
Ok(())
}
Err(err) if durable_snapshot_is_typed_absent(&err) => {
Err(RepairResumeFailure::DurableSnapshotMissing {
detail: err.to_string(),
})
}
Err(err) => {
let kind = classify_resume_error(&err);
tracing::error!(
runtime_id = %runtime_id,
member_id = %member_id,
session_id = %session_id,
kind = ?kind,
error = %err,
"delivery repair resume rejected; durable session preserved, \
delivery fails loudly (refusing fresh-spawn fallback)"
);
Err(RepairResumeFailure::Rejected(BridgeError::ResumeRejected {
kind,
detail: format!("delivery repair resume: {err}"),
}))
}
}
}
fn external_owner_bridge_session_id(&self) -> meerkat_core::types::SessionId {
if let Some(authority) = self.handle.owner_bridge_session_lifecycle_authority() {
return authority.bridge_session_id;
}
self.generated_external_owner_session
.get_or_init(meerkat_core::types::SessionId::new)
.clone()
}
async fn spawn_member_spec(
&self,
spawn_spec: SpawnMemberSpec,
) -> Result<(), meerkat_mob::MobError> {
if spawn_spec_requires_generated_owner_context(&spawn_spec) {
let owner_session_id = self.external_owner_bridge_session_id();
Box::pin(
self.handle
.spawn_spec_with_generated_owner_context(spawn_spec, owner_session_id),
)
.await
.map(|_| ())
} else {
Box::pin(self.handle.spawn_spec(spawn_spec))
.await
.map(|_| ())
}
}
}
fn peer_reachable_count_from_connectivity(
connectivity: Option<&meerkat_contracts::WirePeerConnectivity>,
) -> Option<usize> {
match connectivity {
Some(meerkat_contracts::WirePeerConnectivity::Known { snapshot })
if snapshot.unknown_peer_count == 0 =>
{
Some(snapshot.reachable_peer_count)
}
Some(meerkat_contracts::WirePeerConnectivity::Known { .. }) => None,
Some(
meerkat_contracts::WirePeerConnectivity::NotApplicable
| meerkat_contracts::WirePeerConnectivity::ProbeTimedOut,
)
| None => None,
}
}
pub(crate) fn spawn_spec_requires_generated_owner_context(spawn_spec: &SpawnMemberSpec) -> bool {
matches!(
spawn_spec.binding,
Some(meerkat_mob::RuntimeBinding::External { .. })
)
}
fn spec_uses_external_binding(spec: &DurableAgentSpec) -> bool {
matches!(spec.backend, Some(meerkat_mob::MobBackendKind::External))
|| matches!(
spec.binding.as_ref(),
Some(meerkat_contracts::WireRuntimeBinding::External { .. })
)
}
fn member_id_for_spawn_spec(
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
) -> MobAgentIdentity {
if spec_uses_external_binding(spec) {
crate::member_comms_id::mob_member_id(spec.identity.as_str())
} else {
crate::member_comms_id::mob_member_id(runtime_id.as_str())
}
}
fn identity_prompt_cache_key(identity: &AgentIdentity) -> String {
format!("mobkit:{}", identity.as_str())
}
fn merge_provider_params_missing_from(
target: &mut ProviderParamsOverride,
defaults: &ProviderParamsOverride,
) -> Result<(), BridgeError> {
fn fill<T: Clone>(target: &mut Option<T>, default: &Option<T>) {
if target.is_none()
&& let Some(value) = default
{
*target = Some(value.clone());
}
}
fill(&mut target.temperature, &defaults.temperature);
fill(&mut target.top_p, &defaults.top_p);
fill(&mut target.max_output_tokens, &defaults.max_output_tokens);
fill(&mut target.reasoning, &defaults.reasoning);
fill(
&mut target.thinking_budget_tokens,
&defaults.thinking_budget_tokens,
);
match (target.provider_tag.as_mut(), defaults.provider_tag.as_ref()) {
(Some(tag), Some(default)) => tag
.merge_missing_from(default)
.map_err(|error| BridgeError::InvalidInput(format!("provider params: {error}")))?,
(None, Some(default)) => target.provider_tag = Some(default.clone()),
_ => {}
}
Ok(())
}
fn spawn_is_openai_backed(
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
) -> bool {
if let Some(model) = draft.model.as_deref() {
return matches!(
meerkat_models::infer_provider(model),
Some(meerkat_core::Provider::OpenAI)
);
}
match base_profile {
Some(profile) => match profile.provider {
Some(provider) => matches!(provider, meerkat_core::Provider::OpenAI),
None => matches!(
meerkat_models::infer_provider(&profile.model),
Some(meerkat_core::Provider::OpenAI)
),
},
None => false,
}
}
fn ensure_default_prompt_cache_key(params: &mut ProviderParamsOverride, identity: &AgentIdentity) {
match params.provider_tag.as_mut() {
Some(ProviderTag::OpenAi(tag)) => {
if tag.prompt_cache_key.is_none() {
tag.prompt_cache_key = Some(identity_prompt_cache_key(identity));
}
}
Some(_) => {}
None => {
params.provider_tag = Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_key: Some(identity_prompt_cache_key(identity)),
..Default::default()
}));
}
}
}
fn apply_provider_params(
spawn_spec: &mut SpawnMemberSpec,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
) -> Result<(), BridgeError> {
let declared = draft.provider_params.clone();
if declared.is_none() && spawn_spec.override_profile.is_none() {
return Ok(());
}
let Some(mut profile) = spawn_spec
.override_profile
.take()
.or_else(|| base_profile.cloned())
else {
return Err(BridgeError::InvalidInput(format!(
"identity {} declares provider_params but profile '{}' resolves to no inline \
definition profile to carry them; declare provider_params on the realm profile \
instead",
spec.identity.as_str(),
spec.profile.as_str(),
)));
};
let mut params = match declared {
Some(mut declared) => {
if let Some(profile_params) = profile.provider_params.as_ref() {
merge_provider_params_missing_from(&mut declared, profile_params)?;
}
if !profile
.resume_overrides
.contains(&ResumeOverrideField::ProviderParams)
{
profile
.resume_overrides
.push(ResumeOverrideField::ProviderParams);
}
declared
}
None => profile.provider_params.clone().unwrap_or_default(),
};
if spawn_is_openai_backed(draft, base_profile) {
ensure_default_prompt_cache_key(&mut params, &spec.identity);
}
profile.provider_params = (!params.is_empty()).then_some(params);
spawn_spec.override_profile = Some(profile);
Ok(())
}
pub(crate) fn build_spawn_spec(
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
) -> Result<SpawnMemberSpec, BridgeError> {
let mid = member_id_for_spawn_spec(runtime_id, spec);
let mut spawn_spec = SpawnMemberSpec::new(spec.profile.clone(), mid);
if let Some(message) = spec.initial_message.as_ref() {
spawn_spec = spawn_spec.with_initial_message(message.clone());
}
if let Some(runtime_mode) = spec.runtime_mode_override {
spawn_spec = spawn_spec.with_runtime_mode(runtime_mode);
}
spawn_spec.backend = spec.backend;
if let Some(binding) = spec.binding.clone() {
spawn_spec.binding = runtime_binding_from_wire(binding);
}
if let Some(ref ctx) = draft.app_context {
spawn_spec = spawn_spec.with_context(ctx.clone());
}
let mut labels = draft.labels.clone();
labels.insert(
"agent_identity".to_string(),
spec.identity.as_str().to_string(),
);
labels.insert(
"profile_name".to_string(),
spec.profile.as_str().to_string(),
);
if !labels.is_empty() {
spawn_spec = spawn_spec.with_labels(labels);
}
if !draft.additional_instructions.is_empty() {
spawn_spec = spawn_spec.with_additional_instructions(draft.additional_instructions.clone());
}
if let Some(model) = draft.model.as_ref() {
match base_profile {
Some(base) if base.provider.is_some() || base.self_hosted_server_id.is_some() => {
let mut profile = base.clone();
profile.model = model.clone();
profile.provider = meerkat_models::canonical().infer_provider(model);
profile.self_hosted_server_id = None;
spawn_spec.override_profile = Some(profile);
}
_ => {
spawn_spec.model_override = Some(model.clone());
}
}
}
if let Some(system_prompt) = draft.system_prompt.as_ref() {
spawn_spec.system_prompt_override =
Some(SpawnSystemPromptOverride::Replace(system_prompt.clone()));
}
if let Some(dispatcher) = draft.local_external_tools.dispatcher() {
spawn_spec.external_tools = Some(dispatcher);
}
apply_provider_params(&mut spawn_spec, spec, draft, base_profile)?;
Ok(spawn_spec)
}
pub(crate) fn build_resume_spawn_spec(
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
session_id: &meerkat_core::types::SessionId,
) -> Result<SpawnMemberSpec, BridgeError> {
let mut spawn_spec = build_spawn_spec(runtime_id, spec, draft, base_profile)?;
spawn_spec.launch_mode = MemberLaunchMode::Resume {
bridge_session_id: session_id.clone(),
};
spawn_spec.system_prompt_override = None;
Ok(spawn_spec)
}
fn runtime_binding_from_wire(
binding: meerkat_contracts::WireRuntimeBinding,
) -> Option<meerkat_mob::RuntimeBinding> {
match binding {
meerkat_contracts::WireRuntimeBinding::Session => {
Some(meerkat_mob::RuntimeBinding::Session)
}
meerkat_contracts::WireRuntimeBinding::External {
address,
bootstrap_token,
identity,
} => {
let resolved = identity.resolve().ok()?;
Some(meerkat_mob::RuntimeBinding::External {
peer_id: resolved.peer_id.to_string(),
address,
bootstrap_token,
pubkey: resolved.pubkey,
})
}
}
}
#[async_trait]
impl SessionBridge for MobSessionBridge {
async fn raw_member_alias_exists(&self, alias: &str) -> Result<bool, BridgeError> {
let members = self.handle.list_members_including_retiring().await;
let authoritative_members = self.runtime_members.read().await;
Ok(members.iter().any(|member| {
crate::member_comms_id::runtime_alias_str(member.agent_identity.as_str()) == alias
&& !authoritative_members
.values()
.any(|owned| owned == member.agent_identity.as_str())
&& crate::member_comms_id::durable_identity_label(&member.labels) != Some(alias)
}))
}
fn requires_resume_snapshot(&self) -> bool {
false
}
async fn recover_committed_boundary(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<CommittedBoundaryRepair, BridgeError> {
match self.committed_boundary_recoverer.as_ref() {
Some(recoverer) => recoverer.recover_committed_boundary(session_id).await,
None => Ok(CommittedBoundaryRepair::Unsupported),
}
}
async fn create_session(
&self,
_identity: &AgentIdentity,
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let mid = member_id_for_spawn_spec(runtime_id, spec);
let spawn_spec = build_spawn_spec(
runtime_id,
spec,
draft,
self.base_profile_for_spec(spec).as_ref(),
)?;
self.spawn_member_spec(spawn_spec)
.await
.map_err(|e| BridgeError::Mob(e.to_string()))?;
self.remember_runtime_member(runtime_id, &mid).await;
self.remember_runtime_session(runtime_id, session_id).await;
self.resolve_runtime_session_id(
runtime_id,
&mid,
"member spawned but has no session ID",
&ActorAdmissionDeadline::new(self.actor_admission_budget),
)
.await
}
async fn resume_session(
&self,
identity: &AgentIdentity,
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
session_id: &meerkat_core::types::SessionId,
_snapshot: &SessionSnapshot,
) -> Result<ResumeSessionOutcome, BridgeError> {
self.log_unmasked_resume_divergence(
identity,
spec,
draft,
self.base_profile_for_spec(spec).as_ref(),
session_id,
)
.await;
if spec_uses_external_binding(spec) {
let spawn_spec = build_resume_spawn_spec(
runtime_id,
spec,
draft,
self.base_profile_for_spec(spec).as_ref(),
session_id,
)?;
let mid = member_id_for_spawn_spec(runtime_id, spec);
self.spawn_member_spec(spawn_spec).await.map_err(|error| {
resume_rejected(identity, session_id, &error, "external-binding resume")
})?;
self.remember_runtime_member(runtime_id, &mid).await;
self.remember_runtime_session(runtime_id, session_id).await;
return Ok(ResumeSessionOutcome::Resumed {
session_id: session_id.clone(),
});
}
let spawn_spec = build_resume_spawn_spec(
runtime_id,
spec,
draft,
self.base_profile_for_spec(spec).as_ref(),
session_id,
)?;
let mid = member_id_for_spawn_spec(runtime_id, spec);
match self.spawn_member_spec(spawn_spec.clone()).await {
Ok(()) => {
self.remember_runtime_member(runtime_id, &mid).await;
self.remember_runtime_session(runtime_id, session_id).await;
Ok(ResumeSessionOutcome::Resumed {
session_id: session_id.clone(),
})
}
Err(error) if is_member_already_exists_error(&error) => {
tracing::warn!(
identity = %identity,
session_id = %session_id,
error = %error,
"resume_session hit a roster collision; retiring the stale member and retrying resume"
);
if self.resume_source_confirmed_absent(session_id).await {
return Err(resume_rejected(
identity,
session_id,
&meerkat_mob::MobError::Internal(format!(
"collision retire refused: the resume source for \
{session_id} is confirmed absent, so retiring the stale \
member would destroy the only live copy of the session"
)),
"collision retire precondition",
));
}
let capture = self.capture_pending_member_ingress(session_id).await;
self.log_pending_ingress_before_repair_disposal(&mid, session_id, &capture);
if let Err(err) = self.retire_session_owned_member_to_absence(&mid).await {
return Err(resume_rejected(
identity,
session_id,
&err,
"collision retire before resume retry",
));
}
self.forget_runtime_member(runtime_id).await;
if let Err(error) = self.spawn_member_spec(spawn_spec).await {
if !capture.carryable.is_empty() {
tracing::error!(
identity = %identity,
session_id = %session_id,
lost = capture.carryable.len(),
"the collision retire already destroyed the member's queued \
inputs and the resume retry failed; the captured inputs are \
lost with it"
);
}
self.verify_durable_session_after_rejected_resume(identity, session_id)
.await;
return Err(resume_rejected(
identity,
session_id,
&error,
"resume retry after collision",
));
}
self.remember_runtime_member(runtime_id, &mid).await;
self.remember_runtime_session(runtime_id, session_id).await;
self.readmit_carried_inputs(&mid, session_id, capture).await;
Ok(ResumeSessionOutcome::Resumed {
session_id: session_id.clone(),
})
}
Err(error) => {
if durable_snapshot_is_typed_absent(&error)
&& self.durable_session_row_is_absent(session_id).await
{
tracing::warn!(
identity = %identity,
session_id = %session_id,
error = %error,
"resume target is typed-Absent and the durable store has no row for \
it (never-persisted continuity head): falling back to a FRESH spawn \
under a new session id. External row deletion produces this same \
shape - investigate if unexpected"
);
let capture = self.capture_pending_member_ingress(session_id).await;
self.log_pending_ingress_before_repair_disposal(&mid, session_id, &capture);
if let Err(retire_error) =
self.retire_session_owned_member_to_absence(&mid).await
{
return Err(resume_rejected(
identity,
session_id,
&retire_error,
"never-persisted retire before fresh fallback",
));
}
self.forget_runtime_member(runtime_id).await;
let fresh_session_id = meerkat_core::types::SessionId::new();
let created_session_id = match self
.create_session(identity, runtime_id, spec, draft, &fresh_session_id)
.await
{
Ok(created) => created,
Err(create_error) => {
if !capture.carryable.is_empty() {
tracing::error!(
identity = %identity,
session_id = %session_id,
lost = capture.carryable.len(),
"the never-persisted retire already destroyed the \
member's queued inputs and the fresh spawn failed; \
the captured inputs are lost with it"
);
}
return Err(create_error);
}
};
self.readmit_carried_inputs(&mid, &created_session_id, capture)
.await;
return Ok(ResumeSessionOutcome::FreshSpawned {
session_id: created_session_id,
reason: ResumeFallbackReason::NeverPersisted {
detail: error.to_string(),
},
});
}
self.verify_durable_session_after_rejected_resume(identity, session_id)
.await;
Err(resume_rejected(
identity,
session_id,
&error,
"resume spawn",
))
}
}
}
async fn deliver_admitted(
&self,
runtime_id: &AgentRuntimeId,
delivery: BridgeDelivery,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let content = &delivery.content;
let handling_mode = delivery.handling_mode;
let system_prompt = delivery.system_prompt.as_deref();
let injected_context = delivery.injected_context.as_slice();
let interaction_id = delivery.interaction_id.as_deref();
let delivery_identity = delivery.delivery_identity.as_ref();
let mid = self.member_id_for_runtime_id(runtime_id).await;
let mut deadline = ActorAdmissionDeadline::new(self.actor_admission_budget);
let member_entry_before_delivery = deadline
.bound(
"deliver.get_member.pre_delivery",
&mid,
self.handle.get_member(&mid),
)
.await
.ok()
.and_then(Result::ok)
.flatten()
.map(|entry| (entry.role, entry.labels));
if content_input_has_images(content) {
let member_entry = deadline
.bound(
"deliver.get_member.image_capability",
&mid,
self.handle.get_member(&mid),
)
.await?
.map_err(|err| BridgeError::Mob(err.to_string()))?
.ok_or_else(|| {
BridgeError::Mob("member not found while checking image capability".to_string())
})?;
let caps = deadline
.bound(
"deliver.model_capabilities",
&mid,
model_capabilities_for_member(
&self.handle,
self.session_service.as_ref(),
&member_entry.agent_identity,
),
)
.await?;
if !caps.image_input {
return Err(BridgeError::InvalidInput(
"target member model cannot accept image input".to_string(),
));
}
}
match submit_internal_bridge_work(
&self.handle,
&mid,
InternalBridgeWork {
content,
system_prompt,
injected_context,
interaction_id,
delivery_identity,
},
handling_mode,
&deadline,
)
.await
{
Ok(()) => {}
Err(err @ BridgeError::ActorAdmissionTimeout { .. }) => return Err(err),
Err(err) if is_repairable_bridge_delivery_error(&err.to_string()) => {
tracing::warn!(
runtime_id = %runtime_id,
error = %err,
"identity bridge delivery found stale runtime state; repairing member before retry"
);
Box::pin(self.repair_member_for_delivery(
runtime_id,
&mid,
member_entry_before_delivery,
))
.await?;
deadline = ActorAdmissionDeadline::new(self.actor_admission_budget);
submit_internal_bridge_work(
&self.handle,
&mid,
InternalBridgeWork {
content,
system_prompt,
injected_context,
interaction_id,
delivery_identity,
},
handling_mode,
&deadline,
)
.await?;
}
Err(err) => return Err(BridgeError::Mob(err.to_string())),
}
self.resolve_runtime_session_id(
runtime_id,
&mid,
"member has no bridge session after deliver",
&deadline,
)
.await
}
async fn checkpoint_session(
&self,
_runtime_id: &AgentRuntimeId,
session_id: &meerkat_core::types::SessionId,
) -> Result<SessionSnapshot, BridgeError> {
let store = self.session_store.as_ref().ok_or_else(|| {
BridgeError::InvalidInput(
"checkpoint requires a session store but none was configured".to_string(),
)
})?;
let session = store
.load(session_id)
.await
.map_err(|e| BridgeError::Mob(format!("failed to load session for checkpoint: {e}")))?
.ok_or_else(|| {
BridgeError::Mob(format!(
"session {session_id} not found in store for checkpoint"
))
})?;
let data = serde_json::to_vec(&session)
.map_err(|e| BridgeError::Mob(format!("failed to serialize session: {e}")))?;
Ok(SessionSnapshot { data })
}
async fn retire_member(&self, runtime_id: &AgentRuntimeId) -> Result<(), BridgeError> {
let mid = self.member_id_for_runtime_id(runtime_id).await;
self.retire_session_owned_member_to_absence(&mid)
.await
.map_err(|error| BridgeError::Mob(error.to_string()))?;
self.forget_runtime_member(runtime_id).await;
Ok(())
}
async fn retire_reset_superseded_member(
&self,
runtime_id: &AgentRuntimeId,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
let mid = self.member_id_for_runtime_id(runtime_id).await;
let adapter = self.continuity_session_store.as_ref().ok_or_else(|| {
BridgeError::Mob(format!(
"reset retire cannot abandon superseded member {mid}: the bridge has no continuity session-store authority"
))
})?;
abandon_then_retire_reset_superseded(
&mid,
session_id,
|| adapter.abandon_superseded_session(session_id),
|| self.handle.retire(mid.clone()),
)
.await?;
if self
.handle
.list_all_members()
.await
.iter()
.any(|entry| entry.agent_identity == mid)
{
return Err(BridgeError::Mob(format!(
"reset retire reported success but retained roster anchor {mid}"
)));
}
self.forget_runtime_member(runtime_id).await;
self.unregister_session_runtime_state(session_id).await
}
async fn wire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
let member_a = self.member_id_for_runtime_id(a).await;
let member_b = self.member_id_for_runtime_id(b).await;
self.handle
.wire(
meerkat_mob::AgentIdentity::from(member_a.as_str()),
member_b,
)
.await
.map_err(|e| BridgeError::Mob(e.to_string()))
}
async fn wire_peers_batch(
&self,
edges: &[(AgentRuntimeId, AgentRuntimeId)],
) -> Result<(), BridgeError> {
let mut member_edges = Vec::with_capacity(edges.len());
for (a, b) in edges {
let member_a = self.member_id_for_runtime_id(a).await;
let member_b = self.member_id_for_runtime_id(b).await;
member_edges.push((
meerkat_mob::AgentIdentity::from(member_a.as_str()),
meerkat_mob::AgentIdentity::from(member_b.as_str()),
));
}
match self.handle.wire_members_batch(member_edges.clone()).await {
Ok(_) => Ok(()),
Err(error)
if error
.to_string()
.contains("does not support legacy external (peer-only) members") =>
{
for (member_a, member_b) in member_edges {
self.handle
.wire(member_a, member_b)
.await
.map_err(|error| BridgeError::Mob(error.to_string()))?;
}
Ok(())
}
Err(error) => Err(BridgeError::Mob(error.to_string())),
}
}
async fn current_member_wires(
&self,
) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
self.member_wires(true).await
}
async fn current_member_wires_any_half(
&self,
) -> Result<Vec<(AgentRuntimeId, AgentRuntimeId)>, BridgeError> {
self.member_wires(false).await
}
async fn unwire_peer(&self, a: &AgentRuntimeId, b: &AgentRuntimeId) -> Result<(), BridgeError> {
let member_a = self.member_id_for_runtime_id(a).await;
let member_b = self.member_id_for_runtime_id(b).await;
match self
.handle
.unwire(
meerkat_mob::AgentIdentity::from(member_a.as_str()),
member_b,
)
.await
{
Ok(()) => Ok(()),
Err(err) => {
let message = err.to_string();
if message.contains("peer not found") || message.contains("not wired") {
Ok(())
} else {
Err(BridgeError::Mob(message))
}
}
}
}
async fn inspect_member(
&self,
runtime_id: &AgentRuntimeId,
) -> Result<MemberInspection, BridgeError> {
let mid = self.member_id_for_runtime_id(runtime_id).await;
let snap = self
.handle
.member_status(&mid)
.await
.map_err(|e| BridgeError::Mob(e.to_string()))?;
let peer_reachable_count =
match peer_reachable_count_from_connectivity(snap.peer_connectivity.as_ref()) {
Some(count) => count,
None => self
.handle
.get_member(&mid)
.await
.ok()
.flatten()
.map(|entry| entry.wired_to.len())
.unwrap_or(0),
};
Ok(MemberInspection {
output_preview: snap.output_preview.clone(),
is_final: snap.is_final,
peer_reachable_count,
})
}
async fn register_session_runtime_state(
&self,
session_id: &meerkat_core::types::SessionId,
identity: &AgentIdentity,
generation: ContinuityGeneration,
checkpoint_version: CheckpointVersion,
fencing_token: FencingToken,
) -> Result<CheckpointVersion, BridgeError> {
if let Some(adapter) = self.continuity_session_store.as_ref() {
return adapter
.register_session(
session_id,
SessionRuntimeState {
identity: identity.clone(),
generation,
checkpoint_version,
fencing_token,
},
)
.await
.map_err(|err| BridgeError::Mob(format!("continuity register_session: {err}")));
}
Ok(checkpoint_version)
}
async fn suspend_session_runtime_state(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
if let Some(adapter) = self.continuity_session_store.as_ref() {
adapter
.suspend_session(session_id)
.await
.map_err(|err| BridgeError::Mob(format!("continuity suspend_session: {err}")))?;
}
Ok(())
}
async fn unregister_session_runtime_state(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), BridgeError> {
if let Some(adapter) = self.continuity_session_store.as_ref() {
adapter
.unregister_session(session_id)
.await
.map_err(|err| BridgeError::Mob(format!("continuity unregister_session: {err}")))?;
}
Ok(())
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::panic)]
mod tests {
use std::sync::Arc;
#[test]
fn internal_bridge_work_spec_threads_every_carrier() {
let content = meerkat_core::ContentInput::Text("turn content".to_string());
let injected = vec![meerkat_core::ContentInput::Text(
"ambient recall".to_string(),
)];
let interaction = uuid::Uuid::new_v4();
let spec = super::internal_bridge_work_spec(
&content,
Some("per-turn system message"),
&injected,
Some(&interaction.to_string()),
);
assert_eq!(
spec.system_prompt.as_deref(),
Some("per-turn system message")
);
assert_eq!(spec.injected_context, injected);
assert_eq!(
spec.interaction_id,
Some(meerkat_core::interaction::InteractionId(interaction))
);
assert!(matches!(spec.origin, meerkat_mob::WorkOrigin::Internal));
let bare = super::internal_bridge_work_spec(&content, None, &[], None);
assert_eq!(bare.system_prompt, None, "absent carrier stays absent");
assert!(bare.injected_context.is_empty());
assert_eq!(bare.interaction_id, None);
}
use async_trait::async_trait;
use meerkat_core::agent::AgentToolDispatcher;
use meerkat_core::lifecycle::run_primitive::{OpenAiPromptCacheOptions, ReasoningEffort};
use meerkat_core::model_profile::capabilities::{OpenAiPromptCacheMode, OpenAiPromptCacheTtl};
use meerkat_core::types::ToolCallView;
use meerkat_core::{ToolDef, error::ToolError, ops::ToolDispatchOutcome};
use meerkat_mob::{MobRespawnError, MobRuntimeMode};
use super::*;
use crate::identity_first::{AgentAddressability, LocalExternalToolOverlay};
struct EmptyDispatcher;
#[async_trait]
impl AgentToolDispatcher for EmptyDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::from([])
}
async fn dispatch(
&self,
_call: ToolCallView<'_>,
) -> Result<ToolDispatchOutcome, ToolError> {
Err(ToolError::ExecutionFailed {
message: "not implemented".to_string(),
})
}
}
#[test]
fn unmasked_resume_divergence_flags_only_unmasked_declared_differences() {
let unmasked = meerkat_core::service::ResumeOverrideMask::default();
let divergence = unmasked_resume_divergence(
&unmasked,
"claude-opus-4-8",
Some(meerkat_core::Provider::Anthropic),
"claude-sonnet-4-5",
meerkat_core::Provider::OpenAI,
);
assert!(divergence.model, "unmasked differing model must be flagged");
assert!(
divergence.provider,
"unmasked differing declared provider must be flagged"
);
let masked = meerkat_core::service::ResumeOverrideMask {
model: true,
provider: true,
..Default::default()
};
let divergence = unmasked_resume_divergence(
&masked,
"claude-opus-4-8",
Some(meerkat_core::Provider::Anthropic),
"claude-sonnet-4-5",
meerkat_core::Provider::OpenAI,
);
assert!(
!divergence.model && !divergence.provider,
"a mask covering the field silences the tripwire (the profile wins anyway)"
);
let divergence = unmasked_resume_divergence(
&unmasked,
"claude-opus-4-8",
None,
"claude-opus-4-8",
meerkat_core::Provider::OpenAI,
);
assert!(!divergence.model, "an identical model is not a divergence");
assert!(
!divergence.provider,
"an undeclared provider states no intent to diverge from"
);
}
fn durable_spec() -> DurableAgentSpec {
DurableAgentSpec {
identity: AgentIdentity::parse("agent:alpha").expect("identity"),
profile: meerkat_mob::ProfileName::from("worker"),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: Default::default(),
context: None,
additional_instructions: Vec::new(),
initial_message: Some(meerkat_core::ContentInput::Text("hello".to_string())),
runtime_mode_override: Some(MobRuntimeMode::TurnDriven),
backend: None,
binding: None,
}
}
#[test]
fn peer_reachable_count_tri_state_defers_to_wiring_when_unresolved() {
use meerkat_contracts::{WirePeerConnectivity, WirePeerConnectivitySnapshot};
let known = WirePeerConnectivity::Known {
snapshot: WirePeerConnectivitySnapshot {
reachable_peer_count: 3,
unknown_peer_count: 0,
unreachable_peers: Vec::new(),
},
};
assert_eq!(
peer_reachable_count_from_connectivity(Some(&known)),
Some(3),
"a resolved probe owns the count"
);
let structurally_wired_but_unresolved = WirePeerConnectivity::Known {
snapshot: WirePeerConnectivitySnapshot {
reachable_peer_count: 0,
unknown_peer_count: 3,
unreachable_peers: Vec::new(),
},
};
assert_eq!(
peer_reachable_count_from_connectivity(Some(&structurally_wired_but_unresolved)),
None,
"a partially resolved probe must defer to the machine-owned wiring degree"
);
assert_eq!(
peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::NotApplicable)),
None,
"not-applicable must defer to the wiring fallback"
);
assert_eq!(
peer_reachable_count_from_connectivity(Some(&WirePeerConnectivity::ProbeTimedOut)),
None,
"probe timeout must defer to the wiring fallback"
);
assert_eq!(peer_reachable_count_from_connectivity(None), None);
}
#[test]
fn build_spawn_spec_maps_identity_first_overrides() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let mut labels = std::collections::BTreeMap::new();
labels.insert("team".to_string(), "ops".to_string());
let draft = AgentBuildDraft {
model: Some("gpt-test".to_string()),
system_prompt: Some("system override".to_string()),
additional_instructions: vec!["stay focused".to_string()],
labels,
app_context: Some(serde_json::json!({"ticket": 7})),
external_tools: Vec::new(),
local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
provider_params: None,
};
let base_profile: meerkat_mob::Profile =
serde_json::from_value(serde_json::json!({"model": "base-model"}))
.expect("minimal profile");
let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile))
.expect("spawn spec");
assert_eq!(spawn.model_override.as_deref(), Some("gpt-test"));
assert!(
spawn.override_profile.is_none(),
"unpinned profiles must not freeze the whole profile for a model-only override"
);
assert_eq!(
spawn.system_prompt_override,
Some(SpawnSystemPromptOverride::Replace(
"system override".to_string()
))
);
assert!(spawn.external_tools.is_some());
assert_eq!(spawn.runtime_mode, Some(MobRuntimeMode::TurnDriven));
assert_eq!(
spawn.initial_message,
Some(meerkat_core::ContentInput::Text("hello".to_string()))
);
assert_eq!(
spawn
.labels
.as_ref()
.and_then(|labels| labels.get("team"))
.map(String::as_str),
Some("ops")
);
assert_eq!(
spawn
.labels
.as_ref()
.and_then(|labels| labels.get("agent_identity"))
.map(String::as_str),
Some("agent:alpha")
);
assert_eq!(
spawn
.labels
.as_ref()
.and_then(|labels| labels.get("profile_name"))
.map(String::as_str),
Some("worker"),
"identity-first spawn labels must carry the adopted profile so the \
SDK build callback sees the roster profile, not a checkpoint default"
);
assert_eq!(
spawn.role_name.as_str(),
"worker",
"SpawnMemberSpec role remains the authoritative mob profile"
);
}
#[test]
fn build_spawn_spec_keeps_profile_snapshot_for_pinned_provider() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = AgentBuildDraft {
model: Some("gpt-test".to_string()),
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
provider_params: None,
};
let base_profile: meerkat_mob::Profile = serde_json::from_value(
serde_json::json!({"model": "base-model", "provider": "openai"}),
)
.expect("pinned profile");
let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile))
.expect("spawn spec");
assert!(spawn.model_override.is_none());
let profile = spawn
.override_profile
.as_ref()
.expect("pinned profile keeps the snapshot path");
assert_eq!(profile.model.as_str(), "gpt-test");
assert!(
profile.provider.is_none(),
"a catalog-unknown pin carries no provider (config-entry resolution downstream)"
);
}
#[test]
fn build_spawn_spec_derives_pair_provider_for_catalog_model_pin() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = AgentBuildDraft {
model: Some("claude-opus-4-8".to_string()),
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
provider_params: None,
};
let base_profile: meerkat_mob::Profile = serde_json::from_value(serde_json::json!({
"model": "gpt-5.5",
"provider": "openai",
"resume_overrides": ["model", "provider"],
}))
.expect("pinned profile");
let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile))
.expect("spawn spec");
let profile = spawn
.override_profile
.as_ref()
.expect("pinned profile keeps the snapshot path");
assert_eq!(profile.model.as_str(), "claude-opus-4-8");
assert_eq!(
profile.provider,
Some(meerkat_core::Provider::Anthropic),
"the pin's provider must be the DRAFT model's catalog owner, applied as a pair"
);
assert!(
profile
.resume_overrides
.contains(&ResumeOverrideField::Model)
&& profile
.resume_overrides
.contains(&ResumeOverrideField::Provider),
"the base profile's pair mask must ride the snapshot so both fields apply on resume"
);
}
fn profile_with_provider(provider: &str) -> meerkat_mob::Profile {
serde_json::from_value(serde_json::json!({
"model": "base-model",
"provider": provider,
}))
.expect("profile")
}
fn draft_with_provider_params(params: Option<ProviderParamsOverride>) -> AgentBuildDraft {
AgentBuildDraft {
model: None,
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: Default::default(),
provider_params: params,
}
}
fn spawn_params(spawn: &SpawnMemberSpec) -> ProviderParamsOverride {
spawn
.override_profile
.as_ref()
.expect("provider params require a profile snapshot")
.provider_params
.clone()
.expect("profile carries provider params")
}
fn spawn_openai_tag(spawn: &SpawnMemberSpec) -> OpenAiProviderTag {
match spawn_params(spawn).provider_tag {
Some(ProviderTag::OpenAi(tag)) => tag,
other => panic!("expected an OpenAI provider tag, got {other:?}"),
}
}
fn implicit_30m() -> OpenAiPromptCacheOptions {
OpenAiPromptCacheOptions {
mode: Some(OpenAiPromptCacheMode::Implicit),
ttl: Some(OpenAiPromptCacheTtl::ThirtyMinutes),
}
}
#[test]
fn build_spawn_spec_lands_declared_prompt_cache_options() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_options: Some(implicit_30m()),
..Default::default()
})),
..Default::default()
}));
let spawn = build_spawn_spec(
&runtime_id,
&durable_spec(),
&draft,
Some(&profile_with_provider("openai")),
)
.expect("spawn spec");
assert_eq!(
spawn_openai_tag(&spawn).prompt_cache_options,
Some(implicit_30m())
);
assert!(
spawn
.override_profile
.as_ref()
.expect("profile snapshot")
.resume_overrides
.contains(&ResumeOverrideField::ProviderParams),
"an explicit declaration must win over durable metadata on resume"
);
}
#[test]
fn build_spawn_spec_merges_draft_params_over_profile_declaration() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let mut base_profile = profile_with_provider("openai");
base_profile.provider_params = Some(ProviderParamsOverride {
thinking_budget_tokens: Some(8192),
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
reasoning_effort: Some(ReasoningEffort::High),
..Default::default()
})),
..Default::default()
});
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_options: Some(implicit_30m()),
..Default::default()
})),
..Default::default()
}));
let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile))
.expect("spawn spec");
assert_eq!(
spawn_params(&spawn).thinking_budget_tokens,
Some(8192),
"profile-declared knobs survive a draft that sets an unrelated knob"
);
let tag = spawn_openai_tag(&spawn);
assert_eq!(tag.reasoning_effort, Some(ReasoningEffort::High));
assert_eq!(tag.prompt_cache_options, Some(implicit_30m()));
}
#[test]
fn build_spawn_spec_defaults_prompt_cache_key_from_identity() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_options: Some(implicit_30m()),
..Default::default()
})),
..Default::default()
}));
let spawn = build_spawn_spec(
&runtime_id,
&durable_spec(),
&draft,
Some(&profile_with_provider("openai")),
)
.expect("spawn spec");
assert_eq!(
spawn_openai_tag(&spawn).prompt_cache_key.as_deref(),
Some("mobkit:agent:alpha")
);
}
#[test]
fn build_spawn_spec_keeps_caller_supplied_prompt_cache_key() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_key: Some("tenant-a:shared-prefix".to_string()),
..Default::default()
})),
..Default::default()
}));
let spawn = build_spawn_spec(
&runtime_id,
&durable_spec(),
&draft,
Some(&profile_with_provider("openai")),
)
.expect("spawn spec");
assert_eq!(
spawn_openai_tag(&spawn).prompt_cache_key.as_deref(),
Some("tenant-a:shared-prefix")
);
}
#[test]
fn build_spawn_spec_prompt_cache_key_is_stable_across_builds() {
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_options: Some(implicit_30m()),
..Default::default()
})),
..Default::default()
}));
let profile = profile_with_provider("openai");
let first = build_spawn_spec(
&AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id"),
&durable_spec(),
&draft,
Some(&profile),
)
.expect("first spawn spec");
let second = build_spawn_spec(
&AgentRuntimeId::parse("rt:agent:alpha:7").expect("runtime id"),
&durable_spec(),
&draft,
Some(&profile),
)
.expect("second spawn spec");
assert_eq!(
spawn_openai_tag(&first).prompt_cache_key,
spawn_openai_tag(&second).prompt_cache_key
);
assert_eq!(
spawn_openai_tag(&first).prompt_cache_key.as_deref(),
Some("mobkit:agent:alpha")
);
}
#[test]
fn build_spawn_spec_skips_prompt_cache_key_for_non_openai_identity() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
temperature: Some(0.2),
..Default::default()
}));
let spawn = build_spawn_spec(
&runtime_id,
&durable_spec(),
&draft,
Some(&profile_with_provider("anthropic")),
)
.expect("spawn spec");
let params = spawn_params(&spawn);
assert_eq!(params.temperature, Some(0.2));
assert!(
params.provider_tag.is_none(),
"no OpenAI tag may be fabricated for an Anthropic-backed identity"
);
}
#[test]
fn build_spawn_spec_without_provider_params_keeps_field_scoped_path() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(None);
let spawn = build_spawn_spec(
&runtime_id,
&durable_spec(),
&draft,
Some(&profile_with_provider("openai")),
)
.expect("spawn spec");
assert!(
spawn.override_profile.is_none(),
"an undeclared draft must not mint a profile snapshot"
);
assert!(spawn.model_override.is_none());
}
#[test]
fn build_spawn_spec_rejects_declared_provider_params_without_inline_profile() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
temperature: Some(0.2),
..Default::default()
}));
let error = build_spawn_spec(&runtime_id, &durable_spec(), &draft, None)
.expect_err("no inline profile can carry provider params");
match error {
BridgeError::InvalidInput(detail) => {
assert!(detail.contains("provider_params"), "{detail}");
}
other => panic!("expected InvalidInput, got {other:?}"),
}
}
#[test]
fn build_spawn_spec_rejects_provider_tag_family_conflict() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let mut base_profile = profile_with_provider("openai");
base_profile.provider_params = Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::Anthropic(Default::default())),
..Default::default()
});
let draft = draft_with_provider_params(Some(ProviderParamsOverride {
provider_tag: Some(ProviderTag::OpenAi(OpenAiProviderTag {
prompt_cache_options: Some(implicit_30m()),
..Default::default()
})),
..Default::default()
}));
let error = build_spawn_spec(&runtime_id, &durable_spec(), &draft, Some(&base_profile))
.expect_err("provider families must not be unioned");
match error {
BridgeError::InvalidInput(detail) => {
assert!(detail.contains("provider params"), "{detail}");
}
other => panic!("expected InvalidInput, got {other:?}"),
}
}
#[test]
fn build_spawn_spec_model_override_works_without_base_profile() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = AgentBuildDraft {
model: Some("gpt-test".to_string()),
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: LocalExternalToolOverlay::new(Arc::new(EmptyDispatcher)),
provider_params: None,
};
let spawn =
build_spawn_spec(&runtime_id, &durable_spec(), &draft, None).expect("spawn spec");
assert_eq!(spawn.model_override.as_deref(), Some("gpt-test"));
assert!(spawn.override_profile.is_none());
}
#[test]
fn fresh_fallback_collision_classifier_matches_member_already_exists() {
let error = meerkat_mob::MobError::MemberAlreadyExists(meerkat_mob::AgentIdentity::from(
"rt-agent-alpha-0",
));
assert!(
is_member_already_exists_error(&error),
"fresh fallback must retry recreate-over-running-member collisions"
);
}
#[test]
fn build_spawn_spec_maps_remote_runtime_binding() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let mut spec = durable_spec();
spec.backend = Some(meerkat_mob::MobBackendKind::External);
spec.binding = Some(
serde_json::from_value(serde_json::json!({
"kind": "external",
"address": "tcp://127.0.0.1:4777",
"identity": {
"kind": "ed25519_public_key",
"public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
}
}))
.expect("wire binding"),
);
let draft = AgentBuildDraft {
model: None,
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: Default::default(),
provider_params: None,
};
let spawn = build_spawn_spec(&runtime_id, &spec, &draft, None).expect("spawn spec");
assert_eq!(
spawn.identity.as_str(),
crate::member_comms_id::mob_member_id_str("agent:alpha").as_ref()
);
assert_eq!(spawn.backend, Some(meerkat_mob::MobBackendKind::External));
assert!(
matches!(
spawn.binding,
Some(meerkat_mob::RuntimeBinding::External { .. })
),
"expected external runtime binding, got {:?}",
spawn.binding
);
if let Some(meerkat_mob::RuntimeBinding::External {
address, pubkey, ..
}) = spawn.binding
{
assert_eq!(address.as_str(), "tcp://127.0.0.1:4777");
assert_eq!(pubkey, [7; 32]);
}
}
#[test]
fn external_binding_spawn_specs_require_generated_owner_context() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = AgentBuildDraft {
model: None,
system_prompt: None,
additional_instructions: Vec::new(),
labels: Default::default(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: Default::default(),
provider_params: None,
};
let session_spawn =
build_spawn_spec(&runtime_id, &durable_spec(), &draft, None).expect("spawn spec");
assert!(
!spawn_spec_requires_generated_owner_context(&session_spawn),
"session-backed spawns must not require a generated owner context"
);
let mut spec = durable_spec();
spec.backend = Some(meerkat_mob::MobBackendKind::External);
spec.binding = Some(
serde_json::from_value(serde_json::json!({
"kind": "external",
"address": "tcp://127.0.0.1:4777",
"identity": {
"kind": "ed25519_public_key",
"public_key": "ed25519:BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc="
}
}))
.expect("wire binding"),
);
let external_spawn =
build_spawn_spec(&runtime_id, &spec, &draft, None).expect("spawn spec");
assert!(
spawn_spec_requires_generated_owner_context(&external_spawn),
"external peer-only spawns must carry a generated owner binding on meerkat 0.7.1"
);
}
#[test]
fn bridge_delivery_repair_covers_missing_bridge_session_snapshot() {
let error = "session bridge mob error: member rt:review:singleton:0 failed to restore session 019e5fc2-dad4-77e2-abbe-a8a66bc15f66: missing bridge session snapshot for '019e5fc2-dad4-77e2-abbe-a8a66bc15f66'";
assert!(
is_repairable_bridge_delivery_error(error),
"stale bridge-session bindings should be repaired before retrying delivery"
);
assert!(
is_repairable_bridge_delivery_error("missing event injector capability for member"),
"existing stale event-injector repair path must remain covered"
);
assert!(
is_repairable_bridge_delivery_error(
"mob member rt:us-president:0 missing required capability interaction_event_injector: autonomous member dispatch"
),
"newer autonomous member dispatch wording should repair and retry instead of dropping the event"
);
assert!(
is_repairable_bridge_delivery_error(
"previous member cleanup ambiguous for member rt:deep-investigator:singleton:0"
),
"ambiguous Meerkat respawn cleanup should trigger bridge repair instead of failing delivery"
);
assert!(
!is_repairable_bridge_delivery_error("model provider returned rate limit"),
"ordinary turn failures must not trigger member repair"
);
}
#[test]
fn bridge_delivery_repair_classifies_topology_restore_failure_as_degraded() {
let identity = meerkat_mob::AgentIdentity::from("rt:review:singleton:0");
let receipt = meerkat_mob::MemberRespawnReceipt::new(
identity.clone(),
meerkat_mob::AgentRuntimeId::new(identity, meerkat_mob::ids::Generation::INITIAL),
meerkat_mob::FenceToken::new(1),
meerkat_mob::FenceToken::new(2),
);
let err = MobRespawnError::TopologyRestoreFailed {
receipt,
failed_peer_ids: vec![meerkat_mob::RespawnTopologyPeerId::from(
"initiative:broken",
)],
};
assert_eq!(
classify_member_repair_respawn_failure(&err),
MemberRepairRespawnFailure::DegradedTopologyRestore {
failed_peer_ids: vec!["initiative:broken".to_string()]
},
"failed peer edges should degrade bridge repair instead of bricking delivery"
);
assert!(
matches!(
classify_member_repair_respawn_failure(&MobRespawnError::NoRuntimeControl {
identity: meerkat_mob::AgentIdentity::from("rt:review:singleton:0"),
}),
MemberRepairRespawnFailure::Fatal(_)
),
"ordinary respawn failures must still fail bridge repair"
);
}
#[test]
fn resume_spawn_spec_authors_nothing() {
let runtime_id = AgentRuntimeId::parse("rt:agent:alpha:0").expect("runtime id");
let draft = AgentBuildDraft {
model: None,
system_prompt: Some("explicit customizer prompt".to_string()),
additional_instructions: Vec::new(),
labels: std::collections::BTreeMap::new(),
app_context: None,
external_tools: Vec::new(),
local_external_tools: Default::default(),
provider_params: None,
};
let session_id = meerkat_core::types::SessionId::new();
let spawn =
build_resume_spawn_spec(&runtime_id, &durable_spec(), &draft, None, &session_id)
.expect("resume spawn spec");
assert_eq!(
spawn.system_prompt_override, None,
"resume must not author or re-send prompt configuration; explicit System \
authoring via a turn is the only mid-thread instruction change"
);
match &spawn.launch_mode {
MemberLaunchMode::Resume { bridge_session_id } => {
assert_eq!(bridge_session_id, &session_id);
}
other => panic!("expected Resume launch mode, got {other:?}"),
}
}
#[test]
fn repair_fallback_only_on_missing_durable_snapshot() {
let session_id = meerkat_core::types::SessionId::new();
assert!(durable_snapshot_is_typed_absent(
&meerkat_mob::MobError::SessionUnavailableForResume {
session_id: session_id.clone(),
reason: meerkat_mob::error::SessionResumeUnavailableReason::Absent,
runtime_state: None,
}
));
assert!(!durable_snapshot_is_typed_absent(
&meerkat_mob::MobError::SessionUnavailableForResume {
session_id,
reason: meerkat_mob::error::SessionResumeUnavailableReason::ArchivedNotRevivable,
runtime_state: Some("archived".to_string()),
}
));
let impersonator = meerkat_mob::MobError::Internal(
"missing durable session snapshot for '019e5fc2-dad4-77e2-abbe-a8a66bc15f66'"
.to_string(),
);
assert!(
impersonator
.to_string()
.contains("missing durable session snapshot")
);
assert!(!durable_snapshot_is_typed_absent(&impersonator));
assert!(
matches!(
classify_member_repair_respawn_failure(&MobRespawnError::Mob(impersonator)),
MemberRepairRespawnFailure::Fatal(_)
),
"an Internal impersonator must stay Fatal in the respawn ladder"
);
}
#[test]
fn typed_absent_recovery_spawns_fresh_directly_because_respawn_is_fatal_after_absence() {
let member_id = MobAgentIdentity::from("rt-agent-alpha-0");
let respawn_after_verified_absence =
MobRespawnError::Mob(meerkat_mob::MobError::MemberNotFound(member_id.clone()));
assert!(
matches!(
classify_member_repair_respawn_failure(&respawn_after_verified_absence),
MemberRepairRespawnFailure::Fatal(_)
),
"MemberNotFound stays Fatal in the respawn ladder — the typed-absent arm \
must never fall through to handle.respawn"
);
let mut labels = std::collections::BTreeMap::new();
labels.insert("agent_identity".to_string(), "agent:alpha".to_string());
let spec = fresh_member_spec_from_pre_delivery_entry(
&member_id,
meerkat_mob::ProfileName::from("worker"),
labels.clone(),
);
assert_eq!(
spec.identity.as_str(),
member_id.as_str(),
"typed-absent recovery must rebuild under the SAME identity"
);
assert_eq!(spec.role_name.as_str(), "worker");
assert_eq!(spec.labels, Some(labels));
assert!(
matches!(spec.launch_mode, MemberLaunchMode::Fresh),
"the durable snapshot is typed-absent: the rebuild takes a fresh session, \
never a Resume rebind onto the gone session"
);
let bare = fresh_member_spec_from_pre_delivery_entry(
&member_id,
meerkat_mob::ProfileName::from("worker"),
std::collections::BTreeMap::new(),
);
assert_eq!(bare.labels, None);
}
#[tokio::test]
async fn reset_snapshot_is_abandoned_before_meerkat_retire() {
let member_id = MobAgentIdentity::from("rt-agent-alpha-0");
let session_id = meerkat_core::types::SessionId::new();
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
abandon_then_retire_reset_superseded(
&member_id,
&session_id,
{
let order = Arc::clone(&order);
move || async move {
order.lock().expect("order lock").push("abandon");
Ok(())
}
},
{
let order = Arc::clone(&order);
move || async move {
assert_eq!(
order.lock().expect("order lock").as_slice(),
["abandon"],
"Meerkat retirement must not start until the exact superseded projection is tombstoned"
);
order.lock().expect("order lock").push("retire");
Ok(())
}
},
)
.await
.expect("pre-abandon then retire");
assert_eq!(
order.lock().expect("order lock").as_slice(),
["abandon", "retire"]
);
}
#[tokio::test]
async fn reset_snapshot_cas_failure_stays_visible_and_blocks_retire() {
let member_id = MobAgentIdentity::from("rt-agent-alpha-0");
let session_id = meerkat_core::types::SessionId::new();
let retire_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
let error = abandon_then_retire_reset_superseded(
&member_id,
&session_id,
|| async {
Err(meerkat_store::SessionStoreError::Internal(
"exact snapshot CAS mismatch".to_string(),
))
},
{
let retire_called = Arc::clone(&retire_called);
move || async move {
retire_called.store(true, std::sync::atomic::Ordering::SeqCst);
Ok(())
}
},
)
.await
.expect_err("CAS failure must retain reset cleanup debt");
assert!(
error.to_string().contains("exact snapshot CAS mismatch"),
"the exact abandon failure must remain visible to the cleanup-debt owner: {error}"
);
assert!(
!retire_called.load(std::sync::atomic::Ordering::SeqCst),
"retirement must not acknowledge cleanup after the exact-CAS abandon failed"
);
}
#[test]
fn resume_error_classification_is_typed_first() {
let restore_failed = meerkat_mob::MobError::MemberRestoreFailed {
member_id: meerkat_mob::ids::AgentIdentity::from("agent-alpha"),
session_id: None,
reason: "durable snapshot missing".to_string(),
};
assert_eq!(
classify_resume_error(&restore_failed),
ResumeRejectionKind::MemberRestoreFailed
);
let continuity = meerkat_mob::MobError::Internal(
"session save rejected: incoming transcript is not a continuation of persisted \
revision sha256:d57e07"
.to_string(),
);
assert_eq!(
classify_resume_error(&continuity),
ResumeRejectionKind::TranscriptContinuity
);
let archived = meerkat_mob::MobError::SessionUnavailableForResume {
session_id: meerkat_core::types::SessionId::new(),
reason: meerkat_mob::error::SessionResumeUnavailableReason::ArchivedNotRevivable,
runtime_state: None,
};
assert_eq!(
classify_resume_error(&archived),
ResumeRejectionKind::ArchivedNotRevivable
);
let absent = meerkat_mob::MobError::SessionUnavailableForResume {
session_id: meerkat_core::types::SessionId::new(),
reason: meerkat_mob::error::SessionResumeUnavailableReason::Absent,
runtime_state: None,
};
assert_eq!(classify_resume_error(&absent), ResumeRejectionKind::Other);
let other = meerkat_mob::MobError::WiringError("unrelated".to_string());
assert_eq!(classify_resume_error(&other), ResumeRejectionKind::Other);
}
#[test]
fn admission_budget_defaults_when_unset_or_unparseable() {
for raw in [None, Some(""), Some(" "), Some("soon"), Some("-5")] {
assert_eq!(
parse_bridge_actor_admission_budget(raw),
BRIDGE_ACTOR_ADMISSION_BUDGET,
"unset or unparseable {raw:?} must fall back to the default budget"
);
}
}
#[test]
fn admission_budget_honours_configured_value_within_the_clamp() {
assert_eq!(
parse_bridge_actor_admission_budget(Some("30")),
Duration::from_secs(30)
);
assert_eq!(
parse_bridge_actor_admission_budget(Some(" 30 ")),
Duration::from_secs(30)
);
assert_eq!(
parse_bridge_actor_admission_budget(Some("0")),
Duration::from_secs(1)
);
assert_eq!(
parse_bridge_actor_admission_budget(Some("999999")),
Duration::from_hours(1)
);
}
#[tokio::test]
async fn responsive_round_trip_passes_through_the_bound_unchanged() {
let deadline = ActorAdmissionDeadline::new(Duration::from_mins(10));
let member = MobAgentIdentity::from("rt-agent-alpha-0");
let value = deadline
.bound("test.responsive", &member, async { 7_u32 })
.await
.expect("a ready round trip must pass through the bound untouched");
assert_eq!(value, 7);
assert!(
deadline
.deadline
.saturating_duration_since(tokio::time::Instant::now())
> Duration::from_secs(599)
);
}
#[tokio::test]
async fn stalled_actor_fails_typed_instead_of_hanging() {
let deadline = ActorAdmissionDeadline::new(Duration::from_millis(20));
let member = MobAgentIdentity::from("rt-agent-alpha-0");
let error = deadline
.bound(
"deliver.submit_work",
&member,
std::future::pending::<Result<(), meerkat_mob::MobError>>(),
)
.await
.expect_err("a mob actor that never replies must not hang the delivery");
match error {
BridgeError::ActorAdmissionTimeout {
operation,
identity,
waited,
} => {
assert_eq!(operation, "deliver.submit_work");
assert_eq!(identity.as_str(), "rt-agent-alpha-0");
assert!(waited >= Duration::from_millis(20), "waited {waited:?}");
}
other => panic!("expected a typed admission timeout, got {other:?}"),
}
}
#[tokio::test]
async fn admission_timeout_names_the_operation_and_member_and_is_never_repairable() {
let deadline = ActorAdmissionDeadline::new(Duration::from_millis(5));
let member = MobAgentIdentity::from("rt-agent-alpha-0");
let error = deadline
.bound(
"deliver.get_member",
&member,
std::future::pending::<Result<(), meerkat_mob::MobError>>(),
)
.await
.expect_err("stalled actor");
let rendered = error.to_string();
assert!(rendered.contains("deliver.get_member"), "{rendered}");
assert!(rendered.contains("rt-agent-alpha-0"), "{rendered}");
assert!(
!is_repairable_bridge_delivery_error(&rendered),
"a blocked actor is not stale runtime state: {rendered}"
);
}
#[tokio::test]
async fn serialized_round_trips_share_one_budget() {
let budget = Duration::from_millis(50);
let deadline = ActorAdmissionDeadline::new(budget);
let member = MobAgentIdentity::from("rt-agent-alpha-0");
let mut outcomes = Vec::new();
outcomes.push(
deadline
.bound(
"test.hop",
&member,
std::future::pending::<Result<(), meerkat_mob::MobError>>(),
)
.await,
);
let after_first = tokio::time::Instant::now();
for _ in 0..2 {
outcomes.push(
deadline
.bound(
"test.hop",
&member,
std::future::pending::<Result<(), meerkat_mob::MobError>>(),
)
.await,
);
}
let spent_after_first = after_first.elapsed();
assert!(
outcomes
.iter()
.all(|outcome| matches!(outcome, Err(BridgeError::ActorAdmissionTimeout { .. }))),
"every hop against a dead actor must fail typed"
);
assert!(
spent_after_first < budget,
"the deadline must not be re-armed per hop: two further hops spent \
{spent_after_first:?} against a {budget:?} budget"
);
}
#[test]
fn heal_error_tier_recovery_refused_parks_terminal_unprovable() {
let id = meerkat_core::SessionId::new();
let verdict = map_committed_boundary_recovery_error(
meerkat_core::SessionError::DurableTailRecoveryRefused { id: id.clone() },
);
match verdict {
Ok(CommittedBoundaryRepair::Unprovable { reason }) => {
assert!(
reason.contains(&id.to_string()),
"park reason must carry the session id for the operator: {reason}"
);
assert!(
reason.contains("refused"),
"park reason must state the refusal: {reason}"
);
}
other => panic!("refused must park as Unprovable, got {other:?}"),
}
}
#[test]
fn heal_error_tier_quarantined_evidence_parks_terminal_unprovable() {
let verdict = map_committed_boundary_recovery_error(
meerkat_core::SessionError::DurableEvidenceQuarantined {
id: meerkat_core::SessionId::new(),
},
);
assert!(
matches!(verdict, Ok(CommittedBoundaryRepair::Unprovable { .. })),
"quarantined evidence must park as Unprovable, got {verdict:?}"
);
}
#[test]
fn heal_error_tier_busy_and_held_stay_retryable_errors() {
for error in [
meerkat_core::SessionError::Busy {
id: meerkat_core::SessionId::new(),
},
meerkat_core::SessionError::DurableTailHeldForRecovery {
id: meerkat_core::SessionId::new(),
},
] {
let verdict = map_committed_boundary_recovery_error(error);
assert!(
matches!(verdict, Err(BridgeError::Mob(_))),
"retryable-tier errors must stay bridge errors, got {verdict:?}"
);
}
}
}