use super::*;
use serde_json::json;
fn sample_lifecycle() -> LifecycleEvent {
LifecycleEvent::PmThinking {
session_id: "s1".into(),
text: "considering options".into(),
}
}
fn envelope(payload: HarnessPayload, session: Option<&str>) -> HarnessEvent {
HarnessEvent {
source: HarnessSource::Agents,
session: session.map(str::to_string),
seq: 0,
at: chrono::Utc::now(),
payload,
id: EventId::new(),
parent_id: None,
}
}
#[test]
fn harness_source_round_trips() {
for (src, tag) in [
(HarnessSource::Agents, "\"agents\""),
(HarnessSource::Mpm, "\"mpm\""),
(HarnessSource::Code, "\"code\""),
] {
let s = serde_json::to_string(&src).expect("serialize source");
assert_eq!(s, tag);
let back: HarnessSource = serde_json::from_str(&s).expect("deserialize source");
assert_eq!(back, src);
}
}
#[test]
fn lifecycle_event_serializes_with_type_tag() {
let s = serde_json::to_string(&sample_lifecycle()).expect("serialize");
assert!(s.contains("\"type\":\"pm_thinking\""), "{s}");
assert!(s.contains("\"session_id\":\"s1\""), "{s}");
}
#[test]
fn lifecycle_session_id_returns_correct_field() {
let ev = LifecycleEvent::AgentMessage {
session_id: "abc".into(),
agent: "python".into(),
text: "hi".into(),
};
assert_eq!(ev.session_id(), Some("abc"));
}
#[test]
fn lifecycle_recap_round_trips() {
let ev = LifecycleEvent::RecapGenerated {
session_id: "s9".into(),
summary: "did a thing".into(),
table_rows: vec![("step".into(), "ok".into())],
};
let s = serde_json::to_string(&ev).expect("serialize recap");
let back: LifecycleEvent = serde_json::from_str(&s).expect("deserialize recap");
assert_eq!(back, ev);
}
#[test]
fn payload_lifecycle_round_trips() {
let p = HarnessPayload::Lifecycle(sample_lifecycle());
let s = serde_json::to_string(&p).expect("serialize");
assert!(s.contains("\"domain\":\"lifecycle\""), "{s}");
assert!(s.contains("\"event\":{"), "{s}");
assert!(s.contains("\"type\":\"pm_thinking\""), "{s}");
let back: HarnessPayload = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, p);
}
#[test]
fn payload_hook_round_trips() {
let p = HarnessPayload::Hook {
kind: "pre_tool_use".into(),
data: json!({"tool": "bash", "ok": true}),
};
let s = serde_json::to_string(&p).expect("serialize");
assert!(s.contains("\"domain\":\"hook\""), "{s}");
assert!(s.contains("\"kind\":\"pre_tool_use\""), "{s}");
let back: HarnessPayload = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, p);
}
#[test]
fn payload_ping_round_trips() {
let p = HarnessPayload::Ping;
let s = serde_json::to_string(&p).expect("serialize");
assert_eq!(s, "{\"domain\":\"ping\"}");
let back: HarnessPayload = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, p);
}
#[test]
fn payload_domain_matches_serde_tag() {
assert_eq!(
HarnessPayload::Lifecycle(sample_lifecycle()).domain(),
"lifecycle"
);
assert_eq!(
HarnessPayload::Hook {
kind: "x".into(),
data: json!(null)
}
.domain(),
"hook"
);
assert_eq!(HarnessPayload::Ping.domain(), "ping");
assert_eq!(
HarnessPayload::Action(sample_action_event()).domain(),
"action"
);
}
#[test]
fn harness_payload_pre_action_payload_still_deserializes() {
let json = r#"{"domain":"hook","event":{"kind":"pre_tool_use","data":{"tool":"bash"}}}"#;
let back: HarnessPayload = serde_json::from_str(json).expect("legacy hook payload parses");
assert_eq!(
back,
HarnessPayload::Hook {
kind: "pre_tool_use".into(),
data: json!({"tool": "bash"}),
}
);
}
#[test]
fn harness_payload_action_round_trips() {
let p = HarnessPayload::Action(sample_action_event());
let s = serde_json::to_string(&p).expect("serialize");
assert!(s.contains("\"domain\":\"action\""), "{s}");
assert!(s.contains("\"kind\":\"session\""), "{s}");
let back: HarnessPayload = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, p);
}
fn sample_meta() -> ActionMeta {
ActionMeta {
id: EventId::new(),
at: chrono::Utc::now(),
source: HarnessSource::Mpm,
session: Some("s1".into()),
parent_id: None,
actor: Actor::Agent {
name: "rust-engineer".into(),
agent_id: "agent-1".into(),
},
objects: vec![ObjectRef {
object_type: ObjectType::Session,
id: "s1".into(),
label: "session s1".into(),
}],
schema_version: 1,
}
}
fn sample_action_event() -> ActionEvent {
ActionEvent::Session {
meta: sample_meta(),
phase: SessionPhase::Started,
}
}
#[test]
fn action_event_round_trips_all_six_kinds() {
let meta = sample_meta();
let events = vec![
ActionEvent::Workflow {
meta: meta.clone(),
phase: WorkflowPhase::Spawn,
object: ObjectRef {
object_type: ObjectType::Task,
id: "t1".into(),
label: "build the thing".into(),
},
},
ActionEvent::Agent {
meta: meta.clone(),
phase: AgentPhase::Spawned,
agent_id: "agent-1".into(),
},
ActionEvent::File {
meta: meta.clone(),
phase: FilePhase::Written,
path: PathRef {
path: "src/lib.rs".into(),
diff_ref: Some("diff-1".into()),
},
},
ActionEvent::Tool {
meta: meta.clone(),
phase: CallPhase::Finished,
tool: "bash".into(),
call_id: "call-1".into(),
},
ActionEvent::Session {
meta: meta.clone(),
phase: SessionPhase::Done,
},
ActionEvent::Inference {
meta: meta.clone(),
phase: CallPhase::Started,
model: "claude-sonnet".into(),
},
];
for event in events {
let s = serde_json::to_string(&event).expect("serialize");
let back: ActionEvent = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, event, "round trip for {}: {s}", event.kind());
assert_eq!(back.meta(), event.meta());
}
}
#[test]
fn action_event_wire_shape_matches_kind_tag() {
let event = ActionEvent::Tool {
meta: sample_meta(),
phase: CallPhase::Started,
tool: "bash".into(),
call_id: "call-1".into(),
};
let s = serde_json::to_string(&event).expect("serialize");
assert!(s.contains("\"kind\":\"tool\""), "{s}");
assert!(s.contains("\"tool\":\"bash\""), "{s}");
assert!(s.contains("\"call_id\":\"call-1\""), "{s}");
assert!(
s.contains("\"actor\":{\"type\":\"agent\""),
"meta fields are flattened alongside the variant fields: {s}"
);
}
#[test]
fn action_event_kind_matches_serde_tag() {
let meta = sample_meta();
let cases: Vec<(ActionEvent, &str)> = vec![
(
ActionEvent::Workflow {
meta: meta.clone(),
phase: WorkflowPhase::Start,
object: ObjectRef {
object_type: ObjectType::Task,
id: "t1".into(),
label: "l".into(),
},
},
"workflow",
),
(
ActionEvent::Agent {
meta: meta.clone(),
phase: AgentPhase::Done,
agent_id: "a1".into(),
},
"agent",
),
(
ActionEvent::File {
meta: meta.clone(),
phase: FilePhase::Created,
path: PathRef {
path: "x".into(),
diff_ref: None,
},
},
"file",
),
(
ActionEvent::Tool {
meta: meta.clone(),
phase: CallPhase::Errored,
tool: "t".into(),
call_id: "c1".into(),
},
"tool",
),
(
ActionEvent::Session {
meta: meta.clone(),
phase: SessionPhase::Cancelled,
},
"session",
),
(
ActionEvent::Inference {
meta: meta.clone(),
phase: CallPhase::Finished,
model: "m".into(),
},
"inference",
),
];
for (event, expected) in cases {
let s = serde_json::to_string(&event).expect("serialize");
assert_eq!(event.kind(), expected);
assert!(s.contains(&format!("\"kind\":\"{expected}\"")), "{s}");
}
}
#[test]
fn action_meta_schema_version_defaults_to_one() {
let json = r#"{
"kind": "session",
"id": "018f1e0a-0000-7000-8000-000000000000",
"at": "2026-01-01T00:00:00Z",
"source": "mpm",
"actor": {"type": "system"},
"phase": "started"
}"#;
let event: ActionEvent =
serde_json::from_str(json).expect("deserialize without schema_version");
assert_eq!(event.meta().schema_version, 1);
}
#[test]
fn actor_operator_and_system_round_trip() {
for (actor, tag) in [
(Actor::Operator, "\"type\":\"operator\""),
(Actor::System, "\"type\":\"system\""),
] {
let s = serde_json::to_string(&actor).expect("serialize");
assert_eq!(s, format!("{{{tag}}}"));
let back: Actor = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, actor);
}
}
#[test]
fn object_type_round_trips_every_variant() {
for (ty, tag) in [
(ObjectType::Session, "\"session\""),
(ObjectType::Agent, "\"agent\""),
(ObjectType::Task, "\"task\""),
(ObjectType::Workstream, "\"workstream\""),
(ObjectType::File, "\"file\""),
(ObjectType::ToolCall, "\"tool_call\""),
(ObjectType::Inference, "\"inference\""),
(ObjectType::Issue, "\"issue\""),
(ObjectType::Pr, "\"pr\""),
] {
let s = serde_json::to_string(&ty).expect("serialize");
assert_eq!(s, tag);
let back: ObjectType = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, ty);
}
}
#[test]
fn harness_event_round_trips() {
let ev = envelope(HarnessPayload::Ping, Some("sess-1"));
let s = serde_json::to_string(&ev).expect("serialize");
assert!(s.contains("\"source\":\"agents\""), "{s}");
assert!(s.contains("\"session\":\"sess-1\""), "{s}");
assert!(s.contains("\"id\":\""), "id should be present: {s}");
let back: HarnessEvent = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, ev);
assert_eq!(back.id, ev.id);
}
#[test]
fn harness_event_omits_none_session() {
let ev = envelope(HarnessPayload::Ping, None);
let s = serde_json::to_string(&ev).expect("serialize");
assert!(!s.contains("session"), "session should be omitted: {s}");
}
#[test]
fn harness_event_omits_none_parent_id() {
let ev = envelope(HarnessPayload::Ping, None);
let s = serde_json::to_string(&ev).expect("serialize");
assert!(
!s.contains("parent_id"),
"parent_id should be omitted when None: {s}"
);
}
#[test]
fn harness_event_id_is_unique_per_event() {
let a = envelope(HarnessPayload::Ping, None);
let b = envelope(HarnessPayload::Ping, None);
assert_ne!(a.id, b.id);
}
#[test]
fn harness_event_parent_id_links_to_the_causing_event() {
let root = envelope(HarnessPayload::Ping, Some("s1"));
let mut child = envelope(HarnessPayload::Ping, Some("s1"));
child.parent_id = Some(root.id);
let s = serde_json::to_string(&child).expect("serialize");
assert!(s.contains(&format!("\"parent_id\":\"{}\"", root.id)), "{s}");
let back: HarnessEvent = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back.parent_id, Some(root.id));
assert_ne!(
back.parent_id,
Some(child.id),
"child is not its own parent"
);
}
#[test]
fn harness_event_back_compat_missing_fields_deserializes() {
let ev = envelope(HarnessPayload::Ping, Some("legacy"));
let mut value = serde_json::to_value(&ev).expect("serialize to value");
let obj = value
.as_object_mut()
.expect("envelope serializes as an object");
obj.remove("id");
obj.remove("parent_id");
let json = serde_json::to_string(&value).expect("serialize legacy shape");
let back: HarnessEvent =
serde_json::from_str(&json).expect("legacy payload without id/parent_id should parse");
assert_eq!(back.source, ev.source);
assert_eq!(back.session, ev.session);
assert_eq!(back.seq, ev.seq);
assert!(back.parent_id.is_none());
}
#[test]
fn harness_event_missing_id_mints_a_fresh_id_each_deserialize() {
let ev = envelope(HarnessPayload::Ping, None);
let mut value = serde_json::to_value(&ev).expect("serialize to value");
value
.as_object_mut()
.expect("envelope serializes as an object")
.remove("id");
let json = serde_json::to_string(&value).expect("serialize legacy shape");
let a: HarnessEvent = serde_json::from_str(&json).expect("first deserialize");
let b: HarnessEvent = serde_json::from_str(&json).expect("second deserialize");
assert_ne!(a.id, b.id);
}
#[test]
fn event_id_round_trips() {
let id = EventId::new();
let s = serde_json::to_string(&id).expect("serialize");
let back: EventId = serde_json::from_str(&s).expect("deserialize");
assert_eq!(back, id);
}
#[test]
fn event_id_new_mints_distinct_ids() {
let a = EventId::new();
let b = EventId::new();
assert_ne!(a, b);
}
#[test]
fn event_id_display_matches_serialized_string() {
let id = EventId::new();
let json = serde_json::to_string(&id).expect("serialize");
let quoted = format!("\"{id}\"");
assert_eq!(json, quoted);
}
#[test]
fn filter_default_matches_all() {
let f = Filter::default();
assert!(f.matches(&envelope(HarnessPayload::Ping, None)));
assert!(f.matches(&envelope(
HarnessPayload::Lifecycle(sample_lifecycle()),
Some("x")
)));
}
#[test]
fn filter_by_source() {
let f = Filter {
source: Some(HarnessSource::Mpm),
..Default::default()
};
let mut ev = envelope(HarnessPayload::Ping, None);
ev.source = HarnessSource::Mpm;
assert!(f.matches(&ev));
ev.source = HarnessSource::Agents;
assert!(!f.matches(&ev));
}
#[test]
fn filter_by_session() {
let f = Filter {
session: Some("sess-7".into()),
..Default::default()
};
assert!(f.matches(&envelope(HarnessPayload::Ping, Some("sess-7"))));
assert!(!f.matches(&envelope(HarnessPayload::Ping, Some("other"))));
assert!(!f.matches(&envelope(HarnessPayload::Ping, None)));
}
#[test]
fn filter_by_domain() {
let f = Filter {
domains: Some(vec!["hook", "ping"]),
..Default::default()
};
assert!(f.matches(&envelope(HarnessPayload::Ping, None)));
assert!(f.matches(&envelope(
HarnessPayload::Hook {
kind: "k".into(),
data: json!({})
},
None
)));
assert!(!f.matches(&envelope(
HarnessPayload::Lifecycle(sample_lifecycle()),
Some("x")
)));
}
#[test]
fn filter_combination() {
let f = Filter {
source: Some(HarnessSource::Code),
session: Some("s".into()),
domains: Some(vec!["lifecycle"]),
};
let mut ev = envelope(HarnessPayload::Lifecycle(sample_lifecycle()), Some("s"));
ev.source = HarnessSource::Code;
assert!(f.matches(&ev));
ev.source = HarnessSource::Mpm;
assert!(!f.matches(&ev));
}
const MODULE_SOURCES: &[(&str, &str)] = &[
("control_bus/mod.rs", include_str!("mod.rs")),
("control_bus/action.rs", include_str!("action.rs")),
("control_bus/lifecycle.rs", include_str!("lifecycle.rs")),
("control_bus/envelope.rs", include_str!("envelope.rs")),
("control_bus/event_id.rs", include_str!("event_id.rs")),
("control_bus/filter.rs", include_str!("filter.rs")),
("control_bus/push_client.rs", include_str!("push_client.rs")),
("control_bus/tests.rs", include_str!("tests.rs")),
];
const FORBIDDEN_SUBSTRINGS: &[&str] = &[
concat!("broadcast", "::"),
concat!("Once", "Lock"),
concat!("tokio", "::", "sync"),
concat!("lazy_", "static!"),
concat!("once_", "cell"),
];
fn strip_leading_pub(trimmed: &str) -> &str {
let Some(rest) = trimmed.strip_prefix("pub") else {
return trimmed;
};
if let Some(after_paren) = rest.strip_prefix('(') {
match after_paren.find(')') {
Some(close) => after_paren[close + 1..].trim_start(),
None => rest, }
} else if rest.starts_with(char::is_whitespace) || rest.is_empty() {
rest.trim_start()
} else {
trimmed
}
}
#[test]
fn control_bus_declares_no_transport() {
for (name, src) in MODULE_SOURCES {
for needle in FORBIDDEN_SUBSTRINGS {
assert!(
!src.contains(needle),
"{name} contains `{needle}`: control_bus holds event TYPES only \
— the bus lives in trusty-console (#6846)"
);
}
for (idx, line) in src.lines().enumerate() {
let candidate = strip_leading_pub(line.trim_start());
assert!(
!(candidate.starts_with("static ") || candidate.starts_with("static mut ")),
"{}:{} declares a global `static`: control_bus holds event TYPES \
only — no global state (#6846)\n {line}",
name,
idx + 1
);
}
}
}
#[test]
fn static_scan_catches_every_visibility_form() {
let is_static_declaration = |line: &str| -> bool {
let candidate = strip_leading_pub(line.trim_start());
candidate.starts_with("static ") || candidate.starts_with("static mut ")
};
assert!(
is_static_declaration("pub(crate) static X: u8 = 0;"),
"`pub(crate) static` must be detected as a static declaration"
);
assert!(
is_static_declaration("static mut Y: u8 = 0;"),
"`static mut` with no visibility qualifier must be detected"
);
assert!(
!is_static_declaration("let s: &'static str = \"\";"),
"`&'static` in a type position must not be flagged as a static declaration"
);
}
#[test]
fn module_source_scan_covers_every_submodule() {
let mod_rs = include_str!("mod.rs");
let declared: Vec<&str> = mod_rs
.lines()
.map(str::trim)
.filter_map(|l| {
l.strip_prefix("mod ")
.or_else(|| l.strip_prefix("pub mod "))
})
.filter_map(|rest| rest.strip_suffix(';'))
.collect();
assert!(
!declared.is_empty(),
"found no `mod` declarations in control_bus/mod.rs — the parser above is \
out of date, which would make the transport scan vacuous"
);
for name in &declared {
let expected = format!("control_bus/{name}.rs");
assert!(
MODULE_SOURCES.iter().any(|(n, _)| *n == expected),
"control_bus/mod.rs declares `mod {name};` but MODULE_SOURCES has no \
row for {expected} — add one so the transport scan covers it"
);
}
for expected in ["control_bus/mod.rs", "control_bus/tests.rs"] {
assert!(
MODULE_SOURCES.iter().any(|(n, _)| *n == expected),
"MODULE_SOURCES is missing {expected}"
);
}
}