use std::borrow::Cow;
pub(crate) fn durable_identity_label(
labels: &std::collections::BTreeMap<String, String>,
) -> Option<&str> {
labels
.get("agent_identity")
.map(String::as_str)
.map(str::trim)
.filter(|value| !value.is_empty() && !is_reserved_generated_alias(value))
}
pub(crate) fn validate_raw_identity_labels(
labels: &std::collections::BTreeMap<String, String>,
) -> Result<(), &'static str> {
if labels.contains_key("agent_identity") {
return Err(
"labels.agent_identity is runtime-authoritative and may not be supplied by raw member creation",
);
}
Ok(())
}
const MARKER: &str = "mk--";
fn is_valid_comms_component(s: &str) -> bool {
let mut chars = s.chars();
let Some(first) = chars.next() else {
return false;
};
if !first.is_ascii_alphabetic() && first != '_' {
return false;
}
chars.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
}
fn escape_body(s: &str) -> String {
let mut out = String::with_capacity(s.len() + 8);
for c in s.chars() {
match c {
'_' => out.push_str("__"),
':' => out.push_str("_c"),
c if c.is_ascii_alphanumeric() || c == '-' => out.push(c),
c => {
out.push_str("_x");
out.push_str(&format!("{:x}", c as u32));
out.push('_');
}
}
}
out
}
fn unescape_body(s: &str) -> Option<String> {
let mut out = String::with_capacity(s.len());
let mut chars = s.chars();
while let Some(c) = chars.next() {
if c != '_' {
out.push(c);
continue;
}
match chars.next()? {
'_' => out.push('_'),
'c' => out.push(':'),
'x' => {
let mut hex = String::new();
loop {
match chars.next()? {
'_' => break,
h => hex.push(h),
}
}
let code = u32::from_str_radix(&hex, 16).ok()?;
out.push(char::from_u32(code)?);
}
_ => return None,
}
}
Some(out)
}
pub fn mob_member_id_str(alias: &str) -> Cow<'_, str> {
if is_valid_comms_component(alias) && !alias.starts_with(MARKER) {
Cow::Borrowed(alias)
} else {
Cow::Owned(format!("{MARKER}{}", escape_body(alias)))
}
}
pub fn mob_member_id(alias: &str) -> meerkat_mob::ids::AgentIdentity {
meerkat_mob::ids::AgentIdentity::from(mob_member_id_str(alias).as_ref())
}
pub fn runtime_alias_str(member_id: &str) -> Cow<'_, str> {
match member_id.strip_prefix(MARKER) {
Some(body) => match unescape_body(body) {
Some(alias) => Cow::Owned(alias),
None => Cow::Borrowed(member_id),
},
None => Cow::Borrowed(member_id),
}
}
pub(crate) fn is_reserved_generated_alias(member_id: &str) -> bool {
runtime_alias_str(member_id).starts_with("rt:")
}
pub(crate) fn durable_identity_from_runtime_alias(alias: &str) -> Option<String> {
let rest = alias.strip_prefix("rt:")?;
let (identity, generation) = rest.rsplit_once(':')?;
if identity.is_empty() || generation.parse::<u64>().is_err() {
return None;
}
Some(identity.to_string())
}
pub(crate) fn logical_memory_identity(member_id_or_alias: &str) -> String {
let alias = runtime_alias_str(member_id_or_alias);
durable_identity_from_runtime_alias(&alias).unwrap_or_else(|| alias.into_owned())
}
pub(crate) fn uses_reserved_roster_marker(member_id: &str) -> bool {
member_id.trim().starts_with(MARKER)
}
pub(crate) fn validate_public_member_alias(field: &str, value: &str) -> Result<(), String> {
if uses_reserved_roster_marker(value) {
return Err(format!(
"{field} may not use the reserved encoded roster-id namespace"
));
}
Ok(())
}
pub(crate) fn validate_public_rpc_member_aliases(params: &serde_json::Value) -> Result<(), String> {
const TARGET_FIELDS: &[&str] = &[
"identity",
"member_id",
"agent_id",
"agent_identity",
"meerkat_id",
"local_member_id",
"remote_member_id",
"from_member_id",
"source_member_id",
"runtime_member_id",
"agent_runtime_id",
];
let Some(params) = params.as_object() else {
return Ok(());
};
for field in TARGET_FIELDS {
if let Some(value) = params.get(*field).and_then(serde_json::Value::as_str) {
validate_public_member_alias(field, value)?;
}
}
Ok(())
}
pub(crate) async fn validate_raw_member_target(
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_id: &str,
) -> Result<String, String> {
let member_id = member_id.trim();
let alias = runtime_alias_str(member_id).into_owned();
if uses_reserved_roster_marker(member_id) || is_reserved_generated_alias(&alias) {
return Err(format!(
"member id '{member_id}' uses an identity-runtime reserved namespace"
));
}
if let Some(identity_runtime) = identity_runtime
&& let Some(identity) = identity_runtime.identity_for_member_mutation(&alias).await
{
return Err(format!(
"member id '{alias}' is owned by durable identity '{identity}'"
));
}
Ok(alias)
}
pub(crate) struct RawMemberTargetReservation {
aliases: Vec<String>,
_alias_guards: Vec<tokio::sync::OwnedMutexGuard<()>>,
}
impl RawMemberTargetReservation {
pub(crate) fn alias(&self) -> &str {
self.aliases.first().map(String::as_str).unwrap_or_default()
}
pub(crate) fn aliases(&self) -> &[String] {
&self.aliases
}
}
pub(crate) async fn reserve_raw_member_target(
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_id: &str,
) -> Result<RawMemberTargetReservation, String> {
reserve_raw_member_targets(identity_runtime, std::iter::once(member_id)).await
}
pub(crate) async fn reserve_raw_member_targets<'a>(
identity_runtime: Option<&std::sync::Arc<crate::identity_first::IdentityRuntime>>,
member_ids: impl IntoIterator<Item = &'a str>,
) -> Result<RawMemberTargetReservation, String> {
let member_ids = member_ids
.into_iter()
.map(str::trim)
.map(ToString::to_string)
.collect::<Vec<_>>();
if member_ids.is_empty() {
return Err("raw member reservation requires at least one member id".to_string());
}
for member_id in &member_ids {
validate_raw_member_target(identity_runtime, member_id).await?;
}
let canonical_aliases = member_ids
.iter()
.map(|member_id| runtime_alias_str(member_id).into_owned())
.collect::<Vec<_>>();
let sorted_aliases = canonical_aliases
.iter()
.cloned()
.collect::<std::collections::BTreeSet<_>>();
let mut alias_guards = Vec::with_capacity(sorted_aliases.len());
if let Some(identity_runtime) = identity_runtime {
for alias in sorted_aliases {
let lock = identity_runtime.raw_member_alias_lock(&alias).await;
alias_guards.push(lock.lock_owned().await);
}
}
let mut aliases = Vec::with_capacity(member_ids.len());
for member_id in member_ids {
aliases.push(validate_raw_member_target(identity_runtime, &member_id).await?);
}
debug_assert_eq!(aliases, canonical_aliases);
Ok(RawMemberTargetReservation {
aliases,
_alias_guards: alias_guards,
})
}
pub fn runtime_event_alias(runtime_id: &meerkat_mob::ids::AgentRuntimeId) -> String {
format!(
"{}:{}",
runtime_alias_str(runtime_id.identity.as_str()),
runtime_id.generation.get()
)
}
pub(crate) fn canonical_correlation_id(correlation_id: &str) -> std::borrow::Cow<'_, str> {
if uuid::Uuid::try_parse(correlation_id)
.is_ok_and(|parsed| !parsed.is_nil() && parsed.to_string() == correlation_id)
{
return std::borrow::Cow::Borrowed(correlation_id);
}
let namespace = uuid::Uuid::new_v5(
&uuid::Uuid::NAMESPACE_URL,
b"rkat-mobkit:delivery-correlation",
);
std::borrow::Cow::Owned(uuid::Uuid::new_v5(&namespace, correlation_id.as_bytes()).to_string())
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
#[test]
fn canonical_correlation_id_passes_canonical_uuids_and_canonicalizes_the_rest() {
let occurrence = "7c9e6679-7425-40de-944b-e07fc1f90ae7";
assert_eq!(canonical_correlation_id(occurrence).as_ref(), occurrence);
let a1 = canonical_correlation_id("telegram:primary/i:769307582").into_owned();
let a2 = canonical_correlation_id("telegram:primary/i:769307582").into_owned();
let b = canonical_correlation_id("telegram:primary/i:769307583").into_owned();
assert_eq!(a1, a2);
assert_ne!(a1, b);
let parsed = uuid::Uuid::try_parse(&a1).expect("canonical output");
assert!(!parsed.is_nil());
assert_eq!(parsed.to_string(), a1);
assert_ne!(
canonical_correlation_id("7C9E6679-7425-40DE-944B-E07FC1F90AE7").as_ref(),
"7C9E6679-7425-40DE-944B-E07FC1F90AE7"
);
assert_ne!(
canonical_correlation_id("00000000-0000-0000-0000-000000000000").as_ref(),
"00000000-0000-0000-0000-000000000000"
);
}
fn raw_lock_test_runtime()
-> Result<std::sync::Arc<crate::identity_first::IdentityRuntime>, Box<dyn std::error::Error>>
{
Ok(std::sync::Arc::new(
crate::identity_first::IdentityRuntime::new(
crate::identity_first::IdentityRuntimeConfig {
continuity_store: std::sync::Arc::new(
crate::identity_first::LocalContinuityStore::in_memory()?,
),
lease_provider: std::sync::Arc::new(
crate::identity_first::LocalLeaseProvider::new(),
),
runtime_instance_id: "raw-lock-test".to_string(),
has_runtime_store: true,
durability_policy: crate::identity_first::DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
},
),
))
}
#[test]
fn plain_member_names_pass_through_unchanged() {
for name in ["worker", "worker-one", "_internal", "Agent7"] {
assert_eq!(mob_member_id_str(name), name);
assert_eq!(runtime_alias_str(name), name);
}
}
#[test]
fn colon_aliases_round_trip_and_are_comms_safe() {
for alias in [
"rt:review:singleton:0",
"rt:channel:C0SMOKEOB3:0",
"agent:beta",
"review:singleton",
"rt:agent-x:0",
"with_underscore:and:colons",
] {
let encoded = mob_member_id_str(alias);
assert!(
is_valid_comms_component(&encoded),
"{encoded:?} must satisfy meerkat 0.7 MemberCommsName"
);
assert_eq!(runtime_alias_str(&encoded), alias, "round trip of {alias}");
}
}
#[test]
fn non_ascii_and_punctuation_round_trip() {
for alias in ["user@host", "a.b:c", "9starts-with-digit", "ünïcode:1"] {
let encoded = mob_member_id_str(alias);
assert!(is_valid_comms_component(&encoded), "{encoded:?}");
assert_eq!(runtime_alias_str(&encoded), alias);
}
}
#[test]
fn reserved_marker_names_re_encode_so_round_trip_holds() {
let alias = "mk--rt_creview";
let encoded = mob_member_id_str(alias);
assert_ne!(encoded, alias, "marker-prefixed names must be re-encoded");
assert_eq!(runtime_alias_str(&encoded), alias);
}
#[test]
fn generated_alias_reservation_applies_to_public_and_encoded_forms() {
let alias = "rt:review:singleton:0";
let encoded = mob_member_id_str(alias);
assert!(is_reserved_generated_alias(alias));
assert!(is_reserved_generated_alias(&encoded));
assert!(!is_reserved_generated_alias("review:singleton"));
}
#[tokio::test]
async fn raw_member_targets_reject_reserved_and_registered_durable_aliases()
-> Result<(), Box<dyn std::error::Error>> {
let identity_runtime = std::sync::Arc::new(crate::identity_first::IdentityRuntime::new(
crate::identity_first::IdentityRuntimeConfig {
continuity_store: std::sync::Arc::new(
crate::identity_first::LocalContinuityStore::in_memory()?,
),
lease_provider: std::sync::Arc::new(
crate::identity_first::LocalLeaseProvider::new(),
),
runtime_instance_id: "raw-target-test".to_string(),
has_runtime_store: true,
durability_policy: crate::identity_first::DurabilityPolicy::SyncWriteThrough,
bridge: None,
default_timeout: None,
},
));
let identity = crate::identity_first::AgentIdentity::parse("lead")?;
identity_runtime
.register(
crate::identity_first::DurableAgentSpec {
identity,
profile: meerkat_mob::ProfileName::from("worker"),
addressability: crate::identity_first::AgentAddressability::Addressable,
display_name: None,
labels: std::collections::BTreeMap::new(),
context: None,
additional_instructions: Vec::new(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
placement: None,
},
crate::identity_first::IdentityLifecycleState::Dormant,
None,
None,
)
.await;
for target in [
"lead".to_string(),
"rt:lead:0".to_string(),
mob_member_id_str("rt:lead:0").into_owned(),
] {
let error = match validate_raw_member_target(Some(&identity_runtime), &target).await {
Err(error) => error,
Ok(alias) => {
return Err(format!(
"identity-owned target {target:?} was admitted as {alias:?}"
)
.into());
}
};
assert!(error.contains("identity-runtime") || error.contains("durable identity"));
}
assert_eq!(
validate_raw_member_target(Some(&identity_runtime), "worker").await,
Ok("worker".to_string())
);
let reservation = reserve_raw_member_target(Some(&identity_runtime), "worker").await?;
assert_eq!(reservation.alias(), "worker");
assert!(
tokio::time::timeout(
std::time::Duration::from_millis(25),
reserve_raw_member_target(Some(&identity_runtime), " worker "),
)
.await
.is_err(),
"canonical spellings of one alias must block on the same reservation"
);
let other_reservation = match tokio::time::timeout(
std::time::Duration::from_secs(1),
reserve_raw_member_target(Some(&identity_runtime), "other-worker"),
)
.await
{
Ok(result) => result?,
Err(error) => {
return Err(
format!("different raw aliases must reserve concurrently: {error}").into(),
);
}
};
assert_eq!(other_reservation.alias(), "other-worker");
drop(other_reservation);
drop(reservation);
let same_alias = match tokio::time::timeout(
std::time::Duration::from_secs(1),
reserve_raw_member_target(Some(&identity_runtime), " worker "),
)
.await
{
Ok(result) => result?,
Err(error) => {
return Err(format!("alias lock released with reservation: {error}").into());
}
};
assert_eq!(same_alias.alias(), "worker");
drop(same_alias);
let batch = match tokio::time::timeout(
std::time::Duration::from_secs(1),
reserve_raw_member_targets(Some(&identity_runtime), ["zeta", " alpha ", "zeta"]),
)
.await
{
Ok(result) => result?,
Err(error) => {
return Err(format!(
"sorted deduplicated acquisition must not self-deadlock: {error}"
)
.into());
}
};
assert_eq!(batch.aliases(), ["zeta", "alpha", "zeta"]);
Ok(())
}
#[tokio::test]
async fn rejected_unique_aliases_do_not_allocate_keyed_locks()
-> Result<(), Box<dyn std::error::Error>> {
let identity_runtime = raw_lock_test_runtime()?;
for index in 0..4_096 {
let alias = format!("rt:rejected:{index}:0");
assert!(
reserve_raw_member_target(Some(&identity_runtime), &alias)
.await
.is_err(),
"generated alias {alias} must be rejected"
);
}
let (entry_count, _next_sweep_len, sweep_count) =
identity_runtime.raw_member_alias_lock_metrics().await;
assert_eq!(
entry_count, 0,
"preflight rejection must not allocate locks"
);
assert_eq!(sweep_count, 0, "no allocation means no sweep work");
Ok(())
}
#[tokio::test]
async fn active_alias_lock_survives_sweep_and_blocks_same_alias()
-> Result<(), Box<dyn std::error::Error>> {
let identity_runtime = raw_lock_test_runtime()?;
let held_lock = identity_runtime.raw_member_alias_lock("held-worker").await;
let held_guard = held_lock.clone().lock_owned().await;
for index in 0..600 {
let lock = identity_runtime
.raw_member_alias_lock(&format!("transient-{index}"))
.await;
drop(lock);
}
let (_entry_count, _next_sweep_len, sweep_count) =
identity_runtime.raw_member_alias_lock_metrics().await;
assert!(sweep_count >= 2, "test must exercise opportunistic sweep");
let observed = identity_runtime
.raw_member_alias_lock(" held-worker ")
.await;
assert!(
std::sync::Arc::ptr_eq(&held_lock, &observed),
"sweep must retain the live mutex for canonical spellings"
);
drop(observed);
let waiter_runtime = identity_runtime.clone();
let mut waiter = tokio::spawn(async move {
let lock = waiter_runtime.raw_member_alias_lock("held-worker").await;
let _guard = lock.lock_owned().await;
});
assert!(
tokio::time::timeout(std::time::Duration::from_millis(25), &mut waiter)
.await
.is_err(),
"same-alias waiter must remain blocked on the live pre-sweep lock"
);
drop(held_guard);
tokio::time::timeout(std::time::Duration::from_secs(1), waiter).await??;
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_same_alias_lock_creation_has_one_critical_section()
-> Result<(), Box<dyn std::error::Error>> {
const TASKS: usize = 64;
let identity_runtime = raw_lock_test_runtime()?;
drop(identity_runtime.raw_member_alias_lock("same-worker").await);
let barrier = std::sync::Arc::new(tokio::sync::Barrier::new(TASKS + 1));
let active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let max_active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut tasks = Vec::with_capacity(TASKS);
for _ in 0..TASKS {
let runtime = identity_runtime.clone();
let barrier = barrier.clone();
let active = active.clone();
let max_active = max_active.clone();
tasks.push(tokio::spawn(async move {
barrier.wait().await;
let lock = runtime.raw_member_alias_lock("same-worker").await;
let _guard = lock.lock_owned().await;
let now = active.fetch_add(1, std::sync::atomic::Ordering::SeqCst) + 1;
max_active.fetch_max(now, std::sync::atomic::Ordering::SeqCst);
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
active.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
}));
}
barrier.wait().await;
for task in tasks {
task.await?;
}
assert_eq!(
max_active.load(std::sync::atomic::Ordering::SeqCst),
1,
"write-side recheck must prevent split locks for one alias"
);
Ok(())
}
#[tokio::test]
async fn large_live_alias_set_uses_geometric_sweeps_and_reclaims_after_drop()
-> Result<(), Box<dyn std::error::Error>> {
const ALIASES: usize = 4_096;
let identity_runtime = raw_lock_test_runtime()?;
let mut live_locks = Vec::with_capacity(ALIASES);
for index in 0..ALIASES {
live_locks.push(
identity_runtime
.raw_member_alias_lock(&format!("roster-{index}"))
.await,
);
}
let (entry_count, next_sweep_len, sweep_count) =
identity_runtime.raw_member_alias_lock_metrics().await;
assert_eq!(entry_count, ALIASES);
assert!(next_sweep_len >= ALIASES);
assert!(
sweep_count <= 8,
"live roster acquisition must sweep geometrically, got {sweep_count} sweeps"
);
drop(live_locks);
let post_drop_lock = identity_runtime
.raw_member_alias_lock("post-roster-sweep")
.await;
let (entry_count, next_sweep_len, post_drop_sweeps) =
identity_runtime.raw_member_alias_lock_metrics().await;
assert_eq!(entry_count, 1, "next miss must reclaim the expired cohort");
assert_eq!(next_sweep_len, 256);
assert_eq!(post_drop_sweeps, sweep_count + 1);
drop(post_drop_lock);
Ok(())
}
#[test]
fn raw_identity_labels_reserve_agent_identity_for_the_runtime() {
let mut labels = std::collections::BTreeMap::new();
labels.insert(
"agent_identity".to_string(),
"forged-durable-id".to_string(),
);
assert!(validate_raw_identity_labels(&labels).is_err());
labels.insert(
"agent_identity".to_string(),
mob_member_id_str("rt:forged:0").into_owned(),
);
assert!(validate_raw_identity_labels(&labels).is_err());
assert_eq!(durable_identity_label(&labels), None);
labels.remove("agent_identity");
labels.insert("team".to_string(), "red".to_string());
assert_eq!(validate_raw_identity_labels(&labels), Ok(()));
}
#[test]
fn public_rpc_ingress_rejects_encoded_roster_targets() -> Result<(), Box<dyn std::error::Error>>
{
for field in [
"identity",
"member_id",
"agent_id",
"agent_identity",
"meerkat_id",
"local_member_id",
"remote_member_id",
"from_member_id",
"source_member_id",
] {
let params = serde_json::json!({ (field): " mk--rt_csecret_c0 " });
let error = match validate_public_rpc_member_aliases(¶ms) {
Err(error) => error,
Ok(()) => {
return Err(format!(
"encoded roster id in {field} was admitted as a public alias"
)
.into());
}
};
assert!(error.contains(field), "{field}: {error}");
}
assert_eq!(
validate_public_rpc_member_aliases(&serde_json::json!({
"identity": "rt:secret:0",
"member_id": "plain-member",
})),
Ok(())
);
Ok(())
}
#[test]
fn runtime_event_alias_decodes_encoded_roster_member_ids() {
use meerkat_mob::ids::{AgentIdentity, AgentRuntimeId, Generation};
let encoded = mob_member_id_str("rt:review:singleton:0").into_owned();
let runtime_id =
AgentRuntimeId::new(AgentIdentity::from(encoded.as_str()), Generation::new(1));
assert_eq!(runtime_event_alias(&runtime_id), "rt:review:singleton:0:1");
let runtime_id = AgentRuntimeId::initial(AgentIdentity::from("worker-one"));
assert_eq!(runtime_event_alias(&runtime_id), "worker-one:0");
}
#[test]
fn encoding_is_injective_across_pass_through_and_encoded_forms() {
let inputs = [
"rt:review:singleton:0",
"rt-review-singleton-0",
"rt_creview_csingleton_c0",
"mk--rt_creview_csingleton_c0",
];
let mut seen = std::collections::BTreeSet::new();
for input in inputs {
assert!(
seen.insert(mob_member_id_str(input).into_owned()),
"collision on {input}"
);
}
}
fn defensive_corpus() -> Vec<String> {
let mut v: Vec<String> = [
"worker",
"worker-one",
"_internal",
"Agent7",
"a",
"identity:parent-1",
"identity:parent-2",
"identity:child-1",
"identity:child-2",
"domain:calendar",
"domain:school",
"domain:health",
"domain:home",
"domain:home-automation",
"domain:finance",
"domain:discovery",
"family-group:main",
"triage:main",
"gate:main",
"rt:identity:parent-1:0",
"rt:domain:home-automation:3",
"rt:channel:C0SMOKEOB3:0",
"rt:review:singleton:12",
"a:b_c",
"with_underscore:and:colons",
"user@host",
"a.b:c",
"9starts-with-digit",
"::",
":",
"-",
"mk--collision",
]
.into_iter()
.map(String::from)
.collect();
v.push("ünïcode:1".to_string());
v.push("emoji-😀:2".to_string());
v.push("Ω≈ç:3".to_string());
v.push("tab\tnl\n".to_string());
v.push("\u{0}\u{1f}x".to_string());
v.push("a b".to_string());
v.push(String::new());
v
}
#[test]
fn homecore_corpus_round_trips_and_is_comms_safe() {
for alias in defensive_corpus() {
let encoded = mob_member_id_str(&alias).into_owned();
assert!(
is_valid_comms_component(&encoded),
"encoded {encoded:?} (from {alias:?}) is not a valid comms component"
);
assert!(
!encoded.contains([':', '/']),
"encoded {encoded:?} leaks a routing separator"
);
assert_eq!(
runtime_alias_str(&encoded),
alias,
"round trip of {alias:?}"
);
}
}
#[test]
fn encoded_ids_satisfy_meerkat_member_comms_name_validator() {
use meerkat_core::connection::MemberCommsName;
for alias in defensive_corpus() {
let encoded = mob_member_id_str(&alias).into_owned();
assert!(
MemberCommsName::new("homecore-mob", "worker", encoded.clone()).is_ok(),
"MemberCommsName::new rejected encoded id {encoded:?} (from alias {alias:?})"
);
}
}
#[test]
fn exact_wire_format_for_known_aliases() {
let cases = [
("worker-one", "worker-one"), ("identity:parent-1", "mk--identity_cparent-1"),
("domain:home-automation", "mk--domain_chome-automation"),
("family-group:main", "mk--family-group_cmain"),
("rt:identity:parent-1:0", "mk--rt_cidentity_cparent-1_c0"),
("triage:main", "mk--triage_cmain"),
("a:b_c", "mk--a_cb__c"),
("a b", "mk--a_x20_b"),
("", "mk--"),
];
for (alias, expected) in cases {
assert_eq!(mob_member_id_str(alias), expected, "encode {alias:?}");
assert_eq!(runtime_alias_str(expected), alias, "decode {expected:?}");
}
}
#[test]
fn malformed_encoded_bodies_decode_to_literal_without_panic() {
for bad in [
"mk--_", "mk--_z", "mk--_x", "mk--_x_", "mk--_xZZ_", "mk--_x110000_", "mk--_xffffffff_", "mk--abc_", ] {
assert_eq!(runtime_alias_str(bad), bad, "literal fallback for {bad:?}");
}
}
#[test]
fn empty_and_boundary_strings_round_trip() {
assert_eq!(mob_member_id_str(""), "mk--");
assert_eq!(runtime_alias_str("mk--"), "");
for alias in [":", "_", "-", "::", "_x", "mk", "mk-", "mk--"] {
let encoded = mob_member_id_str(alias).into_owned();
assert!(
is_valid_comms_component(&encoded),
"{encoded:?} not comms-safe"
);
assert_eq!(runtime_alias_str(&encoded), alias, "round trip {alias:?}");
}
}
#[test]
fn decode_then_encode_canonicalizes_both_id_spaces() {
for (input, canonical) in [
("rt:domain:home:0", "mk--rt_cdomain_chome_c0"),
("mk--rt_cdomain_chome_c0", "mk--rt_cdomain_chome_c0"),
("domain:home", "mk--domain_chome"),
("mk--domain_chome", "mk--domain_chome"),
("digest-owner", "digest-owner"),
] {
let key = mob_member_id_str(runtime_alias_str(input).as_ref()).into_owned();
assert_eq!(key, canonical, "canonical roster key for {input:?}");
}
assert_ne!(
mob_member_id_str("mk--rt_cdomain_chome_c0"),
"mk--rt_cdomain_chome_c0",
"encode alone re-encodes roster ids — never use it on binding member ids"
);
}
#[test]
fn decode_is_identity_on_unencoded_ids() {
for id in [
"worker",
"worker-one",
"_internal",
"Agent7",
"rt-review-singleton-0",
"a_cb", ] {
assert_eq!(runtime_alias_str(id), id);
}
}
#[test]
fn runtime_event_alias_across_generations_and_forms() {
use meerkat_mob::ids::{AgentIdentity, AgentRuntimeId, Generation};
let encoded = mob_member_id_str("rt:identity:parent-1:0").into_owned();
let rid = AgentRuntimeId::new(AgentIdentity::from(encoded.as_str()), Generation::new(7));
assert_eq!(runtime_event_alias(&rid), "rt:identity:parent-1:0:7");
let encoded = mob_member_id_str("domain:home-automation").into_owned();
let rid = AgentRuntimeId::new(AgentIdentity::from(encoded.as_str()), Generation::new(1));
assert_eq!(runtime_event_alias(&rid), "domain:home-automation:1");
let rid = AgentRuntimeId::initial(AgentIdentity::from("triage-main"));
assert_eq!(runtime_event_alias(&rid), "triage-main:0");
}
#[test]
fn round_trip_is_total_and_injective_over_generated_corpus() {
use meerkat_core::connection::MemberCommsName;
use std::collections::BTreeMap;
let charset = ['a', 'Z', '0', '-', '_', ':', '.', '@', ' ', 'ü', '😀'];
let mut aliases: Vec<String> = vec![String::new()];
let mut frontier = vec![String::new()];
for _ in 0..3 {
let mut next = Vec::new();
for prefix in &frontier {
for c in charset {
let mut s = prefix.clone();
s.push(c);
aliases.push(s.clone());
next.push(s);
}
}
frontier = next;
}
let mut encoded_to_alias: BTreeMap<String, String> = BTreeMap::new();
for alias in &aliases {
let encoded = mob_member_id_str(alias).into_owned();
assert!(
MemberCommsName::new("m", "r", encoded.clone()).is_ok(),
"encoded {encoded:?} (from {alias:?}) is not comms-safe"
);
assert_eq!(
&runtime_alias_str(&encoded).into_owned(),
alias,
"round trip of {alias:?}"
);
if let Some(prev) = encoded_to_alias.insert(encoded.clone(), alias.clone()) {
assert_eq!(
&prev, alias,
"collision: {prev:?} and {alias:?} both encode to {encoded:?}"
);
}
}
}
#[test]
fn logical_memory_identity_normalizes_every_member_id_shape() {
let encoded_gen0 = mob_member_id_str("rt:identity:parent-1:0").into_owned();
let encoded_gen7 = mob_member_id_str("rt:identity:parent-1:7").into_owned();
assert_eq!(logical_memory_identity(&encoded_gen0), "identity:parent-1");
assert_eq!(logical_memory_identity(&encoded_gen7), "identity:parent-1");
assert_eq!(
logical_memory_identity("rt:identity:parent-1:0"),
"identity:parent-1"
);
assert_eq!(
logical_memory_identity(&mob_member_id_str("review:singleton")),
"review:singleton"
);
assert_eq!(logical_memory_identity("helper"), "helper");
assert_eq!(logical_memory_identity("rt:oddly:named"), "rt:oddly:named");
for id in ["identity:parent-1", "review:singleton", "helper"] {
assert_eq!(logical_memory_identity(id), id);
}
}
}