use crate::definition::MobDefinition;
use crate::ids::{
AgentIdentity, AgentRuntimeId, FenceToken, FlowId, Generation, MobId, ProfileName, RunId,
StepId,
};
use crate::roster::MobMemberKickoffSnapshot;
use crate::runtime_mode::MobRuntimeMode;
use chrono::{DateTime, Utc};
use meerkat_contracts::wire::supervisor_bridge::BridgeBootstrapToken;
use meerkat_core::comms::{PeerName, TrustedPeerDescriptor};
use meerkat_core::event::{AgentEvent, EventEnvelope};
use meerkat_core::service::{MobToolCallerProvenance, OpaquePrincipalToken};
use meerkat_core::types::SessionId;
use serde::ser::SerializeMap;
use serde::{Deserialize, Serialize};
#[cfg(not(target_arch = "wasm32"))]
use serde_json::Value;
use std::collections::BTreeMap;
#[cfg(not(target_arch = "wasm32"))]
use std::io::{Error as IoError, ErrorKind as IoErrorKind};
#[derive(Debug, Clone, Serialize)]
pub struct MobEvent {
pub cursor: u64,
pub timestamp: DateTime<Utc>,
pub mob_id: MobId,
pub kind: MobEventKind,
}
#[derive(Debug, Clone, Deserialize)]
struct MobEventCanonical {
pub cursor: u64,
pub timestamp: DateTime<Utc>,
pub mob_id: MobId,
pub kind: MobEventKind,
}
impl<'de> Deserialize<'de> for MobEvent {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let canonical = MobEventCanonical::deserialize(deserializer)?;
Ok(Self {
cursor: canonical.cursor,
timestamp: canonical.timestamp,
mob_id: canonical.mob_id,
kind: canonical.kind,
})
}
}
#[derive(Debug, Clone)]
pub struct NewMobEvent {
pub(crate) mob_id: MobId,
pub(crate) timestamp: Option<DateTime<Utc>>,
pub(crate) kind: MobEventKind,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum MemberRef {
Session {
session_id: SessionId,
},
BackendPeer {
peer_id: String,
address: String,
pubkey: [u8; 32],
bootstrap_token: Option<BridgeBootstrapToken>,
session_id: Option<SessionId>,
},
}
impl MemberRef {
pub(crate) fn from_bridge_session_id(session_id: SessionId) -> Self {
Self::Session { session_id }
}
pub(crate) fn bridge_session_id(&self) -> Option<&SessionId> {
match self {
Self::Session { session_id } => Some(session_id),
Self::BackendPeer { session_id, .. } => session_id.as_ref(),
}
}
}
impl Serialize for MemberRef {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
match self {
Self::Session { session_id } => {
let mut map = serializer.serialize_map(Some(2))?;
map.serialize_entry("kind", "session")?;
map.serialize_entry("bridge_session_id", session_id)?;
map.end()
}
Self::BackendPeer {
peer_id,
address,
pubkey,
session_id,
..
} => {
let mut map = serializer.serialize_map(None)?;
map.serialize_entry("kind", "backend_peer")?;
map.serialize_entry("peer_id", peer_id)?;
map.serialize_entry("address", address)?;
map.serialize_entry("pubkey", pubkey)?;
if let Some(session_id) = session_id {
map.serialize_entry("bridge_session_id", session_id)?;
}
map.end()
}
}
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
enum MemberRefDe {
Session {
bridge_session_id: SessionId,
},
BackendPeer {
peer_id: String,
address: String,
pubkey: [u8; 32],
#[serde(default)]
bootstrap_token: Option<BridgeBootstrapToken>,
#[serde(default)]
bridge_session_id: Option<SessionId>,
},
}
impl<'de> Deserialize<'de> for MemberRef {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
Ok(match MemberRefDe::deserialize(deserializer)? {
MemberRefDe::Session { bridge_session_id } => Self::Session {
session_id: bridge_session_id,
},
MemberRefDe::BackendPeer {
peer_id,
address,
pubkey,
bootstrap_token,
bridge_session_id,
} => Self::BackendPeer {
peer_id,
address,
pubkey,
bootstrap_token,
session_id: bridge_session_id,
},
})
}
}
#[cfg(not(target_arch = "wasm32"))]
const CURRENT_STORED_MOB_EVENT_SCHEMA_VERSION: u32 = 8;
#[cfg(not(target_arch = "wasm32"))]
fn stored_mob_event_format_error(message: impl Into<String>) -> serde_json::Error {
serde_json::Error::io(IoError::new(IoErrorKind::InvalidData, message.into()))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowFailureClass {
StepTimeout,
TopologyViolation,
SchemaValidation,
RunCanceled,
SupervisorEscalation,
InsufficientTargets,
StepError,
AdmissionFailed,
AlreadyFailed,
Internal,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum MobEventKind {
MobCreated {
definition: Box<MobDefinition>,
},
MobOwnerBridgeSessionBound {
bridge_session_id: SessionId,
destroy_on_owner_archive: bool,
implicit_delegation_mob: bool,
},
MobCompleted,
MobDestroying,
MobDestroyStorageFinalizing,
MobReset,
MemberSpawned(MemberSpawnedEvent),
MemberRetired {
agent_identity: AgentIdentity,
generation: Generation,
role: ProfileName,
},
MemberReset {
agent_identity: AgentIdentity,
previous_generation: Generation,
new_generation: Generation,
fence_token: FenceToken,
agent_runtime_id: AgentRuntimeId,
},
MemberSessionBindingRecovered(MemberSessionBindingRecoveredEvent),
MemberKickoffUpdated {
member: AgentIdentity,
kickoff: MobMemberKickoffSnapshot,
},
MembersWired {
a: AgentIdentity,
b: AgentIdentity,
},
MembersWiredBatch {
edges: Vec<MemberWireEdge>,
},
MembersUnwired {
a: AgentIdentity,
b: AgentIdentity,
},
ExternalPeerWired {
local: AgentIdentity,
spec: TrustedPeerDescriptor,
},
ExternalPeerUnwired {
local: AgentIdentity,
peer_name: PeerName,
},
FlowStarted {
run_id: RunId,
flow_id: FlowId,
params: serde_json::Value,
},
FlowCompleted {
run_id: RunId,
flow_id: FlowId,
#[serde(default, skip_serializing_if = "Option::is_none")]
structured_output: Option<serde_json::Value>,
},
FlowFailed {
run_id: RunId,
flow_id: FlowId,
cause: FlowFailureClass,
reason: String,
},
FlowCanceled { run_id: RunId, flow_id: FlowId },
StepDispatched {
run_id: RunId,
step_id: StepId,
target: AgentRuntimeId,
},
StepTargetCompleted {
run_id: RunId,
step_id: StepId,
target: AgentRuntimeId,
},
StepTargetFailed {
run_id: RunId,
step_id: StepId,
target: AgentRuntimeId,
reason: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
error_report: Option<meerkat_core::event::AgentErrorReport>,
#[serde(default, skip_serializing_if = "Option::is_none")]
error: Option<meerkat_core::event::TurnErrorMetadata>,
},
StepCompleted { run_id: RunId, step_id: StepId },
StepFailed {
run_id: RunId,
step_id: StepId,
reason: String,
},
StepSkipped {
run_id: RunId,
step_id: StepId,
reason: String,
},
TopologyViolation {
from_role: ProfileName,
to_role: ProfileName,
},
SupervisorEscalation {
run_id: RunId,
step_id: StepId,
escalated_to: AgentIdentity,
},
OperatorActionRecorded {
tool_name: String,
principal_token: OpaquePrincipalToken,
#[serde(default, skip_serializing_if = "Option::is_none")]
caller_provenance: Option<MobToolCallerProvenance>,
#[serde(default, skip_serializing_if = "Option::is_none")]
audit_invocation_id: Option<String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct MemberWireEdge {
pub a: AgentIdentity,
pub b: AgentIdentity,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AttributedEvent {
pub source: AgentRuntimeId,
pub source_fence_token: FenceToken,
pub role: ProfileName,
pub envelope: EventEnvelope<AgentEvent>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct MemberSpawnedEvent {
pub agent_identity: AgentIdentity,
pub generation: Generation,
pub fence_token: FenceToken,
pub agent_runtime_id: AgentRuntimeId,
pub role: ProfileName,
pub runtime_mode: MobRuntimeMode,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub labels: BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub effective_profile_override: Option<crate::profile::Profile>,
#[serde(default, skip_serializing_if = "crate::event::is_ephemeral_continuity")]
pub continuity_intent: crate::runtime::SpawnContinuityIntent,
#[serde(skip, default)]
pub(crate) bridge_member_ref: Option<MemberRef>,
}
impl MemberSpawnedEvent {
pub fn new(
agent_identity: AgentIdentity,
generation: Generation,
fence_token: FenceToken,
agent_runtime_id: AgentRuntimeId,
role: ProfileName,
) -> Self {
Self {
agent_identity,
generation,
fence_token,
agent_runtime_id,
role,
runtime_mode: MobRuntimeMode::AutonomousHost,
labels: BTreeMap::new(),
effective_profile_override: None,
continuity_intent: crate::runtime::SpawnContinuityIntent::Ephemeral,
bridge_member_ref: None,
}
}
pub(crate) fn with_bridge_member_ref(mut self, bridge_member_ref: Option<MemberRef>) -> Self {
self.bridge_member_ref = bridge_member_ref;
self
}
#[cfg(any(not(target_arch = "wasm32"), test))]
pub(crate) fn bridge_member_ref(&self) -> Option<&MemberRef> {
self.bridge_member_ref.as_ref()
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct MemberSessionBindingRecoveredEvent {
pub agent_identity: AgentIdentity,
pub agent_runtime_id: AgentRuntimeId,
#[serde(skip, default)]
pub(crate) bridge_session_id: Option<SessionId>,
}
impl MemberSessionBindingRecoveredEvent {
pub(crate) fn new(
agent_identity: AgentIdentity,
agent_runtime_id: AgentRuntimeId,
bridge_session_id: SessionId,
) -> Self {
Self {
agent_identity,
agent_runtime_id,
bridge_session_id: Some(bridge_session_id),
}
}
pub(crate) fn bridge_session_id(&self) -> Option<&SessionId> {
self.bridge_session_id.as_ref()
}
}
fn is_ephemeral_continuity(intent: &crate::runtime::SpawnContinuityIntent) -> bool {
matches!(intent, crate::runtime::SpawnContinuityIntent::Ephemeral)
}
#[cfg(any(not(target_arch = "wasm32"), test))]
impl MobEventKind {
pub(crate) fn member_spawned(&self) -> Option<&MemberSpawnedEvent> {
match self {
Self::MemberSpawned(event) => Some(event),
_ => None,
}
}
pub(crate) fn member_spawned_mut(&mut self) -> Option<&mut MemberSpawnedEvent> {
match self {
Self::MemberSpawned(event) => Some(event),
_ => None,
}
}
pub(crate) fn member_session_binding_recovered(
&self,
) -> Option<&MemberSessionBindingRecoveredEvent> {
match self {
Self::MemberSessionBindingRecovered(event) => Some(event),
_ => None,
}
}
pub(crate) fn member_session_binding_recovered_mut(
&mut self,
) -> Option<&mut MemberSessionBindingRecoveredEvent> {
match self {
Self::MemberSessionBindingRecovered(event) => Some(event),
_ => None,
}
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn encode_stored_mob_event(event: &MobEvent) -> Result<Vec<u8>, serde_json::Error> {
let mut value = serde_json::to_value(event)?;
if let Some(member_spawned) = event.kind.member_spawned()
&& let Some(bridge_member_ref) = member_spawned.bridge_member_ref()
&& let Some(kind) = value.get_mut("kind").and_then(Value::as_object_mut)
{
kind.insert(
"bridge_member_ref".to_string(),
serde_json::to_value(bridge_member_ref)?,
);
}
if let Some(recovered) = event.kind.member_session_binding_recovered()
&& let Some(bridge_session_id) = recovered.bridge_session_id()
&& let Some(kind) = value.get_mut("kind").and_then(Value::as_object_mut)
{
kind.insert(
"bridge_session_id".to_string(),
serde_json::to_value(bridge_session_id)?,
);
}
serde_json::to_vec(&serde_json::json!({
"schema_version": CURRENT_STORED_MOB_EVENT_SCHEMA_VERSION,
"event": value,
}))
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn decode_stored_mob_event(bytes: &[u8]) -> Result<MobEvent, serde_json::Error> {
let mut encoded: Value = serde_json::from_slice(bytes)?;
let encoded_object = encoded.as_object_mut().ok_or_else(|| {
stored_mob_event_format_error("stored mob event envelope must be an object")
})?;
let schema_version = encoded_object
.remove("schema_version")
.and_then(|value| value.as_u64())
.ok_or_else(|| {
stored_mob_event_format_error(
"stored mob event missing schema_version; pre-0.6 mob event history is unsupported",
)
})?;
if schema_version != CURRENT_STORED_MOB_EVENT_SCHEMA_VERSION as u64 {
return Err(stored_mob_event_format_error(format!(
"unsupported stored mob event schema_version={schema_version}; expected {CURRENT_STORED_MOB_EVENT_SCHEMA_VERSION}",
)));
}
let mut value = encoded_object
.remove("event")
.ok_or_else(|| stored_mob_event_format_error("stored mob event missing event payload"))?;
let bridge_member_ref = value
.get_mut("kind")
.and_then(Value::as_object_mut)
.and_then(|kind| kind.remove("bridge_member_ref"))
.map(serde_json::from_value)
.transpose()?;
let recovered_bridge_session_id = value
.get_mut("kind")
.and_then(Value::as_object_mut)
.and_then(|kind| {
if kind.get("type").and_then(Value::as_str) == Some("member_session_binding_recovered")
{
kind.remove("bridge_session_id")
} else {
None
}
})
.map(serde_json::from_value)
.transpose()?;
let mut event: MobEvent = serde_json::from_value(value)?;
if let Some(bridge_member_ref) = bridge_member_ref
&& let Some(member_spawned) = event.kind.member_spawned_mut()
{
member_spawned.bridge_member_ref = Some(bridge_member_ref);
}
if let Some(bridge_session_id) = recovered_bridge_session_id
&& let Some(recovered) = event.kind.member_session_binding_recovered_mut()
{
recovered.bridge_session_id = Some(bridge_session_id);
}
Ok(event)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::definition::MobDefinition;
use crate::profile::{Profile, ProfileBinding, ToolConfig};
use serde_json::json;
use std::collections::BTreeMap;
use uuid::Uuid;
fn sample_definition() -> MobDefinition {
let mut definition = MobDefinition::explicit("test-mob");
definition.profiles.insert(
ProfileName::from("worker"),
ProfileBinding::Inline(Box::new(Profile {
model: "claude-sonnet-4-5".to_string(),
provider: None,
self_hosted_server_id: None,
image_generation_provider: None,
auto_compact_threshold: None,
resume_overrides: Vec::new(),
skills: vec![],
tools: ToolConfig::default(),
peer_description: "A worker".to_string(),
external_addressable: false,
backend: None,
runtime_mode: MobRuntimeMode::AutonomousHost,
max_inline_peer_notifications: None,
output_schema: None,
provider_params: None,
})),
);
definition
}
fn roundtrip(kind: &MobEventKind) {
let json = serde_json::to_string(kind).unwrap();
let parsed: MobEventKind = serde_json::from_str(&json).unwrap();
assert_eq!(&parsed, kind);
}
#[test]
fn test_mob_created_roundtrip() {
roundtrip(&MobEventKind::MobCreated {
definition: Box::new(sample_definition()),
});
}
#[test]
fn test_mob_owner_bridge_session_bound_roundtrip() {
roundtrip(&MobEventKind::MobOwnerBridgeSessionBound {
bridge_session_id: SessionId::from_uuid(Uuid::nil()),
destroy_on_owner_archive: true,
implicit_delegation_mob: true,
});
}
#[test]
fn test_stored_mob_event_roundtrip_preserves_owner_bridge_session_bound() {
let sid = SessionId::from_uuid(Uuid::nil());
let event = MobEvent {
cursor: 1,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind: MobEventKind::MobOwnerBridgeSessionBound {
bridge_session_id: sid.clone(),
destroy_on_owner_archive: true,
implicit_delegation_mob: true,
},
};
let encoded = encode_stored_mob_event(&event).unwrap();
let decoded = decode_stored_mob_event(&encoded).unwrap();
match decoded.kind {
MobEventKind::MobOwnerBridgeSessionBound {
bridge_session_id,
destroy_on_owner_archive,
implicit_delegation_mob,
} => {
assert_eq!(bridge_session_id, sid);
assert!(destroy_on_owner_archive);
assert!(implicit_delegation_mob);
}
other => panic!("expected MobOwnerBridgeSessionBound, got {other:?}"),
}
}
#[test]
fn test_mob_completed_roundtrip() {
roundtrip(&MobEventKind::MobCompleted);
}
#[test]
fn test_mob_destroying_roundtrip() {
roundtrip(&MobEventKind::MobDestroying);
}
#[test]
fn test_mob_destroy_storage_finalizing_roundtrip() {
roundtrip(&MobEventKind::MobDestroyStorageFinalizing);
}
#[test]
fn test_mob_reset_roundtrip() {
roundtrip(&MobEventKind::MobReset);
}
#[test]
fn test_flow_variants_roundtrip() {
let run_id = RunId::new();
let flow_id = FlowId::from("flow-a");
let step_id = StepId::from("step-a");
let runtime_id = AgentRuntimeId::initial(AgentIdentity::from("worker-1"));
let escalated_identity = AgentIdentity::from("worker-1");
roundtrip(&MobEventKind::FlowStarted {
run_id: run_id.clone(),
flow_id: flow_id.clone(),
params: serde_json::json!({"k":"v"}),
});
roundtrip(&MobEventKind::FlowCompleted {
run_id: run_id.clone(),
flow_id: flow_id.clone(),
structured_output: Some(serde_json::json!({
"steps": {
"step-a": {
"ok": true
}
}
})),
});
roundtrip(&MobEventKind::FlowFailed {
run_id: run_id.clone(),
flow_id: flow_id.clone(),
cause: FlowFailureClass::StepError,
reason: "boom".to_string(),
});
roundtrip(&MobEventKind::FlowCanceled {
run_id: run_id.clone(),
flow_id: flow_id.clone(),
});
roundtrip(&MobEventKind::StepDispatched {
run_id: run_id.clone(),
step_id: step_id.clone(),
target: runtime_id.clone(),
});
roundtrip(&MobEventKind::StepTargetCompleted {
run_id: run_id.clone(),
step_id: step_id.clone(),
target: runtime_id.clone(),
});
roundtrip(&MobEventKind::StepTargetFailed {
run_id: run_id.clone(),
step_id: step_id.clone(),
target: runtime_id.clone(),
reason: "fail".to_string(),
error_report: None,
error: None,
});
roundtrip(&MobEventKind::StepCompleted {
run_id: run_id.clone(),
step_id: step_id.clone(),
});
roundtrip(&MobEventKind::StepFailed {
run_id: run_id.clone(),
step_id: step_id.clone(),
reason: "timeout after 1000ms".to_string(),
});
roundtrip(&MobEventKind::StepSkipped {
run_id: run_id.clone(),
step_id: step_id.clone(),
reason: "branch lost".to_string(),
});
roundtrip(&MobEventKind::TopologyViolation {
from_role: ProfileName::from("lead"),
to_role: ProfileName::from("worker"),
});
roundtrip(&MobEventKind::SupervisorEscalation {
run_id,
step_id,
escalated_to: escalated_identity,
});
}
#[test]
fn test_members_wired_roundtrip() {
roundtrip(&MobEventKind::MembersWired {
a: AgentIdentity::from("l-1"),
b: AgentIdentity::from("w-2"),
});
}
#[test]
fn test_members_wired_batch_roundtrip() {
roundtrip(&MobEventKind::MembersWiredBatch {
edges: vec![
MemberWireEdge {
a: AgentIdentity::from("l-1"),
b: AgentIdentity::from("w-2"),
},
MemberWireEdge {
a: AgentIdentity::from("l-1"),
b: AgentIdentity::from("w-3"),
},
],
});
}
#[test]
fn test_members_unwired_roundtrip() {
roundtrip(&MobEventKind::MembersUnwired {
a: AgentIdentity::from("l-1"),
b: AgentIdentity::from("w-2"),
});
}
#[test]
fn test_external_peer_wired_roundtrip() {
let pubkey = [8u8; 32];
let peer_id = meerkat_core::comms::PeerId::from_ed25519_pubkey(&pubkey);
let spec = TrustedPeerDescriptor::unsigned_with_pubkey(
"remote-mob/worker/agent-b",
peer_id.to_string(),
pubkey,
"inproc://remote-mob/worker/agent-b",
)
.expect("valid external peer");
roundtrip(&MobEventKind::ExternalPeerWired {
local: AgentIdentity::from("l-1"),
spec,
});
}
#[test]
fn test_external_peer_unwired_roundtrip() {
roundtrip(&MobEventKind::ExternalPeerUnwired {
local: AgentIdentity::from("l-1"),
peer_name: PeerName::new("remote-mob/worker/agent-b").unwrap(),
});
}
#[test]
fn test_operator_action_recorded_roundtrip() {
roundtrip(&MobEventKind::OperatorActionRecorded {
tool_name: "mob_create".to_string(),
principal_token: OpaquePrincipalToken::new("opaque-principal"),
caller_provenance: Some(
MobToolCallerProvenance::new()
.with_session_id(SessionId::from_uuid(Uuid::nil()))
.with_mob_id("test-mob")
.with_member_id("lead-1"),
),
audit_invocation_id: Some("audit-123".to_string()),
});
}
#[test]
fn test_mob_event_full_roundtrip() {
let event = MobEvent {
cursor: 42,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind: MobEventKind::MobCompleted,
};
let json = serde_json::to_string(&event).unwrap();
let parsed: MobEvent = serde_json::from_str(&json).unwrap();
assert_eq!(parsed.cursor, 42);
assert_eq!(parsed.mob_id.as_str(), "test-mob");
}
#[test]
fn test_member_ref_rejects_legacy_session_only_payload() {
let sid = SessionId::from_uuid(Uuid::nil());
let parsed = serde_json::from_value::<MemberRef>(json!({
"session_id": sid,
}));
assert!(parsed.is_err());
let parsed = serde_json::from_value::<MemberRef>(json!({
"kind": "session",
"session_id": sid,
}));
assert!(parsed.is_err());
}
#[test]
fn test_member_ref_serializes_deterministically() {
let sid = SessionId::from_uuid(Uuid::nil());
let member_ref = MemberRef::from_bridge_session_id(sid);
let first = serde_json::to_string(&member_ref).unwrap();
let second = serde_json::to_string(&member_ref).unwrap();
assert_eq!(first, second);
assert_eq!(
first,
r#"{"kind":"session","bridge_session_id":"00000000-0000-0000-0000-000000000000"}"#
);
}
#[test]
fn test_member_ref_deserializes_bridge_session_only_payload() {
let sid = SessionId::from_uuid(Uuid::nil());
let parsed: MemberRef = serde_json::from_value(json!({
"kind": "session",
"bridge_session_id": sid,
}))
.unwrap();
assert_eq!(parsed, MemberRef::from_bridge_session_id(sid));
}
#[test]
fn test_member_ref_roundtrip_backend_peer() {
let member_ref = MemberRef::BackendPeer {
peer_id: "peer-123".to_string(),
address: "https://backend.example/peers/peer-123".to_string(),
pubkey: [9u8; 32],
bootstrap_token: None,
session_id: Some(SessionId::from_uuid(Uuid::nil())),
};
let json = serde_json::to_string(&member_ref).unwrap();
let value: serde_json::Value = serde_json::from_str(&json).unwrap();
assert_eq!(
value["bridge_session_id"],
"00000000-0000-0000-0000-000000000000"
);
let parsed: MemberRef = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, member_ref);
}
#[test]
fn test_member_ref_backend_peer_omits_bootstrap_token_in_serialized_output() {
let member_ref = MemberRef::BackendPeer {
peer_id: "peer-123".to_string(),
address: "https://backend.example/peers/peer-123".to_string(),
pubkey: [9u8; 32],
bootstrap_token: Some("secret-bootstrap-proof".into()),
session_id: None,
};
let value = serde_json::to_value(&member_ref).unwrap();
assert_eq!(value["kind"], "backend_peer");
assert_eq!(value["peer_id"], "peer-123");
assert_eq!(value["address"], "https://backend.example/peers/peer-123");
assert!(
value.get("bootstrap_token").is_none(),
"bootstrap proof must not be exposed through serialized member refs"
);
}
#[test]
fn test_stored_mob_event_roundtrip_preserves_bridge_member_ref() {
let sid = SessionId::from_uuid(Uuid::nil());
let identity = AgentIdentity::from("researcher");
let event = MobEvent {
cursor: 1,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind: MobEventKind::MemberSpawned(
MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(1),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
)
.with_bridge_member_ref(Some(MemberRef::from_bridge_session_id(sid.clone()))),
),
};
let encoded = encode_stored_mob_event(&event).unwrap();
let decoded = decode_stored_mob_event(&encoded).unwrap();
match decoded.kind {
MobEventKind::MemberSpawned(member_spawned) => {
assert_eq!(
member_spawned
.bridge_member_ref()
.and_then(MemberRef::bridge_session_id),
Some(&sid)
);
}
other => panic!("expected MemberSpawned, got {other:?}"),
}
}
#[test]
fn test_stored_mob_event_roundtrip_preserves_recovered_member_session_binding() {
let sid = SessionId::from_uuid(Uuid::nil());
let identity = AgentIdentity::from("researcher");
let event = MobEvent {
cursor: 1,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind: MobEventKind::MemberSessionBindingRecovered(
MemberSessionBindingRecoveredEvent::new(
identity.clone(),
AgentRuntimeId::initial(identity),
sid.clone(),
),
),
};
let encoded = encode_stored_mob_event(&event).unwrap();
let decoded = decode_stored_mob_event(&encoded).unwrap();
match decoded.kind {
MobEventKind::MemberSessionBindingRecovered(recovered) => {
assert_eq!(recovered.bridge_session_id(), Some(&sid));
}
other => panic!("expected MemberSessionBindingRecovered, got {other:?}"),
}
}
#[test]
fn member_spawned_public_wire_shape_excludes_bridge_member_ref() {
let identity = AgentIdentity::from("researcher");
let sid = SessionId::from_uuid(Uuid::nil());
let kind = MobEventKind::MemberSpawned(
MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(1),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
)
.with_bridge_member_ref(Some(MemberRef::from_bridge_session_id(sid))),
);
let value = serde_json::to_value(&kind).expect("serialize mob event kind");
let object = value
.as_object()
.expect("MemberSpawned serializes as an object");
let serialized = serde_json::to_string(&value).unwrap();
assert!(
!serialized.contains("bridge_member_ref"),
"public MemberSpawned wire shape must never expose bridge_member_ref; got: {serialized}",
);
let payload_carrier = object.values().find(|v| v.is_object()).unwrap_or(&value);
let payload_object = payload_carrier.as_object().unwrap_or(object);
assert!(
payload_object.contains_key("agent_identity") || serialized.contains("agent_identity"),
"public MemberSpawned must carry agent_identity: {serialized}",
);
}
#[test]
fn member_session_binding_recovered_public_wire_shape_excludes_bridge_session_id() {
let identity = AgentIdentity::from("researcher");
let sid = SessionId::from_uuid(Uuid::nil());
let kind =
MobEventKind::MemberSessionBindingRecovered(MemberSessionBindingRecoveredEvent::new(
identity.clone(),
AgentRuntimeId::initial(identity),
sid,
));
let value = serde_json::to_value(&kind).expect("serialize mob event kind");
let serialized = serde_json::to_string(&value).unwrap();
assert!(
!serialized.contains("bridge_session_id"),
"public MemberSessionBindingRecovered must not expose bridge_session_id: {serialized}",
);
assert!(
serialized.contains("agent_identity"),
"public MemberSessionBindingRecovered must carry agent_identity: {serialized}",
);
}
#[test]
fn test_decode_stored_mob_event_rejects_unversioned_payload() {
let raw = serde_json::to_vec(&MobEvent {
cursor: 1,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind: MobEventKind::MobCompleted,
})
.unwrap();
let error =
decode_stored_mob_event(&raw).expect_err("unversioned payload must be rejected");
assert!(
error
.to_string()
.contains("pre-0.6 mob event history is unsupported")
);
}
#[test]
fn test_decode_stored_mob_event_rejects_unsupported_schema_version() {
let raw = serde_json::to_vec(&json!({
"schema_version": CURRENT_STORED_MOB_EVENT_SCHEMA_VERSION + 1,
"event": {
"cursor": 1,
"timestamp": "2026-02-19T00:00:00Z",
"mob_id": "test-mob",
"kind": {
"type": "mob_completed"
}
}
}))
.unwrap();
let error =
decode_stored_mob_event(&raw).expect_err("unsupported schema version must be rejected");
assert!(
error
.to_string()
.contains("unsupported stored mob event schema_version")
);
}
#[test]
fn test_member_spawned_roundtrip() {
let identity = AgentIdentity::from("researcher");
roundtrip(&MobEventKind::MemberSpawned(MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(1),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
)));
}
#[test]
fn test_member_spawned_with_labels_roundtrip() {
let identity = AgentIdentity::from("coder");
let mut labels = BTreeMap::new();
labels.insert("team".to_string(), "backend".to_string());
let mut event = MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(42),
AgentRuntimeId::initial(identity),
ProfileName::from("coder"),
);
event.runtime_mode = MobRuntimeMode::TurnDriven;
event.labels = labels;
roundtrip(&MobEventKind::MemberSpawned(event));
}
fn override_profile_with_mcp_servers() -> Profile {
Profile {
model: "claude-sonnet-4-5".to_string(),
provider: None,
self_hosted_server_id: None,
image_generation_provider: None,
auto_compact_threshold: None,
resume_overrides: Vec::new(),
skills: vec![],
tools: ToolConfig {
mcp_servers: vec![meerkat_core::mcp_config::McpServerConfig::stdio(
"planner",
"/bin/echo",
vec![],
std::collections::HashMap::new(),
)],
..ToolConfig::default()
},
peer_description: "A worker".to_string(),
external_addressable: false,
backend: None,
runtime_mode: MobRuntimeMode::AutonomousHost,
max_inline_peer_notifications: None,
output_schema: None,
provider_params: None,
}
}
#[test]
fn test_member_spawned_effective_profile_override_roundtrip() {
let identity = AgentIdentity::from("coder");
let mut event = MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(7),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
);
event.effective_profile_override = Some(override_profile_with_mcp_servers());
roundtrip(&MobEventKind::MemberSpawned(event));
}
#[test]
fn test_member_spawned_without_override_serializes_without_the_field() {
let identity = AgentIdentity::from("coder");
let event = MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(7),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
);
let serialized =
serde_json::to_string(&MobEventKind::MemberSpawned(event)).expect("serialize");
assert!(
!serialized.contains("effective_profile_override"),
"None override must keep the wire shape byte-identical to pre-wave-2 events: {serialized}"
);
}
#[test]
fn test_member_spawned_pre_wave2_event_without_override_deserializes_to_none() {
let event: MobEvent = serde_json::from_value(json!({
"cursor": 1,
"timestamp": "2026-02-19T00:00:00Z",
"mob_id": "test-mob",
"kind": {
"type": "member_spawned",
"agent_identity": "researcher",
"generation": 0,
"fence_token": 1,
"agent_runtime_id": {
"identity": "researcher",
"generation": 0
},
"role": "worker",
"runtime_mode": "autonomous_host"
},
}))
.expect("pre-wave-2 stored event must stay readable");
match event.kind {
MobEventKind::MemberSpawned(member_spawned) => {
assert_eq!(member_spawned.effective_profile_override, None);
}
other => panic!("expected MemberSpawned, got {other:?}"),
}
}
#[test]
fn test_member_spawned_rejects_missing_runtime_mode() {
let result = serde_json::from_value::<MobEvent>(json!({
"cursor": 1,
"timestamp": "2026-02-19T00:00:00Z",
"mob_id": "test-mob",
"kind": {
"type": "member_spawned",
"agent_identity": "researcher",
"generation": 0,
"fence_token": 1,
"agent_runtime_id": {
"identity": "researcher",
"generation": 0
},
"role": "worker"
},
}));
assert!(
result.is_err(),
"member_spawned without runtime_mode must be rejected"
);
}
#[test]
fn test_member_retired_roundtrip() {
roundtrip(&MobEventKind::MemberRetired {
agent_identity: AgentIdentity::from("researcher"),
generation: Generation::new(2),
role: ProfileName::from("worker"),
});
}
#[test]
fn test_member_reset_roundtrip() {
let identity = AgentIdentity::from("worker-1");
roundtrip(&MobEventKind::MemberReset {
agent_identity: identity.clone(),
previous_generation: Generation::new(0),
new_generation: Generation::new(1),
fence_token: FenceToken::new(2),
agent_runtime_id: AgentRuntimeId::new(identity, Generation::new(1)),
});
}
#[test]
fn test_member_reset_generation_advancement() {
let identity = AgentIdentity::from("agent-x");
let event = MobEventKind::MemberReset {
agent_identity: identity.clone(),
previous_generation: Generation::new(3),
new_generation: Generation::new(4),
fence_token: FenceToken::new(99),
agent_runtime_id: AgentRuntimeId::new(identity, Generation::new(4)),
};
let json = serde_json::to_string(&event).unwrap();
let parsed: MobEventKind = serde_json::from_str(&json).unwrap();
if let MobEventKind::MemberReset {
previous_generation,
new_generation,
..
} = parsed
{
assert!(new_generation > previous_generation);
} else {
panic!("expected MemberReset");
}
}
}