use std::collections::BTreeMap;
use std::sync::Arc;
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use crate::unified_runtime::ConsoleEventStore;
pub(crate) type SharedConsoleSpawnSinkSlot = Arc<std::sync::RwLock<Option<ConsoleSpawnSink>>>;
pub(crate) fn new_console_spawn_sink_slot() -> SharedConsoleSpawnSinkSlot {
Arc::new(std::sync::RwLock::new(None))
}
pub(crate) const SPAWNED_BY_LABEL: &str = "spawned_by";
pub(crate) const VIA_TOOL_LABEL: &str = "via_tool";
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ConsoleSpawnSeed {
pub(crate) mob_id: Option<String>,
pub(crate) member_id: String,
pub(crate) identity: String,
pub(crate) initial_message: Option<Value>,
pub(crate) labels: BTreeMap<String, String>,
pub(crate) spawned_by: Option<String>,
pub(crate) via_tool: String,
}
#[derive(Clone)]
pub(crate) struct ConsoleSpawnSink {
console_events: ConsoleEventStore,
}
impl ConsoleSpawnSink {
pub(crate) fn new(console_events: ConsoleEventStore) -> Self {
Self { console_events }
}
pub(crate) async fn identity_labels_snapshot(
&self,
) -> BTreeMap<String, BTreeMap<String, String>> {
self.console_events.identity_labels_snapshot().await
}
pub(crate) async fn project_spawned_member(&self, seed: &ConsoleSpawnSeed) {
let identity = seed.identity.trim();
let member_id = seed.member_id.trim();
if identity.is_empty() || member_id.is_empty() {
return;
}
self.console_events
.register_runtime_identity(member_id, identity)
.await;
self.console_events
.register_runtime_identity(format!("rt:{member_id}"), identity)
.await;
let mut labels = seed.labels.clone();
match seed.spawned_by.as_deref() {
Some(spawned_by) => {
labels.insert(SPAWNED_BY_LABEL.to_string(), spawned_by.to_string());
}
None => {
labels.remove(SPAWNED_BY_LABEL);
self.console_events
.unregister_identity_label(identity, SPAWNED_BY_LABEL)
.await;
}
}
labels.insert(VIA_TOOL_LABEL.to_string(), seed.via_tool.clone());
self.console_events
.register_identity_labels(identity, labels)
.await;
let Some(initial_message) = seed.initial_message.as_ref() else {
return;
};
let content = match initial_message {
Value::String(text) => json!([{ "type": "text", "text": text }]),
other => other.clone(),
};
let mut data = json!({
"content": content,
"message": { "role": "user", "content": initial_message },
"source_event_type": "spawn_initial_message",
"via_tool": seed.via_tool,
"member_id": member_id,
});
if let Some(object) = data.as_object_mut() {
if let Some(mob_id) = seed.mob_id.as_deref().filter(|id| !id.trim().is_empty()) {
object.insert("mob_id".to_string(), json!(mob_id));
}
if let Some(parent) = seed.spawned_by.as_deref() {
object.insert("parent_identity".to_string(), json!(parent));
}
}
self.console_events
.append_envelope(crate::console_contracts::ConsoleIdentityEventEnvelope {
event_id: kickoff_event_id(seed.mob_id.as_deref(), member_id, initial_message),
interaction_id: None,
identity: identity.to_string(),
event_type: "user_input".to_string(),
timestamp_ms: current_time_ms(),
data,
})
.await;
}
}
pub(crate) fn is_console_spawn_tool(name: &str) -> bool {
matches!(
name,
"mob_spawn_member" | "spawn_member" | "spawn_many_members" | "delegate"
)
}
pub(crate) fn spawned_by_from_comms_name(comms_name: &str) -> Option<String> {
let comms_name = comms_name.trim();
if comms_name.is_empty() {
return None;
}
if !comms_name.contains('/') {
return Some(comms_name.to_string());
}
let mut parts = comms_name.splitn(3, '/');
let (_mob, _role) = (parts.next()?, parts.next()?);
let identity = parts.next()?.trim();
(!identity.is_empty()).then(|| identity.to_string())
}
pub(crate) fn sanitize_unverified_lineage_labels(labels: &mut BTreeMap<String, String>) {
labels.remove(SPAWNED_BY_LABEL);
labels.remove(VIA_TOOL_LABEL);
}
struct SpawnArgRecord {
member_id: Option<String>,
mob_id: Option<String>,
initial_message: Option<Value>,
labels: BTreeMap<String, String>,
}
fn text_field(value: &Value, key: &str) -> Option<String> {
value
.get(key)
.and_then(Value::as_str)
.map(str::trim)
.filter(|text| !text.is_empty())
.map(ToString::to_string)
}
fn member_id_field(value: &Value) -> Option<String> {
text_field(value, "member_id")
.or_else(|| text_field(value, "agent_identity"))
.or_else(|| text_field(value, "identity"))
}
fn labels_field(value: &Value) -> BTreeMap<String, String> {
value
.get("labels")
.and_then(Value::as_object)
.map(|labels| {
labels
.iter()
.filter_map(|(key, value)| {
value.as_str().map(|text| (key.clone(), text.to_string()))
})
.collect()
})
.unwrap_or_default()
}
fn spawn_arg_record(value: &Value, default_mob_id: Option<&str>) -> Option<SpawnArgRecord> {
if !value.is_object() {
return None;
}
let member_id = member_id_field(value);
let initial_message = value
.get("initial_message")
.or_else(|| value.get("task"))
.filter(|message| !message.is_null())
.cloned();
if member_id.is_none() && initial_message.is_none() {
return None;
}
let mut labels = labels_field(value);
if let Some(display_name) = text_field(value, "display_name") {
labels
.entry("display_name".to_string())
.or_insert(display_name);
}
Some(SpawnArgRecord {
member_id,
mob_id: text_field(value, "mob_id").or_else(|| default_mob_id.map(ToString::to_string)),
initial_message,
labels,
})
}
fn spawn_arg_records(args: &Value) -> Vec<SpawnArgRecord> {
let default_mob_id = text_field(args, "mob_id");
let mut records = Vec::new();
for key in ["specs", "members"] {
let Some(values) = args.get(key).and_then(Value::as_array) else {
continue;
};
for value in values {
if let Some(record) = spawn_arg_record(value, default_mob_id.as_deref()) {
records.push(record);
}
}
}
if records.is_empty()
&& let Some(record) = spawn_arg_record(args, default_mob_id.as_deref())
{
records.push(record);
}
records
}
struct SpawnOutcomeTarget {
member_id: String,
mob_id: Option<String>,
}
fn collect_outcome_targets(
value: &Value,
default_mob_id: Option<&str>,
targets: &mut Vec<SpawnOutcomeTarget>,
) {
if let Some(member_id) = member_id_field(value)
&& !targets.iter().any(|target| target.member_id == member_id)
{
targets.push(SpawnOutcomeTarget {
member_id,
mob_id: text_field(value, "mob_id").or_else(|| default_mob_id.map(ToString::to_string)),
});
}
for key in ["members", "specs", "spawned", "results"] {
let Some(values) = value.get(key).and_then(Value::as_array) else {
continue;
};
let nested_default = text_field(value, "mob_id");
for nested in values {
collect_outcome_targets(
nested,
nested_default.as_deref().or(default_mob_id),
targets,
);
}
}
}
pub(crate) fn console_spawn_seeds(
tool_name: &str,
args: &Value,
outcome_text: &str,
spawner_comms_name: Option<&str>,
) -> Vec<ConsoleSpawnSeed> {
if !is_console_spawn_tool(tool_name) {
return Vec::new();
}
let spawned_by = spawner_comms_name.and_then(spawned_by_from_comms_name);
let arg_records = spawn_arg_records(args);
let mut outcome_targets = Vec::new();
if let Ok(outcome) = serde_json::from_str::<Value>(outcome_text) {
collect_outcome_targets(&outcome, None, &mut outcome_targets);
}
let seed_for = |member_id: String,
mob_id: Option<String>,
record: Option<&SpawnArgRecord>|
-> ConsoleSpawnSeed {
let labels = record
.map(|record| record.labels.clone())
.unwrap_or_default();
let identity = labels
.get("agent_identity")
.map(|value| value.trim())
.filter(|value| !value.is_empty())
.map(ToString::to_string)
.unwrap_or_else(|| member_id.clone());
ConsoleSpawnSeed {
mob_id: mob_id
.or_else(|| record.and_then(|record| record.mob_id.clone()))
.or_else(|| text_field(args, "mob_id")),
member_id,
identity,
initial_message: record.and_then(|record| record.initial_message.clone()),
labels,
spawned_by: spawned_by.clone(),
via_tool: tool_name.to_string(),
}
};
if !outcome_targets.is_empty() {
return outcome_targets
.into_iter()
.map(|target| {
let record = arg_records
.iter()
.find(|record| record.member_id.as_deref() == Some(target.member_id.as_str()))
.or_else(|| match arg_records.as_slice() {
[only] if only.member_id.is_none() => Some(only),
_ => None,
});
seed_for(target.member_id, target.mob_id, record)
})
.collect();
}
arg_records
.iter()
.filter_map(|record| {
let member_id = record.member_id.clone()?;
Some(seed_for(member_id, record.mob_id.clone(), Some(record)))
})
.collect()
}
fn hash_short(input: &str) -> String {
use std::fmt::Write as _;
let mut hasher = Sha256::new();
hasher.update(input.as_bytes());
let digest = hasher.finalize();
digest[..8]
.iter()
.fold(String::with_capacity(16), |mut out, byte| {
let _ = write!(out, "{byte:02x}");
out
})
}
pub(crate) fn merge_registered_labels(
labels: &mut BTreeMap<String, String>,
registered: &BTreeMap<String, String>,
) {
for (key, value) in registered {
if key == SPAWNED_BY_LABEL || key == VIA_TOOL_LABEL {
labels.insert(key.clone(), value.clone());
} else {
labels.entry(key.clone()).or_insert_with(|| value.clone());
}
}
}
fn kickoff_event_id(mob_id: Option<&str>, member_id: &str, initial_message: &Value) -> String {
format!(
"spawn-kickoff:{}:{}:{}",
mob_id.unwrap_or("-"),
member_id,
hash_short(&initial_message.to_string())
)
}
fn current_time_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or_default()
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::types::{EventEnvelope, UnifiedEvent};
fn seed(member_id: &str, initial_message: Option<Value>) -> ConsoleSpawnSeed {
ConsoleSpawnSeed {
mob_id: Some("ob3".to_string()),
member_id: member_id.to_string(),
identity: member_id.to_string(),
initial_message,
labels: BTreeMap::from([("group".to_string(), "workers".to_string())]),
spawned_by: Some("ops-lead".to_string()),
via_tool: "mob_spawn_member".to_string(),
}
}
#[tokio::test]
async fn projects_kickoff_event_under_member_identity() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
sink.project_spawned_member(&seed("worker-3", Some(json!("Find the person"))))
.await;
let replay = store.replay_all(None).await.expect("replay");
let kickoff = replay
.iter()
.find(|event| event.event_type == "user_input")
.expect("kickoff user_input event projected");
assert_eq!(kickoff.identity, "worker-3");
assert!(
kickoff.event_id.starts_with("spawn-kickoff:ob3:worker-3:"),
"deterministic kickoff id, got {}",
kickoff.event_id
);
assert_eq!(kickoff.data["content"][0]["type"], "text");
assert_eq!(kickoff.data["content"][0]["text"], "Find the person");
assert_eq!(kickoff.data["message"]["role"], "user");
assert_eq!(kickoff.data["source_event_type"], "spawn_initial_message");
assert_eq!(kickoff.data["via_tool"], "mob_spawn_member");
assert_eq!(kickoff.data["parent_identity"], "ops-lead");
}
#[tokio::test]
async fn double_projection_dedupes_to_one_kickoff() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
let seed = seed("worker-3", Some(json!("Find the person")));
sink.project_spawned_member(&seed).await;
sink.project_spawned_member(&seed).await;
let replay = store.replay_all(None).await.expect("replay");
let kickoffs = replay
.iter()
.filter(|event| event.event_type == "user_input")
.count();
assert_eq!(kickoffs, 1, "spawn retries must not duplicate the kickoff");
}
#[tokio::test]
async fn no_initial_message_registers_identity_without_chat_frames() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
sink.project_spawned_member(&seed("worker-9", None)).await;
let replay = store.replay_all(None).await.expect("replay");
assert!(
replay.iter().all(|event| event.identity != "worker-9"),
"no initial_message means no chat frames until first activity"
);
let labels = store
.identity_labels("worker-9")
.await
.expect("identity labels registered");
assert_eq!(labels.get("group").map(String::as_str), Some("workers"));
assert_eq!(
labels.get(SPAWNED_BY_LABEL).map(String::as_str),
Some("ops-lead")
);
assert_eq!(
labels.get(VIA_TOOL_LABEL).map(String::as_str),
Some("mob_spawn_member")
);
}
#[tokio::test]
async fn registers_runtime_identity_so_live_events_join_the_same_chat() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
sink.project_spawned_member(&seed("worker-3", Some(json!("go"))))
.await;
store
.project_unified_event(&EventEnvelope {
event_id: "evt-live-1".to_string(),
source: "test".to_string(),
timestamp_ms: 7,
event: UnifiedEvent::Agent {
agent_id: "worker-3:0:1".to_string(),
event_type: "text_delta".to_string(),
payload: Some(json!({ "delta": "on it" })),
},
})
.await;
let replay = store.replay_all(None).await.expect("replay");
let live = replay
.iter()
.find(|event| event.event_id == "evt-live-1")
.expect("live event projected");
assert_eq!(
live.identity, "worker-3",
"live runtime events must land in the same chat as the kickoff"
);
}
#[tokio::test]
async fn caller_labels_cannot_spoof_runtime_derived_lineage() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
let mut spoofed = seed("worker-3", None);
spoofed
.labels
.insert(SPAWNED_BY_LABEL.to_string(), "victim-agent".to_string());
spoofed
.labels
.insert(VIA_TOOL_LABEL.to_string(), "console".to_string());
sink.project_spawned_member(&spoofed).await;
let labels = store
.identity_labels("worker-3")
.await
.expect("identity labels registered");
assert_eq!(
labels.get(SPAWNED_BY_LABEL).map(String::as_str),
Some("ops-lead"),
"ABAC lineage must record the actual spawner, not a caller-supplied label"
);
assert_eq!(
labels.get(VIA_TOOL_LABEL).map(String::as_str),
Some("mob_spawn_member")
);
let unknown_parent = ConsoleSpawnSeed {
spawned_by: None,
labels: BTreeMap::from([(SPAWNED_BY_LABEL.to_string(), "victim-agent".to_string())]),
..seed("worker-9", None)
};
sink.project_spawned_member(&unknown_parent).await;
let labels = store
.identity_labels("worker-9")
.await
.expect("identity labels registered");
assert!(
!labels.contains_key(SPAWNED_BY_LABEL),
"unknown spawner must not fall back to a caller-supplied lineage claim"
);
}
#[test]
fn spawn_tools_are_recognized() {
for tool in [
"mob_spawn_member",
"spawn_member",
"spawn_many_members",
"delegate",
] {
assert!(is_console_spawn_tool(tool), "{tool} spawns members");
}
assert!(!is_console_spawn_tool("mob_retire_member"));
assert!(!is_console_spawn_tool("send_message"));
}
#[test]
fn spawned_by_uses_comms_name_identity_segment() {
assert_eq!(
spawned_by_from_comms_name("ob3/orchestrator/ops-lead").as_deref(),
Some("ops-lead")
);
assert_eq!(
spawned_by_from_comms_name("ops-lead").as_deref(),
Some("ops-lead")
);
assert_eq!(
spawned_by_from_comms_name("ob3/orchestrator/people/finder").as_deref(),
Some("people/finder")
);
assert_eq!(spawned_by_from_comms_name("mob/role"), None);
assert_eq!(spawned_by_from_comms_name(""), None);
assert_eq!(spawned_by_from_comms_name("mob/role/"), None);
}
#[tokio::test]
async fn respawn_with_unknown_spawner_clears_stale_lineage() {
let store = ConsoleEventStore::new();
let sink = ConsoleSpawnSink::new(store.clone());
sink.project_spawned_member(&seed("worker-3", None)).await;
assert_eq!(
store
.identity_labels("worker-3")
.await
.and_then(|labels| labels.get(SPAWNED_BY_LABEL).cloned())
.as_deref(),
Some("ops-lead")
);
let unknown_spawner = ConsoleSpawnSeed {
spawned_by: None,
via_tool: "spawn_member".to_string(),
..seed("worker-3", None)
};
sink.project_spawned_member(&unknown_spawner).await;
let labels = store
.identity_labels("worker-3")
.await
.expect("identity labels registered");
assert!(
!labels.contains_key(SPAWNED_BY_LABEL),
"a respawn with an unknown spawner must not keep stale lineage alive"
);
assert_eq!(
labels.get(VIA_TOOL_LABEL).map(String::as_str),
Some("spawn_member")
);
}
#[test]
fn seeds_from_single_spawn_args_and_outcome() {
let args = json!({
"mob_id": "ob3",
"profile": "person-worker",
"member_id": "worker-3",
"initial_message": "Find the person",
"labels": { "group": "workers", "display_name": "Worker 3" }
});
let outcome = json!({
"agent_identity": "worker-3",
"member_ref": "opaque-ref"
});
let seeds = console_spawn_seeds(
"mob_spawn_member",
&args,
&outcome.to_string(),
Some("ob3/orchestrator/ops-lead"),
);
assert_eq!(seeds.len(), 1);
let seed = &seeds[0];
assert_eq!(seed.mob_id.as_deref(), Some("ob3"));
assert_eq!(seed.member_id, "worker-3");
assert_eq!(seed.identity, "worker-3");
assert_eq!(seed.initial_message, Some(json!("Find the person")));
assert_eq!(
seed.labels.get("group").map(String::as_str),
Some("workers")
);
assert_eq!(
seed.labels.get("display_name").map(String::as_str),
Some("Worker 3")
);
assert_eq!(seed.spawned_by.as_deref(), Some("ops-lead"));
assert_eq!(seed.via_tool, "mob_spawn_member");
}
#[test]
fn seeds_honor_agent_identity_label_override() {
let args = json!({
"mob_id": "ob3",
"member_id": "worker-3",
"initial_message": "go",
"labels": { "agent_identity": "people/finder" }
});
let outcome = json!({ "agent_identity": "worker-3" });
let seeds = console_spawn_seeds("mob_spawn_member", &args, &outcome.to_string(), None);
assert_eq!(seeds.len(), 1);
assert_eq!(seeds[0].member_id, "worker-3");
assert_eq!(
seeds[0].identity, "people/finder",
"kickoff must land under the same identity the event projection derives"
);
}
#[test]
fn seeds_from_specs_array_enrich_each_member() {
let args = json!({
"mob_id": "ob3",
"specs": [
{
"profile": "person-worker",
"member_id": "worker-a",
"initial_message": "task a",
"labels": { "group": "workers" }
},
{
"profile": "person-worker",
"member_id": "worker-b",
"initial_message": "task b"
}
]
});
let outcome = json!({
"members": [
{ "agent_identity": "worker-a" },
{ "agent_identity": "worker-b" }
]
});
let seeds = console_spawn_seeds("spawn_many_members", &args, &outcome.to_string(), None);
assert_eq!(seeds.len(), 2);
let worker_a = seeds
.iter()
.find(|seed| seed.member_id == "worker-a")
.unwrap();
assert_eq!(worker_a.initial_message, Some(json!("task a")));
assert_eq!(
worker_a.labels.get("group").map(String::as_str),
Some("workers")
);
let worker_b = seeds
.iter()
.find(|seed| seed.member_id == "worker-b")
.unwrap();
assert_eq!(worker_b.initial_message, Some(json!("task b")));
}
#[test]
fn delegate_seed_pairs_generated_member_with_task() {
let args = json!({ "task": "Review the diff" });
let outcome = json!({
"agent_identity": "helper-3f2a",
"member_ref": "opaque",
"mob_id": "implicit-1",
"wired": true
});
let seeds = console_spawn_seeds(
"delegate",
&args,
&outcome.to_string(),
Some("ob3/orchestrator/ops-lead"),
);
assert_eq!(seeds.len(), 1);
let seed = &seeds[0];
assert_eq!(seed.member_id, "helper-3f2a");
assert_eq!(seed.mob_id.as_deref(), Some("implicit-1"));
assert_eq!(seed.initial_message, Some(json!("Review the diff")));
assert_eq!(seed.spawned_by.as_deref(), Some("ops-lead"));
assert_eq!(seed.via_tool, "delegate");
}
#[test]
fn unparsable_outcome_falls_back_to_args_members() {
let args = json!({
"mob_id": "ob3",
"member_id": "worker-3",
"initial_message": "go"
});
let seeds = console_spawn_seeds("mob_spawn_member", &args, "spawned ok", None);
assert_eq!(seeds.len(), 1);
assert_eq!(seeds[0].member_id, "worker-3");
}
#[test]
fn merge_registered_labels_fills_gaps_without_overriding() {
let mut labels = BTreeMap::from([
("group".to_string(), "roster-group".to_string()),
("role".to_string(), "worker".to_string()),
]);
let registered = BTreeMap::from([
("group".to_string(), "registry-group".to_string()),
(SPAWNED_BY_LABEL.to_string(), "ops-lead".to_string()),
]);
merge_registered_labels(&mut labels, ®istered);
assert_eq!(
labels.get("group").map(String::as_str),
Some("roster-group")
);
assert_eq!(
labels.get(SPAWNED_BY_LABEL).map(String::as_str),
Some("ops-lead")
);
assert_eq!(labels.get("role").map(String::as_str), Some("worker"));
}
#[test]
fn merge_registered_labels_keeps_lineage_keys_runtime_derived() {
let mut labels = BTreeMap::from([
(SPAWNED_BY_LABEL.to_string(), "victim-agent".to_string()),
(VIA_TOOL_LABEL.to_string(), "console".to_string()),
]);
let registered = BTreeMap::from([
(SPAWNED_BY_LABEL.to_string(), "ops-lead".to_string()),
(VIA_TOOL_LABEL.to_string(), "mob_spawn_member".to_string()),
]);
merge_registered_labels(&mut labels, ®istered);
assert_eq!(
labels.get(SPAWNED_BY_LABEL).map(String::as_str),
Some("ops-lead"),
"registry lineage must displace spec-supplied spawned_by claims"
);
assert_eq!(
labels.get(VIA_TOOL_LABEL).map(String::as_str),
Some("mob_spawn_member")
);
}
}