use std::collections::BTreeMap;
use std::future::Future;
use std::sync::Arc;
use meerkat_mob::MobError;
use meerkat_mob::event::{AttributedEvent, MobEvent, MobEventKind};
use meerkat_mob::ids::{AgentRuntimeId, FenceToken};
use meerkat_mob::runtime::MobEventsView;
use serde_json::Value;
use tokio::sync::broadcast;
use crate::runtime::{MetadataScope, RuntimeMetadataTable};
use crate::types::MobStructuralEventEnvelope;
use crate::unified_runtime::EventQuery;
pub(crate) const QUERY_BATCH_SIZE: usize = 128;
pub(crate) const DEFAULT_QUERY_LIMIT: usize = 256;
const MOB_EVENTS_CHANNEL_CAP: usize = 512;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct MobEventMemberAuthority {
pub(crate) runtime_id: AgentRuntimeId,
pub(crate) fence_token: FenceToken,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct ProjectedMobEvent {
pub(crate) envelope: MobStructuralEventEnvelope,
pub(crate) member_authority: Option<MobEventMemberAuthority>,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct MobEventsQueryPage<T> {
pub(crate) items: Vec<T>,
pub(crate) resume_after_seq: u64,
}
#[derive(Clone)]
pub struct MobEventsStore {
event_tx: broadcast::Sender<MobStructuralEventEnvelope>,
metadata_table: Option<Arc<RuntimeMetadataTable>>,
}
impl Default for MobEventsStore {
fn default() -> Self {
Self::new()
}
}
impl MobEventsStore {
pub fn new() -> Self {
let (event_tx, _) = broadcast::channel(MOB_EVENTS_CHANNEL_CAP);
Self {
event_tx,
metadata_table: None,
}
}
#[must_use]
pub fn with_metadata_table(mut self, table: Arc<RuntimeMetadataTable>) -> Self {
self.metadata_table = Some(table);
self
}
pub fn subscribe(&self) -> broadcast::Receiver<MobStructuralEventEnvelope> {
self.event_tx.subscribe()
}
pub async fn project_attributed_event(
&self,
_event: &AttributedEvent,
) -> Option<MobStructuralEventEnvelope> {
None
}
pub async fn project_mob_event(&self, event: &MobEvent) -> MobStructuralEventEnvelope {
let envelope = self.build_envelope(event).await;
let _ = self.event_tx.send(envelope.clone());
envelope
}
pub async fn project_event_for_query(&self, event: &MobEvent) -> MobStructuralEventEnvelope {
self.build_envelope(event).await
}
pub(crate) async fn project_event_with_authority(&self, event: &MobEvent) -> ProjectedMobEvent {
ProjectedMobEvent {
envelope: self.build_envelope(event).await,
member_authority: extract_member_authority(&event.kind),
}
}
async fn build_envelope(&self, event: &MobEvent) -> MobStructuralEventEnvelope {
let cursor = event.cursor;
let mob_id = event.mob_id.as_str().to_string();
let timestamp_ms = event.timestamp.timestamp_millis().max(0) as u64;
let kind = event_kind_label(&event.kind).to_string();
let (run_id, step_id, agent_identity) = extract_structural_fields(&event.kind);
let data = serde_json::to_value(&event.kind).unwrap_or(Value::Null);
let (mob_labels, run_labels) = self.lookup_labels(&mob_id, run_id.as_deref()).await;
MobStructuralEventEnvelope {
event_id: format!("mob-evt-{cursor}"),
cursor,
mob_id,
timestamp_ms,
kind,
run_id,
step_id,
agent_identity,
mob_labels,
run_labels,
data,
}
}
async fn lookup_labels(
&self,
mob_id: &str,
run_id: Option<&str>,
) -> (BTreeMap<String, String>, BTreeMap<String, String>) {
let Some(table) = &self.metadata_table else {
return (BTreeMap::new(), BTreeMap::new());
};
let mob_labels = table
.get_labels(&MetadataScope::Mob(mob_id.to_string()))
.await;
let run_labels = match run_id {
Some(run_id) => {
table
.get_labels(&MetadataScope::Run(mob_id.to_string(), run_id.to_string()))
.await
}
None => BTreeMap::new(),
};
(mob_labels, run_labels)
}
}
pub const MOB_EVENTS_STREAM_PATH: &str = "/mobkit/mob_events/stream";
pub(crate) fn build_subscribe_url(
query: &EventQuery,
next_after_seq: Option<u64>,
fallback_cursor: u64,
) -> String {
let after_seq = next_after_seq
.or(query.after_seq)
.unwrap_or(fallback_cursor);
let mut serializer = form_urlencoded::Serializer::new(String::new());
serializer.append_pair("after_seq", &after_seq.to_string());
if let Some(value) = query.mob_id.as_deref() {
serializer.append_pair("mob_id", value);
}
if let Some(value) = query.run_id.as_deref() {
serializer.append_pair("run_id", value);
}
if let Some(value) = query.step_id.as_deref() {
serializer.append_pair("step_id", value);
}
if let Some(value) = query.identity.as_deref() {
serializer.append_pair("identity", value);
}
if let Some(value) = query.member_id.as_deref() {
serializer.append_pair("member_id", value);
}
if let Some(value) = query.since_ms {
serializer.append_pair("since_ms", &value.to_string());
}
if let Some(value) = query.until_ms {
serializer.append_pair("until_ms", &value.to_string());
}
if !query.event_types.is_empty() {
serializer.append_pair("event_types", &query.event_types.join(","));
}
format!("{MOB_EVENTS_STREAM_PATH}?{}", serializer.finish())
}
#[derive(Debug)]
pub enum MobEventsQueryError {
Stale {
after_cursor: u64,
latest_cursor: u64,
},
Backend(MobError),
}
impl std::fmt::Display for MobEventsQueryError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Stale {
after_cursor,
latest_cursor,
} => write!(
f,
"stale mob event cursor: requested {after_cursor}, latest {latest_cursor}"
),
Self::Backend(err) => write!(f, "{err}"),
}
}
}
impl std::error::Error for MobEventsQueryError {}
impl From<MobError> for MobEventsQueryError {
fn from(err: MobError) -> Self {
if let MobError::StaleEventCursor {
after_cursor,
latest_cursor,
} = err
{
Self::Stale {
after_cursor,
latest_cursor,
}
} else {
Self::Backend(err)
}
}
}
pub(crate) fn envelope_matches(envelope: &MobStructuralEventEnvelope, query: &EventQuery) -> bool {
if let Some(since) = query.since_ms
&& envelope.timestamp_ms < since
{
return false;
}
if let Some(until) = query.until_ms
&& envelope.timestamp_ms >= until
{
return false;
}
if let Some(mob_id) = query.mob_id.as_deref()
&& envelope.mob_id != mob_id
{
return false;
}
if let Some(run_id) = query.run_id.as_deref()
&& envelope.run_id.as_deref() != Some(run_id)
{
return false;
}
if let Some(step_id) = query.step_id.as_deref()
&& envelope.step_id.as_deref() != Some(step_id)
{
return false;
}
let identity_filter = query.identity.as_deref().or(query.member_id.as_deref());
if let Some(identity) = identity_filter
&& envelope.agent_identity.as_deref() != Some(identity)
{
return false;
}
if !query.event_types.is_empty() && !query.event_types.iter().any(|ty| ty == &envelope.kind) {
return false;
}
true
}
pub(crate) async fn query_ledger_with_filter(
events: &MobEventsView,
store: &MobEventsStore,
query: &EventQuery,
) -> Result<Vec<MobStructuralEventEnvelope>, MobEventsQueryError> {
Ok(
query_ledger_with_filter_selected(events, store, query, |projection| async move {
Some(projection.envelope)
})
.await?
.items,
)
}
pub(crate) async fn query_ledger_with_filter_selected<T, F, Fut>(
events: &MobEventsView,
store: &MobEventsStore,
query: &EventQuery,
mut selector: F,
) -> Result<MobEventsQueryPage<T>, MobEventsQueryError>
where
F: FnMut(ProjectedMobEvent) -> Fut,
Fut: Future<Output = Option<T>>,
{
let limit = query.limit.unwrap_or(DEFAULT_QUERY_LIMIT);
if limit == 0 {
let resume_after_seq = match query.after_seq {
Some(after_seq) => after_seq,
None => events.latest_cursor().await?,
};
return Ok(MobEventsQueryPage {
items: Vec::new(),
resume_after_seq,
});
}
if let Some(after_seq) = query.after_seq {
return scan_forward(events, store, query, after_seq, limit, &mut selector).await;
}
scan_backward(events, store, query, limit, &mut selector).await
}
async fn scan_forward<T, F, Fut>(
events: &MobEventsView,
store: &MobEventsStore,
query: &EventQuery,
after_seq: u64,
limit: usize,
selector: &mut F,
) -> Result<MobEventsQueryPage<T>, MobEventsQueryError>
where
F: FnMut(ProjectedMobEvent) -> Fut,
Fut: Future<Output = Option<T>>,
{
let mut results: Vec<T> = Vec::with_capacity(limit.min(QUERY_BATCH_SIZE));
let mut cursor = after_seq;
loop {
let batch = events.poll_strict(cursor, QUERY_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let cursor_before_batch = cursor;
for event in batch {
cursor = cursor.max(event.cursor);
let projection = store.project_event_with_authority(&event).await;
if envelope_matches(&projection.envelope, query)
&& let Some(selected) = selector(projection).await
{
results.push(selected);
if results.len() >= limit {
return Ok(MobEventsQueryPage {
items: results,
resume_after_seq: cursor,
});
}
}
}
if cursor <= cursor_before_batch {
break;
}
}
Ok(MobEventsQueryPage {
items: results,
resume_after_seq: cursor,
})
}
async fn scan_backward<T, F, Fut>(
events: &MobEventsView,
store: &MobEventsStore,
query: &EventQuery,
limit: usize,
selector: &mut F,
) -> Result<MobEventsQueryPage<T>, MobEventsQueryError>
where
F: FnMut(ProjectedMobEvent) -> Fut,
Fut: Future<Output = Option<T>>,
{
let latest = events.latest_cursor().await?;
if latest == 0 {
return Ok(MobEventsQueryPage {
items: Vec::new(),
resume_after_seq: 0,
});
}
let batch_size = QUERY_BATCH_SIZE as u64;
let mut window_end = latest;
let mut accumulator: Vec<T> = Vec::new();
loop {
let from = window_end.saturating_sub(batch_size);
let take = (window_end - from) as usize;
if take == 0 {
break;
}
let batch = events.poll_strict(from, take).await?;
if batch.is_empty() {
break;
}
let mut window_matches: Vec<T> = Vec::with_capacity(batch.len());
for event in batch {
let projection = store.project_event_with_authority(&event).await;
if envelope_matches(&projection.envelope, query)
&& let Some(selected) = selector(projection).await
{
window_matches.push(selected);
}
}
let mut combined = Vec::with_capacity(window_matches.len() + accumulator.len());
combined.append(&mut window_matches);
combined.append(&mut accumulator);
accumulator = combined;
if accumulator.len() >= limit || from == 0 {
break;
}
window_end = from;
}
if accumulator.len() > limit {
let drop = accumulator.len() - limit;
accumulator.drain(0..drop);
}
Ok(MobEventsQueryPage {
items: accumulator,
resume_after_seq: latest,
})
}
fn extract_member_authority(kind: &MobEventKind) -> Option<MobEventMemberAuthority> {
let (runtime_id, fence_token) = match kind {
MobEventKind::MemberSpawned(event) => (&event.agent_runtime_id, event.fence_token),
MobEventKind::MemberReset {
agent_runtime_id,
fence_token,
..
}
| MobEventKind::RemoteMemberRuntimeRetired {
agent_runtime_id,
fence_token,
..
}
| MobEventKind::RemoteMemberSupervisorRevoked {
agent_runtime_id,
fence_token,
..
} => (agent_runtime_id, *fence_token),
_ => return None,
};
Some(MobEventMemberAuthority {
runtime_id: runtime_id.clone(),
fence_token,
})
}
fn event_kind_label(kind: &MobEventKind) -> &'static str {
match kind {
MobEventKind::MobCreated { .. } => "mob_created",
MobEventKind::MobOwnerBridgeSessionBound { .. } => "mob_owner_bridge_session_bound",
MobEventKind::MobCompleted => "mob_completed",
MobEventKind::MobDestroying => "mob_destroying",
MobEventKind::MobDestroyStorageFinalizing => "mob_destroy_storage_finalizing",
MobEventKind::MobReset => "mob_reset",
MobEventKind::MemberSpawned(_) => "member_spawned",
MobEventKind::MemberSessionBindingRecovered(_) => "member_session_binding_recovered",
MobEventKind::MemberRetirementStarted { .. } => "member_retirement_started",
MobEventKind::MemberRetired { .. } => "member_retired",
MobEventKind::RespawnTopologyAbandoned { .. } => "respawn_topology_abandoned",
MobEventKind::RemoteMemberRuntimeRetired { .. } => "remote_member_runtime_retired",
MobEventKind::RemoteMemberSupervisorRevoked { .. } => "remote_member_supervisor_revoked",
MobEventKind::RemoteMemberReleaseConfirmed { .. } => "remote_member_release_confirmed",
MobEventKind::MemberReset { .. } => "member_reset",
MobEventKind::MemberKickoffUpdated { .. } => "member_kickoff_updated",
MobEventKind::MembersWired { .. } => "members_wired",
MobEventKind::MembersWiredBatch { .. } => "members_wired_batch",
MobEventKind::MembersUnwired { .. } => "members_unwired",
MobEventKind::ExternalPeerWired { .. } => "external_peer_wired",
MobEventKind::ExternalPeerUnwired { .. } => "external_peer_unwired",
MobEventKind::FlowStarted { .. } => "flow_started",
MobEventKind::FlowCompleted { .. } => "flow_completed",
MobEventKind::FlowFailed { .. } => "flow_failed",
MobEventKind::FlowCanceled { .. } => "flow_canceled",
MobEventKind::StepDispatched { .. } => "step_dispatched",
MobEventKind::StepTargetCompleted { .. } => "step_target_completed",
MobEventKind::StepTargetFailed { .. } => "step_target_failed",
MobEventKind::RemoteTurnObligationRecorded { .. } => "remote_turn_obligation_recorded",
MobEventKind::RemoteTurnOutcomeResolved { .. } => "remote_turn_outcome_resolved",
MobEventKind::RemoteTurnOutcomeAcknowledged { .. } => "remote_turn_outcome_acknowledged",
MobEventKind::RemoteTurnOutcomeDisposed { .. } => "remote_turn_outcome_disposed",
MobEventKind::PlacedCompletionLifecycleQuiesceStarted { .. } => {
"placed_completion_lifecycle_quiesce_started"
}
MobEventKind::MobStopped => "mob_stopped",
MobEventKind::PlacedCompletionLifecycleQuiesceEnded { .. } => {
"placed_completion_lifecycle_quiesce_ended"
}
MobEventKind::PlacedCompletionObligationRecorded { .. } => {
"placed_completion_obligation_recorded"
}
MobEventKind::PlacedCompletionCancellationRequested { .. } => {
"placed_completion_cancellation_requested"
}
MobEventKind::PlacedCompletionOutcomeResolved { .. } => {
"placed_completion_outcome_resolved"
}
MobEventKind::PlacedCompletionOutcomeClosed { .. } => "placed_completion_outcome_closed",
MobEventKind::PlacedCompletionOutcomeAcknowledged { .. } => {
"placed_completion_outcome_acknowledged"
}
MobEventKind::PlacedCompletionOutcomeDisposed { .. } => {
"placed_completion_outcome_disposed"
}
MobEventKind::PlacedKickoffObligationRecorded { .. } => {
"placed_kickoff_obligation_recorded"
}
MobEventKind::PlacedKickoffOutcomeResolved { .. } => "placed_kickoff_outcome_resolved",
MobEventKind::PlacedKickoffRejectedNoEffect { .. } => "placed_kickoff_rejected_no_effect",
MobEventKind::PlacedKickoffOutcomeAcknowledged { .. } => {
"placed_kickoff_outcome_acknowledged"
}
MobEventKind::PlacedKickoffOutcomeDisposed { .. } => "placed_kickoff_outcome_disposed",
MobEventKind::RemoteHostBindStarted { .. } => "remote_host_bind_started",
MobEventKind::RemoteHostBindConfirmed { .. } => "remote_host_bind_confirmed",
MobEventKind::RemoteHostBindCompleted { .. } => "remote_host_bind_completed",
MobEventKind::RemoteHostBindAbortedNoEffect { .. } => "remote_host_bind_aborted_no_effect",
MobEventKind::RemoteHostRevokeStarted { .. } => "remote_host_revoke_started",
MobEventKind::RemoteHostRevokeConfirmed { .. } => "remote_host_revoke_confirmed",
MobEventKind::RemoteHostRevokeCompleted { .. } => "remote_host_revoke_completed",
MobEventKind::StepCompleted { .. } => "step_completed",
MobEventKind::StepFailed { .. } => "step_failed",
MobEventKind::StepSkipped { .. } => "step_skipped",
MobEventKind::TopologyViolation { .. } => "topology_violation",
MobEventKind::SupervisorEscalation { .. } => "supervisor_escalation",
MobEventKind::OperatorActionRecorded { .. } => "operator_action_recorded",
MobEventKind::ObjectiveOwnerBound { .. } => "objective_owner_bound",
MobEventKind::ObjectiveConcluded { .. } => "objective_concluded",
}
}
fn decode_member_id(member_id: &str) -> String {
crate::member_comms_id::runtime_alias_str(member_id).into_owned()
}
pub(crate) fn extract_structural_fields(
kind: &MobEventKind,
) -> (Option<String>, Option<String>, Option<String>) {
match kind {
MobEventKind::FlowStarted { run_id, .. }
| MobEventKind::FlowCompleted { run_id, .. }
| MobEventKind::FlowFailed { run_id, .. }
| MobEventKind::FlowCanceled { run_id, .. } => (Some(run_id.to_string()), None, None),
MobEventKind::StepDispatched {
run_id,
step_id,
target,
}
| MobEventKind::StepTargetCompleted {
run_id,
step_id,
target,
..
} => (
Some(run_id.to_string()),
Some(step_id.as_str().to_string()),
Some(decode_member_id(target.identity.as_str())),
),
MobEventKind::StepTargetFailed {
run_id,
step_id,
target,
..
} => (
Some(run_id.to_string()),
Some(step_id.as_str().to_string()),
Some(decode_member_id(target.identity.as_str())),
),
MobEventKind::StepCompleted { run_id, step_id }
| MobEventKind::StepFailed {
run_id, step_id, ..
}
| MobEventKind::StepSkipped {
run_id, step_id, ..
} => (
Some(run_id.to_string()),
Some(step_id.as_str().to_string()),
None,
),
MobEventKind::SupervisorEscalation {
run_id,
step_id,
escalated_to,
} => (
Some(run_id.to_string()),
Some(step_id.as_str().to_string()),
Some(decode_member_id(escalated_to.as_str())),
),
MobEventKind::RemoteTurnObligationRecorded { obligation }
| MobEventKind::RemoteTurnOutcomeResolved { obligation }
| MobEventKind::RemoteTurnOutcomeAcknowledged { obligation }
| MobEventKind::RemoteTurnOutcomeDisposed { obligation } => (
Some(obligation.run_id.to_string()),
Some(obligation.step_id.as_str().to_string()),
Some(decode_member_id(obligation.agent_identity.as_str())),
),
MobEventKind::PlacedCompletionObligationRecorded { obligation }
| MobEventKind::PlacedCompletionCancellationRequested { obligation }
| MobEventKind::PlacedCompletionOutcomeResolved { obligation, .. }
| MobEventKind::PlacedCompletionOutcomeClosed { obligation, .. }
| MobEventKind::PlacedCompletionOutcomeAcknowledged { obligation }
| MobEventKind::PlacedCompletionOutcomeDisposed { obligation } => (
None,
None,
Some(decode_member_id(obligation.agent_identity.as_str())),
),
MobEventKind::PlacedKickoffObligationRecorded { obligation }
| MobEventKind::PlacedKickoffOutcomeResolved { obligation, .. }
| MobEventKind::PlacedKickoffRejectedNoEffect { obligation, .. }
| MobEventKind::PlacedKickoffOutcomeAcknowledged { obligation }
| MobEventKind::PlacedKickoffOutcomeDisposed { obligation } => (
None,
None,
Some(decode_member_id(obligation.agent_identity.as_str())),
),
MobEventKind::MemberSpawned(event) => (
None,
None,
Some(decode_member_id(event.agent_identity.as_str())),
),
MobEventKind::MemberSessionBindingRecovered(event) => (
None,
None,
Some(decode_member_id(event.agent_identity.as_str())),
),
MobEventKind::MemberRetirementStarted { agent_identity, .. }
| MobEventKind::MemberRetired { agent_identity, .. }
| MobEventKind::RespawnTopologyAbandoned { agent_identity, .. }
| MobEventKind::MemberReset { agent_identity, .. }
| MobEventKind::RemoteMemberRuntimeRetired { agent_identity, .. }
| MobEventKind::RemoteMemberSupervisorRevoked { agent_identity, .. }
| MobEventKind::RemoteMemberReleaseConfirmed { agent_identity, .. } => {
(None, None, Some(decode_member_id(agent_identity.as_str())))
}
MobEventKind::MemberKickoffUpdated { member, .. } => {
(None, None, Some(decode_member_id(member.as_str())))
}
MobEventKind::ObjectiveOwnerBound { owner, .. } => {
(None, None, Some(decode_member_id(owner.as_str())))
}
MobEventKind::ObjectiveConcluded { member, .. } => {
(None, None, Some(decode_member_id(member.as_str())))
}
MobEventKind::ExternalPeerWired { local, .. }
| MobEventKind::ExternalPeerUnwired { local, .. } => {
(None, None, Some(decode_member_id(local.as_str())))
}
MobEventKind::MobCreated { .. }
| MobEventKind::MobOwnerBridgeSessionBound { .. }
| MobEventKind::MobCompleted
| MobEventKind::MobDestroying
| MobEventKind::MobDestroyStorageFinalizing
| MobEventKind::MobReset
| MobEventKind::PlacedCompletionLifecycleQuiesceStarted { .. }
| MobEventKind::MobStopped
| MobEventKind::PlacedCompletionLifecycleQuiesceEnded { .. }
| MobEventKind::RemoteHostBindStarted { .. }
| MobEventKind::RemoteHostBindConfirmed { .. }
| MobEventKind::RemoteHostBindCompleted { .. }
| MobEventKind::RemoteHostBindAbortedNoEffect { .. }
| MobEventKind::RemoteHostRevokeStarted { .. }
| MobEventKind::RemoteHostRevokeConfirmed { .. }
| MobEventKind::RemoteHostRevokeCompleted { .. }
| MobEventKind::MembersWired { .. }
| MobEventKind::MembersWiredBatch { .. }
| MobEventKind::MembersUnwired { .. }
| MobEventKind::TopologyViolation { .. }
| MobEventKind::OperatorActionRecorded { .. } => (None, None, None),
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
mod tests {
use super::*;
use chrono::Utc;
use meerkat_mob::event::{MemberSpawnedEvent, MemberWireEdge};
use meerkat_mob::ids::{
AgentIdentity, AgentRuntimeId, FenceToken, FlowId, Generation, MobId, ProfileName, RunId,
StepId,
};
fn mob_event(cursor: u64, kind: MobEventKind) -> MobEvent {
MobEvent {
cursor,
timestamp: Utc::now(),
mob_id: MobId::from("test-mob"),
kind,
}
}
#[tokio::test]
async fn projects_flow_started_with_run_id_and_upstream_cursor() {
let store = MobEventsStore::new();
let run_id = RunId::new();
let envelope = store
.project_mob_event(&mob_event(
42,
MobEventKind::FlowStarted {
run_id: run_id.clone(),
flow_id: FlowId::from("flow-a"),
params: serde_json::json!({}),
},
))
.await;
assert_eq!(envelope.kind, "flow_started");
assert_eq!(envelope.cursor, 42);
assert_eq!(envelope.event_id, "mob-evt-42");
assert_eq!(
envelope.run_id.as_deref(),
Some(run_id.to_string().as_str())
);
assert_eq!(envelope.step_id, None);
assert_eq!(envelope.mob_id, "test-mob");
}
#[tokio::test]
async fn projects_step_dispatched_with_run_step_target() {
let store = MobEventsStore::new();
let identity = AgentIdentity::from("worker-1");
let run_id = RunId::new();
let envelope = store
.project_mob_event(&mob_event(
7,
MobEventKind::StepDispatched {
run_id: run_id.clone(),
step_id: StepId::from("step-a"),
target: AgentRuntimeId::initial(identity),
},
))
.await;
assert_eq!(envelope.kind, "step_dispatched");
assert_eq!(envelope.cursor, 7);
assert_eq!(
envelope.run_id.as_deref(),
Some(run_id.to_string().as_str())
);
assert_eq!(envelope.step_id.as_deref(), Some("step-a"));
assert_eq!(envelope.agent_identity.as_deref(), Some("worker-1"));
}
#[tokio::test]
async fn projects_member_spawned_with_identity() {
let store = MobEventsStore::new();
let identity = AgentIdentity::from("researcher");
let envelope = store
.project_mob_event(&mob_event(
3,
MobEventKind::MemberSpawned(MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(1),
AgentRuntimeId::initial(identity),
ProfileName::from("worker"),
)),
))
.await;
assert_eq!(envelope.kind, "member_spawned");
assert_eq!(envelope.agent_identity.as_deref(), Some("researcher"));
}
#[tokio::test]
async fn internal_projection_binds_spawn_to_exact_runtime_and_fence() {
let store = MobEventsStore::new();
let identity = AgentIdentity::from("generation-bound");
let runtime_id = AgentRuntimeId::initial(identity.clone());
let fence_token = FenceToken::new(41);
let event = mob_event(
9,
MobEventKind::MemberSpawned(MemberSpawnedEvent::new(
identity,
Generation::INITIAL,
fence_token,
runtime_id.clone(),
ProfileName::from("worker"),
)),
);
let public = store.project_event_for_query(&event).await;
let projected = store.project_event_with_authority(&event).await;
assert_eq!(projected.envelope, public);
assert_eq!(
projected.member_authority,
Some(MobEventMemberAuthority {
runtime_id,
fence_token,
})
);
let wire = serde_json::to_value(&projected.envelope).expect("public envelope wire value");
assert!(
wire.get("member_authority").is_none(),
"internal authority must not alter the public envelope contract"
);
}
#[tokio::test]
async fn internal_projection_leaves_fenceless_agent_event_unbound() {
let store = MobEventsStore::new();
let identity = AgentIdentity::from("fenceless-target");
let projected = store
.project_event_with_authority(&mob_event(
10,
MobEventKind::StepDispatched {
run_id: RunId::new(),
step_id: StepId::from("step-a"),
target: AgentRuntimeId::initial(identity),
},
))
.await;
assert_eq!(
projected.envelope.agent_identity.as_deref(),
Some("fenceless-target")
);
assert_eq!(projected.member_authority, None);
}
#[tokio::test]
async fn projects_identity_first_member_in_public_alias_space_and_filters_round_trip() {
let alias = "rt:review:singleton:0";
let encoded = crate::member_comms_id::mob_member_id_str(alias).into_owned();
assert_ne!(encoded, alias, "alias must actually encode for this test");
let identity = AgentIdentity::from(encoded.as_str());
for kind in [
MobEventKind::MemberSpawned(MemberSpawnedEvent::new(
identity.clone(),
Generation::INITIAL,
FenceToken::new(1),
AgentRuntimeId::initial(identity.clone()),
ProfileName::from("review"),
)),
MobEventKind::MemberRetired {
agent_identity: identity.clone(),
generation: Generation::INITIAL,
role: ProfileName::from("review"),
},
MobEventKind::StepDispatched {
run_id: RunId::new(),
step_id: StepId::from("step-a"),
target: AgentRuntimeId::initial(identity.clone()),
},
] {
let store = MobEventsStore::new();
let envelope = store.project_mob_event(&mob_event(1, kind)).await;
assert_eq!(
envelope.agent_identity.as_deref(),
Some(alias),
"structural envelope must decode the roster id to the alias"
);
assert!(
!envelope
.agent_identity
.as_deref()
.unwrap()
.starts_with("mk--"),
"encoded mk-- id must not leak onto the public surface"
);
let query = EventQuery {
identity: Some(alias.to_string()),
..EventQuery::default()
};
assert!(
envelope_matches(&envelope, &query),
"identity-filter by the public alias must match"
);
let member_query = EventQuery {
member_id: Some(alias.to_string()),
..EventQuery::default()
};
assert!(
envelope_matches(&envelope, &member_query),
"member_id-filter by the public alias must match"
);
}
}
#[tokio::test]
async fn projects_members_wired_batch_as_compact_structural_event() {
let store = MobEventsStore::new();
let envelope = store
.project_mob_event(&mob_event(
8,
MobEventKind::MembersWiredBatch {
edges: vec![MemberWireEdge {
a: AgentIdentity::from("alpha"),
b: AgentIdentity::from("beta"),
}],
},
))
.await;
assert_eq!(envelope.kind, "members_wired_batch");
assert_eq!(envelope.agent_identity, None);
assert_eq!(envelope.data["edges"][0]["a"], serde_json::json!("alpha"));
assert_eq!(envelope.data["edges"][0]["b"], serde_json::json!("beta"));
}
#[tokio::test]
async fn project_event_for_query_does_not_broadcast() {
let store = MobEventsStore::new();
let mut rx = store.subscribe();
let _ = store
.project_event_for_query(&mob_event(
1,
MobEventKind::FlowStarted {
run_id: RunId::new(),
flow_id: FlowId::from("flow-a"),
params: serde_json::json!({}),
},
))
.await;
assert!(rx.try_recv().is_err());
}
}