use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use tower::dispatch::{
DispatchAuditEvent, DispatchContext, DispatchEngine, DispatchError, DispatchOutcome,
DispatchRecorder, NoopTriggerHandler, TriggerHandler, TriggerKind,
};
use tower::event::{Event, EventScope, Level};
use tower::supervisor::FiredEvent;
use tower_rules::*;
#[derive(Clone)]
struct VecRecorder {
events: Arc<Mutex<Vec<DispatchAuditEvent>>>,
}
impl VecRecorder {
fn new() -> Self {
Self { events: Arc::new(Mutex::new(Vec::new())) }
}
fn events(&self) -> Vec<DispatchAuditEvent> {
self.events.lock().unwrap().clone()
}
}
impl DispatchRecorder for VecRecorder {
fn record(&self, event: DispatchAuditEvent) {
self.events.lock().unwrap().push(event);
}
}
struct FailingHandler;
impl TriggerHandler for FailingHandler {
fn handle(&self, _ctx: &DispatchContext) -> Result<(), DispatchError> {
Err(DispatchError::Failed("simulated failure".into()))
}
}
#[derive(Clone)]
struct RecordingHandler {
calls: Arc<Mutex<Vec<DispatchContext>>>,
}
impl RecordingHandler {
fn new() -> Self {
Self { calls: Arc::new(Mutex::new(Vec::new())) }
}
fn calls(&self) -> Vec<DispatchContext> {
self.calls.lock().unwrap().clone()
}
}
impl TriggerHandler for RecordingHandler {
fn handle(&self, ctx: &DispatchContext) -> Result<(), DispatchError> {
self.calls.lock().unwrap().push(ctx.clone());
Ok(())
}
}
fn notification_rule(id: &str) -> TowerRule {
TowerRule {
schema_version: SchemaVersion::V1,
id: RuleId(id.into()),
name: id.into(),
predicate: Predicate::EventMatch {
scope: ScopeFilter::Any,
level: None,
target: None,
fields: vec![],
rate: None,
},
trigger: Trigger::Notification {
channels: vec![NotificationChannel::DesktopBadge],
severity: Severity::Warning,
},
debounce_ms: None,
federation: FederationPolicy::LocalOnly,
enabled: true,
}
}
fn agent_dispatch_rule(id: &str) -> TowerRule {
TowerRule {
schema_version: SchemaVersion::V1,
id: RuleId(id.into()),
name: id.into(),
predicate: Predicate::ScryerHealth { signal: HealthSignal::YubabaRaftQuorumLost },
trigger: Trigger::AgentDispatch {
agent_class: AgentClassRef("gnome/fixer".into()),
prompt_template: PromptTemplate("Fix the thing: {{event}}".into()),
placement: TaskPlacement::new(TaskLocation::Local, TaskRuntime::Native),
},
debounce_ms: None,
federation: FederationPolicy::LocalOnly,
enabled: true,
}
}
fn yubaba_action_rule(id: &str) -> TowerRule {
TowerRule {
schema_version: SchemaVersion::V1,
id: RuleId(id.into()),
name: id.into(),
predicate: Predicate::EventMatch {
scope: ScopeFilter::Any,
level: None,
target: None,
fields: vec![],
rate: None,
},
trigger: Trigger::YubabaAction {
kind: YubabaActionKind::RestartWorkload,
target: MeshIdent("api.pdx".into()),
},
debounce_ms: None,
federation: FederationPolicy::LocalOnly,
enabled: true,
}
}
fn service_event(ident: &str, seq: u64) -> Event {
Event {
scope: EventScope::Service(MeshIdent(ident.into())),
level: Level::Error,
target: "test.event".into(),
msg: String::new(),
fields: HashMap::new(),
seq,
}
}
fn fired(rule: TowerRule, event: Event) -> FiredEvent {
let rule_id = rule.id.clone();
FiredEvent { rule_id, rule, event, peer: None }
}
fn fired_from_peer(rule: TowerRule, event: Event, peer: &str) -> FiredEvent {
let rule_id = rule.id.clone();
FiredEvent { rule_id, rule, event, peer: Some(peer.into()) }
}
fn noop_engine(recorder: VecRecorder) -> DispatchEngine {
DispatchEngine::new(5, Box::new(NoopTriggerHandler), Box::new(recorder))
}
fn failing_engine(recorder: VecRecorder) -> DispatchEngine {
DispatchEngine::new(5, Box::new(FailingHandler), Box::new(recorder))
}
#[test]
fn process_empty_is_noop() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![]);
assert!(recorder.events().is_empty());
}
#[test]
fn process_calls_trigger_handler() {
let handler = RecordingHandler::new();
let recorder = VecRecorder::new();
let mut engine = DispatchEngine::new(5, Box::new(handler.clone()), Box::new(recorder));
engine.process(vec![fired(notification_rule("rule-1"), service_event("api.pdx", 1))]);
let calls = handler.calls();
assert_eq!(calls.len(), 1);
assert_eq!(calls[0].matched_event.seq, 1);
assert_eq!(calls[0].rule_id, RuleId("rule-1".into()));
}
#[test]
fn process_multiple_events_calls_handler_per_event() {
let handler = RecordingHandler::new();
let mut engine = DispatchEngine::new(5, Box::new(handler.clone()), Box::new(VecRecorder::new()));
let rule = notification_rule("r");
engine.process(vec![
fired(rule.clone(), service_event("x", 1)),
fired(rule.clone(), service_event("x", 2)),
fired(rule, service_event("x", 3)),
]);
assert_eq!(handler.calls().len(), 3);
}
mod audit_trail {
use super::*;
#[test]
fn dispatched_record_written_on_success() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(notification_rule("rule-1"), service_event("api.pdx", 42))]);
let events = recorder.events();
assert_eq!(events.len(), 1);
assert_eq!(events[0].rule_id, RuleId("rule-1".into()));
assert_eq!(events[0].outcome, DispatchOutcome::Dispatched);
assert_eq!(events[0].matched_seq, 42);
assert_eq!(events[0].matched_scope, "service:api.pdx");
assert!(events[0].ts_ms > 0, "ts_ms should be set");
}
#[test]
fn failure_writes_dispatched_then_failed() {
let recorder = VecRecorder::new();
let mut engine = failing_engine(recorder.clone());
engine.process(vec![fired(notification_rule("rule-1"), service_event("api.pdx", 1))]);
let events = recorder.events();
assert_eq!(events.len(), 2, "intent record + failure record");
assert_eq!(events[0].outcome, DispatchOutcome::Dispatched,
"intent record must come first");
assert!(
matches!(&events[1].outcome, DispatchOutcome::Failed { .. }),
"failure record must follow"
);
}
#[test]
fn failure_record_includes_cause() {
let recorder = VecRecorder::new();
let mut engine = failing_engine(recorder.clone());
engine.process(vec![fired(notification_rule("r"), service_event("x", 1))]);
let failed = recorder.events()
.into_iter()
.find(|e| matches!(e.outcome, DispatchOutcome::Failed { .. }))
.expect("should have a Failed event");
match &failed.outcome {
DispatchOutcome::Failed { cause } => {
assert!(cause.contains("simulated failure"), "cause = {cause}");
}
_ => panic!("expected Failed"),
}
}
#[test]
fn dispatched_record_precedes_trigger_invocation() {
let recorder = VecRecorder::new();
let mut engine = failing_engine(recorder.clone());
engine.process(vec![fired(notification_rule("r"), service_event("x", 1))]);
let events = recorder.events();
assert_eq!(events.len(), 2);
assert_eq!(events[0].outcome, DispatchOutcome::Dispatched,
"Dispatched must be recorded before trigger fires");
assert!(matches!(events[1].outcome, DispatchOutcome::Failed { .. }),
"Failed follows Dispatched");
}
#[test]
fn audit_event_scope_format() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(notification_rule("r"), service_event("yubaba.local", 7))]);
assert_eq!(recorder.events()[0].matched_scope, "service:yubaba.local");
}
#[test]
fn to_scryer_event_shape_on_dispatch() {
let audit = DispatchAuditEvent {
rule_id: RuleId("test-rule".into()),
trigger_kind: TriggerKind::Notification,
outcome: DispatchOutcome::Dispatched,
matched_seq: 42,
matched_scope: "service:yubaba.local".into(),
context_event_count: 3,
peer: None,
ts_ms: 1_000,
};
let event = audit.to_scryer_event(99);
assert_eq!(event.target, "tower.dispatch");
assert_eq!(event.scope, EventScope::Service(MeshIdent("tower.local".into())));
assert_eq!(event.seq, 99);
assert_eq!(event.fields["rule_id"], serde_json::json!("test-rule"));
assert_eq!(event.fields["matched_scope"], serde_json::json!("service:yubaba.local"));
assert_eq!(event.fields["trigger_kind"], serde_json::json!("notification"));
assert_eq!(event.fields["context_event_count"], serde_json::json!(3));
assert!(!event.fields.contains_key("peer"), "no peer field for local events");
}
#[test]
fn to_scryer_event_shape_on_failure() {
let audit = DispatchAuditEvent {
rule_id: RuleId("test-rule".into()),
trigger_kind: TriggerKind::AgentDispatch,
outcome: DispatchOutcome::Failed { cause: "bad template".into() },
matched_seq: 1,
matched_scope: "service:api".into(),
context_event_count: 0,
peer: None,
ts_ms: 1_000,
};
let event = audit.to_scryer_event(1);
assert_eq!(event.target, "tower.dispatch.failed");
assert_eq!(event.fields["cause"], serde_json::json!("bad template"));
}
#[test]
fn to_scryer_event_peer_field_set_when_federated() {
let audit = DispatchAuditEvent {
rule_id: RuleId("r".into()),
trigger_kind: TriggerKind::Notification,
outcome: DispatchOutcome::Dispatched,
matched_seq: 1,
matched_scope: "service:db".into(),
context_event_count: 0,
peer: Some("peer-1".into()),
ts_ms: 1_000,
};
let event = audit.to_scryer_event(1);
assert_eq!(event.fields["peer"], serde_json::json!("peer-1"));
}
}
mod trigger_kind {
use super::*;
#[test]
fn notification_rule_emits_notification_kind() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(notification_rule("r"), service_event("x", 1))]);
assert_eq!(recorder.events()[0].trigger_kind, TriggerKind::Notification);
}
#[test]
fn agent_dispatch_rule_emits_agent_dispatch_kind() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(agent_dispatch_rule("r"), service_event("x", 1))]);
assert_eq!(recorder.events()[0].trigger_kind, TriggerKind::AgentDispatch);
}
#[test]
fn yubaba_action_rule_emits_yubaba_action_kind() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(yubaba_action_rule("r"), service_event("x", 1))]);
assert_eq!(recorder.events()[0].trigger_kind, TriggerKind::YubabaAction);
}
}
mod context_ring {
use super::*;
#[test]
fn first_event_has_empty_context() {
let handler = RecordingHandler::new();
let mut engine =
DispatchEngine::new(5, Box::new(handler.clone()), Box::new(VecRecorder::new()));
engine.process(vec![fired(notification_rule("r"), service_event("x", 1))]);
assert_eq!(handler.calls()[0].context_events.len(), 0);
}
#[test]
fn second_event_sees_first_as_context() {
let handler = RecordingHandler::new();
let mut engine =
DispatchEngine::new(5, Box::new(handler.clone()), Box::new(VecRecorder::new()));
let rule = notification_rule("r");
engine.process(vec![
fired(rule.clone(), service_event("x", 1)),
fired(rule, service_event("x", 2)),
]);
let calls = handler.calls();
assert_eq!(calls[1].context_events.len(), 1);
assert_eq!(calls[1].context_events[0].seq, 1, "context should be the first event");
}
#[test]
fn context_window_limits_ring_size() {
let handler = RecordingHandler::new();
let window = 3usize;
let mut engine =
DispatchEngine::new(window, Box::new(handler.clone()), Box::new(VecRecorder::new()));
let rule = notification_rule("r");
for seq in 1..=6u64 {
engine.process(vec![fired(rule.clone(), service_event("x", seq))]);
}
engine.process(vec![fired(rule, service_event("x", 7))]);
let calls = handler.calls();
let last = calls.last().unwrap();
assert!(
last.context_events.len() <= window,
"context ({}) must not exceed window ({})",
last.context_events.len(),
window
);
}
#[test]
fn context_window_preserves_most_recent_events() {
let handler = RecordingHandler::new();
let window = 2usize;
let mut engine =
DispatchEngine::new(window, Box::new(handler.clone()), Box::new(VecRecorder::new()));
let rule = notification_rule("r");
for seq in 1..=4u64 {
engine.process(vec![fired(rule.clone(), service_event("x", seq))]);
}
engine.process(vec![fired(rule, service_event("x", 5))]);
let calls = handler.calls();
let last = calls.last().unwrap();
let seqs: Vec<u64> = last.context_events.iter().map(|e| e.seq).collect();
assert_eq!(seqs, vec![3, 4], "ring should hold the two most recent events");
}
#[test]
fn context_rings_are_independent_per_rule() {
let handler = RecordingHandler::new();
let mut engine =
DispatchEngine::new(5, Box::new(handler.clone()), Box::new(VecRecorder::new()));
engine.process(vec![
fired(notification_rule("rule-a"), service_event("x", 1)),
fired(notification_rule("rule-b"), service_event("y", 2)),
fired(notification_rule("rule-a"), service_event("x", 3)),
]);
let calls = handler.calls();
let rule_a_second = calls
.iter()
.find(|c| c.rule_id == RuleId("rule-a".into()) && c.matched_event.seq == 3)
.expect("rule-a seq=3 call not found");
assert_eq!(rule_a_second.context_events.len(), 1);
assert_eq!(rule_a_second.context_events[0].seq, 1);
}
#[test]
fn context_event_count_in_audit_event_matches_context() {
let recorder = VecRecorder::new();
let handler = RecordingHandler::new();
let mut engine = DispatchEngine::new(5, Box::new(handler.clone()), Box::new(recorder.clone()));
let rule = notification_rule("r");
engine.process(vec![
fired(rule.clone(), service_event("x", 1)),
fired(rule.clone(), service_event("x", 2)),
fired(rule, service_event("x", 3)),
]);
let audit = recorder.events();
assert_eq!(audit[0].context_event_count, 0, "first event has no context");
assert_eq!(audit[1].context_event_count, 1);
assert_eq!(audit[2].context_event_count, 2);
}
}
#[test]
fn peer_propagated_to_audit_event() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired_from_peer(
notification_rule("r"),
service_event("db.remote", 1),
"peer-1",
)]);
assert_eq!(recorder.events()[0].peer, Some("peer-1".into()));
}
#[test]
fn peer_propagated_to_dispatch_context() {
let handler = RecordingHandler::new();
let mut engine = DispatchEngine::new(5, Box::new(handler.clone()), Box::new(VecRecorder::new()));
engine.process(vec![fired_from_peer(
notification_rule("r"),
service_event("db.remote", 1),
"peer-2",
)]);
assert_eq!(handler.calls()[0].peer, Some("peer-2".into()));
}
#[test]
fn local_event_has_no_peer() {
let recorder = VecRecorder::new();
let mut engine = noop_engine(recorder.clone());
engine.process(vec![fired(notification_rule("r"), service_event("x", 1))]);
assert!(recorder.events()[0].peer.is_none());
}