use std::collections::{HashSet, VecDeque};
use tower_rules::{Predicate, RatePredicate, RuleId, TowerRule, Trigger, YubabaActionKind};
use crate::dispatch::DispatchContext;
use crate::event::Event;
use crate::rules::compiler::{self, CompileError};
use crate::supervisor::scope_id;
use crate::triggers::{
event_to_json, render_prompt, AgentDispatchSpec, NotificationPayload, YubabaCommand,
};
#[derive(Debug, Clone)]
pub struct SimulateOptions {
pub event_spacing_ms: u64,
pub context_window: usize,
}
impl Default for SimulateOptions {
fn default() -> Self {
Self { event_spacing_ms: 1_000, context_window: 5 }
}
}
#[derive(Debug, Clone)]
pub enum RenderedTrigger {
AgentDispatch(AgentDispatchSpec),
Notification(NotificationPayload),
YubabaAction(YubabaCommand),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DebounceDecision {
WouldFire,
WouldSuppress {
debounce_remaining_ms: u64,
},
BelowRateThreshold {
seen: usize,
required: u32,
},
}
#[derive(Debug)]
pub struct SimulateMatch {
pub event: Event,
pub rendered_trigger: RenderedTrigger,
pub decision: DebounceDecision,
}
#[derive(Debug)]
pub struct SimulateReport {
pub rule_id: RuleId,
pub rule_name: String,
pub total_events: usize,
pub filter_matches: usize,
pub would_fire: usize,
pub would_suppress: usize,
pub matches: Vec<SimulateMatch>,
}
#[derive(Debug, thiserror::Error)]
pub enum SimulateError {
#[error("compile error: {0}")]
Compile(#[from] CompileError),
}
struct SimRateWindow {
window_ms: u64,
min_count: u32,
timestamps: VecDeque<u64>,
}
impl SimRateWindow {
fn new(rate: &RatePredicate) -> Self {
Self { window_ms: rate.window_ms, min_count: rate.min_count, timestamps: VecDeque::new() }
}
fn record(&mut self, sim_ms: u64) -> (usize, bool) {
self.timestamps.retain(|&t| sim_ms.saturating_sub(t) < self.window_ms);
self.timestamps.push_back(sim_ms);
let count = self.timestamps.len();
(count, count >= self.min_count as usize)
}
}
pub fn simulate(
rule: &TowerRule,
events: &[Event],
options: &SimulateOptions,
) -> Result<SimulateReport, SimulateError> {
let filter = compiler::compile(rule)?;
let debounce_ms = rule.debounce_ms;
let mut last_fire_sim_ms: Option<u64> = None;
let mut rate_window = extract_rate_window(&rule.predicate);
let mut context_ring: VecDeque<Event> = VecDeque::new();
let mut seen: HashSet<(String, u64)> = HashSet::new();
let mut matches = Vec::new();
let mut filter_matches = 0usize;
let mut sim_ms: u64 = 0;
for event in events {
sim_ms = sim_ms.saturating_add(options.event_spacing_ms);
if !filter.matches(event) {
continue;
}
filter_matches += 1;
let eid = (scope_id(&event.scope), event.seq);
if seen.contains(&eid) {
continue;
}
let rate_decision = if let Some(rw) = rate_window.as_mut() {
let (seen_count, threshold_met) = rw.record(sim_ms);
if !threshold_met {
Some(DebounceDecision::BelowRateThreshold {
seen: seen_count,
required: rw.min_count,
})
} else {
None
}
} else {
None
};
let debounce_decision = if rate_decision.is_none() {
if let Some(d_ms) = debounce_ms {
if let Some(last_ms) = last_fire_sim_ms {
let elapsed = sim_ms.saturating_sub(last_ms);
if elapsed < d_ms {
Some(DebounceDecision::WouldSuppress {
debounce_remaining_ms: d_ms - elapsed,
})
} else {
None
}
} else {
None
}
} else {
None
}
} else {
None
};
let decision =
rate_decision.or(debounce_decision).unwrap_or(DebounceDecision::WouldFire);
if decision == DebounceDecision::WouldFire {
last_fire_sim_ms = Some(sim_ms);
seen.insert(eid);
}
let context_events: Vec<Event> = context_ring.iter().cloned().collect();
if context_ring.len() >= options.context_window {
context_ring.pop_front();
}
context_ring.push_back(event.clone());
let ctx = DispatchContext {
rule_id: rule.id.clone(),
rule: rule.clone(),
matched_event: event.clone(),
context_events,
peer: None,
};
let rendered_trigger = render_trigger_dry(&ctx);
matches.push(SimulateMatch { event: event.clone(), rendered_trigger, decision });
}
let would_fire =
matches.iter().filter(|m| m.decision == DebounceDecision::WouldFire).count();
let would_suppress = matches.len() - would_fire;
Ok(SimulateReport {
rule_id: rule.id.clone(),
rule_name: rule.name.clone(),
total_events: events.len(),
filter_matches,
would_fire,
would_suppress,
matches,
})
}
fn extract_rate_window(predicate: &Predicate) -> Option<SimRateWindow> {
match predicate {
Predicate::EventMatch { rate: Some(rate), .. } => Some(SimRateWindow::new(rate)),
_ => None,
}
}
fn render_trigger_dry(ctx: &DispatchContext) -> RenderedTrigger {
match &ctx.rule.trigger {
Trigger::AgentDispatch { agent_class, prompt_template, placement } => {
let prompt = render_prompt(prompt_template, ctx);
RenderedTrigger::AgentDispatch(AgentDispatchSpec {
agent_class: agent_class.clone(),
prompt,
placement: placement.clone(),
})
}
Trigger::Notification { channels: _, severity } => {
let event_json = event_to_json(&ctx.matched_event)
.map(|v| v.to_string())
.unwrap_or_else(|_| "{}".into());
let matched_scope = scope_id(&ctx.matched_event.scope);
RenderedTrigger::Notification(NotificationPayload {
rule_name: ctx.rule.name.clone(),
severity: *severity,
event_json,
matched_scope,
})
}
Trigger::YubabaAction { kind, target } => {
let cmd = match kind {
YubabaActionKind::RestartWorkload => {
YubabaCommand::RestartWorkload { target: target.clone() }
}
YubabaActionKind::Drain => YubabaCommand::Drain { target: target.clone() },
YubabaActionKind::Scale { replicas } => {
YubabaCommand::Scale { target: target.clone(), replicas: *replicas }
}
};
RenderedTrigger::YubabaAction(cmd)
}
}
}