use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tower_rules::{FederationPolicy, Predicate, RatePredicate, RuleId, TowerRule};
use crate::event::{Event, EventScope};
use crate::rules::compiler::{self, CompileError};
use crate::scryer_filter::ScryerFilter;
pub trait EventStream: Send {
fn try_next(&mut self) -> Option<Event>;
}
pub trait SubscriptionSource: Send + Sync {
fn subscribe(&self, filter: ScryerFilter, rule_id: &RuleId) -> Box<dyn EventStream>;
}
pub trait PeerRegistry: Send + Sync {
fn peers(&self) -> Vec<String>;
fn source_for(&self, peer: &str) -> Option<Arc<dyn SubscriptionSource>>;
}
pub struct NoPeers;
impl PeerRegistry for NoPeers {
fn peers(&self) -> Vec<String> {
vec![]
}
fn source_for(&self, _peer: &str) -> Option<Arc<dyn SubscriptionSource>> {
None
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct DedupKey {
pub rule_id: RuleId,
pub scope_id: String,
pub seq: u64,
}
pub(crate) fn scope_id(scope: &EventScope) -> String {
match scope {
EventScope::TaskRun(id) => format!("task:{id}"),
EventScope::Service(ident) => format!("service:{}", ident.0),
EventScope::Forge(id) => format!("forge:{}", id.0),
EventScope::Health => "health".to_string(),
}
}
struct DedupRing {
capacity: usize,
entries: VecDeque<(DedupKey, Instant)>,
}
impl DedupRing {
fn new(capacity: usize) -> Self {
Self { capacity, entries: VecDeque::with_capacity(capacity) }
}
fn contains(&self, key: &DedupKey, debounce: Option<Duration>) -> bool {
self.entries.iter().any(|(k, fired_at)| {
k == key && debounce.map_or(true, |d| fired_at.elapsed() < d)
})
}
fn insert(&mut self, key: DedupKey) {
if self.entries.len() >= self.capacity {
self.entries.pop_front();
}
self.entries.push_back((key, Instant::now()));
}
}
struct RateWindow {
window_ms: u64,
min_count: u32,
timestamps: VecDeque<Instant>,
}
impl RateWindow {
fn new(min_count: u32, window_ms: u64) -> Self {
Self { window_ms, min_count, timestamps: VecDeque::new() }
}
fn record(&mut self) -> bool {
let now = Instant::now();
let window = Duration::from_millis(self.window_ms);
self.timestamps.retain(|t| now.duration_since(*t) < window);
self.timestamps.push_back(now);
self.timestamps.len() >= self.min_count as usize
}
}
struct RuleEntry {
rule: TowerRule,
filter: ScryerFilter,
local_stream: Box<dyn EventStream>,
peer_streams: HashMap<String, Box<dyn EventStream>>,
rate_window: Option<RateWindow>,
}
#[derive(Debug)]
pub struct FiredEvent {
pub rule_id: RuleId,
pub rule: TowerRule,
pub event: Event,
pub peer: Option<String>,
}
#[derive(Debug)]
pub struct SubscriptionHealth {
pub rule_id: RuleId,
pub kind: SubscriptionHealthKind,
}
#[derive(Debug)]
pub enum SubscriptionHealthKind {
Opened,
Closed,
Lag { dropped: u64 },
}
pub struct Supervisor {
local_source: Arc<dyn SubscriptionSource>,
peer_registry: Arc<dyn PeerRegistry>,
rules: HashMap<RuleId, RuleEntry>,
dedup: DedupRing,
backpressure_limit: usize,
}
impl Supervisor {
pub fn new(
local_source: Arc<dyn SubscriptionSource>,
peer_registry: Arc<dyn PeerRegistry>,
dedup_capacity: usize,
backpressure_limit: usize,
) -> Self {
Self {
local_source,
peer_registry,
rules: HashMap::new(),
dedup: DedupRing::new(dedup_capacity),
backpressure_limit,
}
}
pub fn load_rule(&mut self, rule: TowerRule) -> Result<bool, CompileError> {
if !rule.enabled {
return Ok(false);
}
let filter = compiler::compile(&rule)?;
let local_stream = self.local_source.subscribe(filter.clone(), &rule.id);
let peer_streams = open_peer_streams(&rule, &filter, &*self.peer_registry);
let rate_window = make_rate_window(&rule.predicate);
self.rules.insert(
rule.id.clone(),
RuleEntry { rule, filter, local_stream, peer_streams, rate_window },
);
Ok(true)
}
pub fn disable_rule(&mut self, id: &RuleId) -> bool {
self.rules.remove(id).is_some()
}
pub fn update_rule(&mut self, rule: TowerRule) -> Result<(), CompileError> {
let id = rule.id.clone();
self.disable_rule(&id);
self.load_rule(rule)?;
Ok(())
}
pub fn poll(&mut self) -> (Vec<FiredEvent>, Vec<SubscriptionHealth>) {
let mut candidates: Vec<(RuleId, Event, Option<String>, Option<Duration>)> = Vec::new();
let mut health: Vec<SubscriptionHealth> = Vec::new();
for (rule_id, entry) in self.rules.iter_mut() {
let debounce = entry.rule.debounce_ms.map(Duration::from_millis);
let mut read_count = 0usize;
let mut dropped = 0u64;
loop {
if read_count >= self.backpressure_limit {
while entry.local_stream.try_next().is_some() {
dropped += 1;
}
break;
}
match entry.local_stream.try_next() {
None => break,
Some(event) => {
read_count += 1;
if entry.filter.matches(&event) {
candidates.push((rule_id.clone(), event, None, debounce));
}
}
}
}
if dropped > 0 {
health.push(SubscriptionHealth {
rule_id: rule_id.clone(),
kind: SubscriptionHealthKind::Lag { dropped },
});
}
for (peer, stream) in entry.peer_streams.iter_mut() {
while let Some(event) = stream.try_next() {
if entry.filter.matches(&event) {
candidates.push((
rule_id.clone(),
event,
Some(peer.clone()),
debounce,
));
}
}
}
}
let mut fired = Vec::new();
for (rule_id, event, peer, debounce) in candidates {
let key = DedupKey {
rule_id: rule_id.clone(),
scope_id: scope_id(&event.scope),
seq: event.seq,
};
if self.dedup.contains(&key, debounce) {
continue;
}
let rate_ok = self
.rules
.get_mut(&rule_id)
.and_then(|e| e.rate_window.as_mut())
.map_or(true, |rw| rw.record());
if rate_ok {
self.dedup.insert(key);
let rule = self.rules[&rule_id].rule.clone();
fired.push(FiredEvent { rule_id, rule, event, peer });
}
}
(fired, health)
}
pub fn active_count(&self) -> usize {
self.rules.len()
}
pub fn is_active(&self, id: &RuleId) -> bool {
self.rules.contains_key(id)
}
}
fn open_peer_streams(
rule: &TowerRule,
filter: &ScryerFilter,
peers: &dyn PeerRegistry,
) -> HashMap<String, Box<dyn EventStream>> {
let peer_names: Vec<String> = match &rule.federation {
FederationPolicy::LocalOnly => return HashMap::new(),
FederationPolicy::MeshWide => peers.peers(),
FederationPolicy::MirrorSet { mirrors } => mirrors.clone(),
};
peer_names
.into_iter()
.filter_map(|peer| {
peers
.source_for(&peer)
.map(|src| (peer.clone(), src.subscribe(filter.clone(), &rule.id)))
})
.collect()
}
fn make_rate_window(predicate: &Predicate) -> Option<RateWindow> {
match predicate {
Predicate::EventMatch { rate: Some(RatePredicate { min_count, window_ms }), .. } => {
Some(RateWindow::new(*min_count, *window_ms))
}
_ => None,
}
}