use std::collections::{BTreeMap, HashMap};
use std::future::Future;
use std::sync::Arc;
use async_trait::async_trait;
use meerkat_core::types::HandlingMode;
use meerkat_mob::ids::AgentIdentity as MobAgentIdentity;
use meerkat_mob::launch::MemberLaunchMode;
use meerkat_mob::{
MobHandle, MobSessionService, 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 is_missing_durable_session_snapshot_error(error: &str) -> bool {
error.contains("missing durable session snapshot")
}
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())
}
#[derive(Debug)]
pub enum BridgeError {
Mob(String),
InvalidInput(String),
ResumeRejected {
kind: ResumeRejectionKind,
detail: String,
},
}
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"
),
}
}
}
impl std::error::Error for BridgeError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ResumeRejectionKind {
MemberRestoreFailed,
TranscriptContinuity,
Other,
}
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::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 },
}
#[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),
}
}
}
async fn submit_internal_bridge_work(
handle: &MobHandle,
member_id: &MobAgentIdentity,
content: &meerkat_core::ContentInput,
injected_context: &[meerkat_core::ContentInput],
handling_mode: HandlingMode,
interaction_id: Option<&str>,
) -> Result<(), BridgeError> {
let entry = 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 mut spec = WorkSpec::new(content.clone(), WorkOrigin::Internal);
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"
);
}
}
}
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()))
}
#[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(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
) -> Result<meerkat_core::types::SessionId, BridgeError>;
async fn deliver_with_mode(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
handling_mode: HandlingMode,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let _ = handling_mode;
self.deliver(runtime_id, content).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 _ = injected_context;
let _ = interaction_id;
self.deliver_with_mode(runtime_id, content, handling_mode)
.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(())
}
}
#[derive(Debug, Clone)]
pub struct MemberInspection {
pub output_preview: Option<String>,
pub is_final: bool,
pub peer_reachable_count: usize,
}
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>,
}
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(),
}
}
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(),
}
}
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(),
}
}
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(),
}
}
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(),
}
}
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(())
}
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 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 resolve_runtime_session_id(
&self,
runtime_id: &AgentRuntimeId,
member_id: &MobAgentIdentity,
missing_message: &'static str,
) -> 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()
&& 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(),
) {
match self
.resume_repair_member(runtime_id, member_id, role, labels, &session_id)
.await
{
Ok(()) => 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 respawn \
(no transcript left to preserve)"
);
}
Err(RepairResumeFailure::Rejected(err)) => 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
{
let mut spec = SpawnMemberSpec::new(role, member_id.clone());
if !labels.is_empty() {
spec = spec.with_labels(labels);
}
self.handle
.ensure_member(spec)
.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 is_missing_durable_session_snapshot_error(&err.to_string()) => {
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())
}
}
pub(crate) fn build_spawn_spec(
runtime_id: &AgentRuntimeId,
spec: &DurableAgentSpec,
draft: &AgentBuildDraft,
base_profile: Option<&meerkat_mob::Profile>,
) -> SpawnMemberSpec {
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 = None;
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);
}
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,
) -> SpawnMemberSpec {
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;
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 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")
.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> {
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 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 {
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;
Ok(ResumeSessionOutcome::Resumed {
session_id: session_id.clone(),
})
}
Err(error) => {
self.verify_durable_session_after_rejected_resume(identity, session_id)
.await;
Err(resume_rejected(
identity,
session_id,
&error,
"resume spawn",
))
}
}
}
async fn deliver(
&self,
runtime_id: &AgentRuntimeId,
content: &meerkat_core::ContentInput,
) -> Result<meerkat_core::types::SessionId, BridgeError> {
let mid = self.member_id_for_runtime_id(runtime_id).await;
let member_entry_before_delivery = self
.handle
.get_member(&mid)
.await
.ok()
.flatten()
.map(|entry| (entry.role, entry.labels));
if content_input_has_images(content) {
let member_entry = 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 = 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,
content,
&[],
HandlingMode::Queue,
None,
)
.await
{
Ok(()) => {}
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?;
submit_internal_bridge_work(
&self.handle,
&mid,
content,
&[],
HandlingMode::Queue,
None,
)
.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",
)
.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_with_mode_and_context(runtime_id, content, &[], handling_mode, None)
.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 mid = self.member_id_for_runtime_id(runtime_id).await;
let member_entry_before_delivery = self
.handle
.get_member(&mid)
.await
.ok()
.flatten()
.map(|entry| (entry.role, entry.labels));
if content_input_has_images(content) {
let member_entry = 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 = 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,
content,
injected_context,
handling_mode,
interaction_id,
)
.await
{
Ok(()) => {}
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?;
submit_internal_bridge_work(
&self.handle,
&mid,
content,
injected_context,
handling_mode,
interaction_id,
)
.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",
)
.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;
use async_trait::async_trait;
use meerkat_core::agent::AgentToolDispatcher;
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(),
})
}
}
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)),
};
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));
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)),
};
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));
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(),
"stale provider pin must be cleared for catalog re-inference"
);
}
#[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)),
};
let spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, None);
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(),
};
let spawn = build_spawn_spec(&runtime_id, &spec, &draft, None);
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(),
};
let session_spawn = build_spawn_spec(&runtime_id, &durable_spec(), &draft, None);
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);
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_inherits_persisted_system_prompt() {
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(),
};
let session_id = meerkat_core::types::SessionId::new();
let spawn =
build_resume_spawn_spec(&runtime_id, &durable_spec(), &draft, None, &session_id);
assert_eq!(
spawn.system_prompt_override, None,
"resume must inherit the persisted System message, never re-send the base prompt"
);
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() {
assert!(is_missing_durable_session_snapshot_error(
"missing durable session snapshot for '019e5fc2-dad4-77e2-abbe-a8a66bc15f66'"
));
assert!(!is_missing_durable_session_snapshot_error(
"missing bridge session snapshot for '019e5fc2-dad4-77e2-abbe-a8a66bc15f66'"
));
assert!(!is_missing_durable_session_snapshot_error(
"session save rejected: incoming transcript is not a continuation of persisted revision"
));
assert!(!is_missing_durable_session_snapshot_error(
"model provider returned rate limit"
));
}
#[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 other = meerkat_mob::MobError::WiringError("unrelated".to_string());
assert_eq!(classify_resume_error(&other), ResumeRejectionKind::Other);
}
}