use std::collections::{BTreeSet, HashSet};
use std::io::Write;
use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::bus::monitor::{FleetEvent, SampleView, StreamItem};
use crate::judge::condition::{Condition, RuleSet, SweepOutcome};
use crate::model::decode::SchemaStore;
use crate::model::registry::SliceSet;
use crate::model::retain::RetentionBudget;
use crate::report::{
CondState, PreRollInfo, PreambleInfo, PreambleSemantics, RecordReport, Transition, ZrecHeader,
};
use crate::tape::record::{RecordBounds, ZREC_VERSION, ZrecSink, record};
use crate::{Error, Result};
#[derive(Debug, Clone)]
pub struct TriggerSpec {
pub selectors: Vec<String>,
pub pre: Duration,
pub post: Duration,
pub rules: Vec<Condition>,
pub tick: Duration,
pub timeout: Duration,
pub give_up: Option<Duration>,
pub preamble: Option<PreambleSemantics>,
pub max_samples: Option<u64>,
pub max_replies: usize,
}
#[derive(Debug, Clone)]
pub enum TriggerEvent<'a> {
Armed {
watched: &'a [String],
pre: Duration,
},
Transition(&'a Transition),
Fired(&'a Transition),
Preamble(&'a PreambleInfo),
Progress { samples: u64, dropped: u64 },
GaveUp { after: Duration },
}
const FIRE_BUFFER: usize = 4096;
pub fn state_projection(base: &str, selector: &str) -> Option<String> {
use zenkey::grammar::{CLASS_STATE, VERSION_CHUNK};
let rel = zenkey::grammar::strip_base(base, selector)?;
let chunks: Vec<&str> = rel.split('/').collect();
let projected: Vec<String> = match chunks.as_slice() {
["**"] | [VERSION_CHUNK, "**"] => {
vec![
VERSION_CHUNK.into(),
"*".into(),
CLASS_STATE.into(),
"**".into(),
]
}
[VERSION_CHUNK, origin, "**"] => {
vec![
VERSION_CHUNK.into(),
(*origin).into(),
CLASS_STATE.into(),
"**".into(),
]
}
[VERSION_CHUNK, origin, class, rest @ ..] if *class == CLASS_STATE || *class == "*" => {
let mut v = vec![
VERSION_CHUNK.to_string(),
(*origin).into(),
CLASS_STATE.into(),
];
v.extend(rest.iter().map(|c| (*c).to_string()));
if rest.is_empty() {
v.push("**".into());
}
v
}
_ => return None,
};
Some(zenkey::grammar::with_base(base, projected.join("/")))
}
pub fn watch_cover<'s>(selectors: impl IntoIterator<Item = &'s String>) -> Vec<String> {
let mut distinct: Vec<String> = Vec::new();
for sel in selectors {
if !distinct.contains(sel) {
distinct.push(sel.clone());
}
}
let parsed: Vec<Option<zenoh::key_expr::KeyExpr<'static>>> = distinct
.iter()
.map(|s| zenoh::key_expr::KeyExpr::try_from(s.clone()).ok())
.collect();
distinct
.iter()
.enumerate()
.filter(|(i, _)| {
let Some(mine) = &parsed[*i] else { return true };
!parsed
.iter()
.enumerate()
.any(|(j, other)| j != *i && other.as_ref().is_some_and(|o| o.includes(mine)))
})
.map(|(_, s)| s.clone())
.collect()
}
pub async fn record_on<W, F>(
fleet: &crate::Fleet<'_>,
slices: Option<&SliceSet>,
store: &SchemaStore,
spec: &TriggerSpec,
open: impl FnOnce() -> F,
mut on_event: impl FnMut(TriggerEvent<'_>),
) -> Result<RecordReport>
where
W: Write + Send + 'static,
F: std::future::Future<Output = Result<W>>,
{
let (session, base) = (fleet.session(), fleet.base());
if spec.rules.is_empty() {
return Err(Error::unaskable(
"--on",
"a trigger capture needs at least one rule to fire on",
));
}
let mut rules = RuleSet::new(&spec.rules, base, slices)?;
let watched = watch_cover(spec.selectors.iter().chain(rules.watched()));
let (wants_doctor, wants_roster, wants_decode) = (
rules.wants_doctor(),
rules.wants_roster(),
rules.wants_decode(),
);
if wants_decode {
crate::model::decode::prewarm(fleet, store, slices).await;
}
let _sealed = store.seal();
let monitor = crate::Monitor::start(session, crate::MonitorSpec::default()).await?;
monitor.core().set_retention_budget(RetentionBudget {
max_age: spec.pre,
..RetentionBudget::default()
});
let mut events = monitor.events();
let monitor = monitor.watching(&watched).await?;
let core = Arc::clone(monitor.core());
on_event(TriggerEvent::Armed {
watched: &watched,
pre: spec.pre,
});
let armed = tokio::time::Instant::now();
let give_up_at = spec.give_up.map(|d| armed + d);
let mut facts_cache = crate::model::facts::FactsCache::default();
let fired: Option<Transition> = 'armed: loop {
let deadline = rules.last_eval() + spec.tick;
let sweep = async {
let doctor = if wants_doctor {
Some(
crate::judge::doctor::run_doctor(
fleet,
slices,
&crate::judge::doctor::DoctorSpec {
deep: false,
sample: None,
timeout: spec.timeout,
listen: None,
},
)
.await
.map_err(|e| e.to_string()),
)
} else {
None
};
let roster = if wants_roster {
Some(
crate::bus::roster::roster(fleet, spec.timeout)
.await
.map_err(|e| e.to_string()),
)
} else {
None
};
if wants_decode {
crate::model::decode::prewarm(fleet, store, slices).await;
}
(doctor, roster)
};
let mut sweep = std::pin::pin!(sweep);
let mut swept = None;
let tick_over = tokio::time::sleep_until(deadline);
tokio::pin!(tick_over);
let give_up = async {
match give_up_at {
Some(at) => tokio::time::sleep_until(at).await,
None => std::future::pending().await,
}
};
tokio::pin!(give_up);
let mut closed = false;
let mut gave_up = false;
while !closed {
let item = tokio::select! {
item = events.recv() => item,
outcome = &mut sweep, if swept.is_none() => {
swept = Some(outcome);
continue;
}
() = &mut tick_over, if swept.is_some() => break,
() = &mut give_up => {
gave_up = true;
break;
}
};
match item {
Some(StreamItem::Event(FleetEvent::Sample(s))) => {
let verdict = if rules.wants_verdict(&s) {
Some(
crate::model::decode::decode_sample(
fleet,
store,
slices,
&s.key,
Some(&s.encoding),
&s.payload.to_bytes(),
)
.await
.verdict,
)
} else {
None
};
rules.observe_sample(&s, &mut facts_cache, verdict.as_ref());
}
Some(StreamItem::Dropped(n)) => rules.observe_drop(n),
Some(_) => {}
None => closed = true,
}
}
if gave_up {
break 'armed None;
}
if closed {
break 'armed None;
}
let (doctor_outcome, roster_outcome) = match swept {
Some(outcome) => outcome,
None => sweep.await,
};
let now = tokio::time::Instant::now();
let at = crate::tape::record::rfc3339_now();
let transitions = rules.evaluate(
now,
&at,
SweepOutcome {
doctor: doctor_outcome
.as_ref()
.map(|o| o.as_ref().map_err(String::as_str)),
roster: roster_outcome
.as_ref()
.map(|o| o.as_ref().map_err(String::as_str)),
},
);
for t in transitions {
on_event(TriggerEvent::Transition(&t));
if t.to == CondState::Firing {
break 'armed Some(t);
}
}
};
let Some(trigger) = fired else {
monitor.shutdown().await?;
let after = armed.elapsed();
on_event(TriggerEvent::GaveUp { after });
return Ok(RecordReport {
header: ZrecHeader {
zrec: ZREC_VERSION,
selectors: watched,
base: base.to_string(),
captured_at: crate::tape::record::rfc3339_now(),
preamble: None,
pre_roll: None,
},
out: None,
samples: 0,
dropped: 0,
duration_ms: u64::try_from(after.as_millis()).unwrap_or(u64::MAX),
trigger: None,
preamble: None,
pre_roll: None,
preamble_rows: 0,
});
};
on_event(TriggerEvent::Fired(&trigger));
let fired_at = Instant::now();
let captured_at = crate::tape::record::rfc3339_now();
let ring = core.retained();
let stats = core.retention();
let epoch = ring.first().map_or(fired_at, |v| v.received);
let pre_roll = PreRollInfo {
asked_s: spec.pre.as_secs_f64(),
covered_s: stats.span.as_secs_f64(),
watched: watched.clone(),
evicted: stats.evicted,
expired: stats.expired,
};
let already: HashSet<usize> = ring
.iter()
.rev()
.take(crate::MonitorSpec::default().capacity)
.map(|v| Arc::as_ptr(v) as usize)
.collect();
let mut buffered: Vec<StreamItem> = Vec::new();
let mut buffer_dropped = 0u64;
let drain_into =
|item: Option<StreamItem>, buffered: &mut Vec<StreamItem>, dropped: &mut u64| match item {
Some(StreamItem::Event(FleetEvent::Sample(s))) => {
if already.contains(&(Arc::as_ptr(&s) as usize)) {
return;
}
if buffered.len() < FIRE_BUFFER {
buffered.push(StreamItem::Event(FleetEvent::Sample(s)));
} else {
*dropped += 1;
}
}
Some(StreamItem::Dropped(n)) => buffered.push(StreamItem::Dropped(n)),
_ => {}
};
let preamble = async {
let semantics = spec.preamble?;
let started = Instant::now();
let mut selectors = Vec::new();
let mut failed = Vec::new();
for sel in &watched {
match state_projection(base, sel) {
Some(p) if !selectors.contains(&p) => selectors.push(p),
Some(_) => {}
None => failed.push(sel.clone()),
}
}
let opts = crate::GetOpts::new(spec.timeout).max_replies(spec.max_replies);
let gets = futures_util::future::join_all(selectors.iter().map(|selector| {
let opts = &opts;
async move {
(
selector.clone(),
crate::bus::query::snapshot_get(session, selector, opts).await,
)
}
}))
.await;
let mut values = Vec::new();
let mut errors = 0u64;
for (selector, replies) in gets {
match replies {
Ok(r) => {
errors += r.errors;
values.extend(r.values);
}
Err(e) => {
tracing::warn!(selector, error = %e, "preamble GET could not be issued");
failed.push(selector);
}
}
}
let (kept, _superseded) = crate::model::snapshot::fold_latest(values);
let in_ring: BTreeSet<&str> = ring.iter().map(|v| v.key.as_str()).collect();
let rows: Vec<Arc<SampleView>> = kept
.into_values()
.filter(|(view, _)| match semantics {
PreambleSemantics::AbsentFromWindow => !in_ring.contains(view.key.as_str()),
PreambleSemantics::Full => true,
})
.map(|(view, _)| Arc::new(view))
.collect();
Some((
PreambleInfo {
count: rows.len() as u64,
collected_over_s: started.elapsed().as_secs_f64(),
selectors,
semantics,
incomplete: errors + opts.elided(),
failed,
},
rows,
))
};
let mut preamble = std::pin::pin!(preamble);
let fetched = loop {
tokio::select! {
item = events.recv() => drain_into(item, &mut buffered, &mut buffer_dropped),
fetched = &mut preamble => break fetched,
}
};
let (preamble_info, preamble_rows) = match fetched {
Some((info, rows)) => (Some(info), rows),
None => (None, Vec::new()),
};
if let Some(info) = &preamble_info {
on_event(TriggerEvent::Preamble(info));
}
let header = ZrecHeader {
zrec: ZREC_VERSION,
selectors: watched.clone(),
base: base.to_string(),
captured_at,
preamble: preamble_info.clone(),
pre_roll: Some(pre_roll.clone()),
};
let out = match open().await {
Ok(out) => out,
Err(e) => {
if let Err(teardown) = monitor.shutdown().await {
tracing::warn!("after a failed open: {teardown}");
}
return Err(e);
}
};
let sink = ZrecSink::spawn_at(out, &header, epoch).await?;
for row in &preamble_rows {
sink.write_preamble(Arc::clone(row)).await?;
}
for view in ring.iter() {
sink.write_sample(Arc::clone(view)).await?;
}
sink.write_trigger(trigger.clone()).await?;
while let Ok(item) = tokio::time::timeout(Duration::ZERO, events.recv()).await {
drain_into(item, &mut buffered, &mut buffer_dropped);
}
for item in buffered.drain(..) {
match item {
StreamItem::Event(FleetEvent::Sample(s)) => sink.write_sample(s).await?,
StreamItem::Dropped(n) => sink.write_dropped(n).await?,
StreamItem::Event(_) => {}
}
}
if buffer_dropped > 0 {
sink.write_dropped(buffer_dropped).await?;
}
let pre_written = sink.counts().samples;
let bounds = RecordBounds {
max_samples: spec.max_samples.map(|n| n + pre_written),
max_duration: Some(spec.post),
};
let mut last_line = Instant::now();
let recorded = record(&mut events, &sink, bounds, |samples, dropped| {
if last_line.elapsed() >= Duration::from_secs(1) {
last_line = Instant::now();
on_event(TriggerEvent::Progress { samples, dropped });
}
})
.await;
let closed = monitor.shutdown().await;
recorded?;
let counts = sink.finish().await?;
closed?;
Ok(RecordReport {
header,
out: None,
samples: counts.samples,
dropped: counts.dropped,
duration_ms: u64::try_from((stats.span + fired_at.elapsed()).as_millis())
.unwrap_or(u64::MAX),
trigger: Some(trigger),
preamble: preamble_info,
pre_roll: Some(pre_roll),
preamble_rows: counts.preamble,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_watch_cover_declares_no_included_selector_twice() {
let s = |v: &[&str]| v.iter().map(|s| s.to_string()).collect::<Vec<_>>();
assert_eq!(
watch_cover(&s(&["v1/**", "v1/h-aaaaaaaaaaaa/state/p/health"])),
s(&["v1/**"])
);
assert_eq!(
watch_cover(&s(&["v1/h-aaaaaaaaaaaa/state/p/health", "v1/**"])),
s(&["v1/**"])
);
assert_eq!(
watch_cover(&s(&["v1/a/state/**", "v1/b/state/**", "v1/a/state/**"])),
s(&["v1/a/state/**", "v1/b/state/**"])
);
assert_eq!(watch_cover(&s(&["v1/x/**", "v1/x/**"])), s(&["v1/x/**"]));
}
#[test]
fn the_state_projection_narrows_widens_or_declines() {
let p = |s| state_projection("", s);
assert_eq!(p("v1/**").as_deref(), Some("v1/*/state/**"));
assert_eq!(p("**").as_deref(), Some("v1/*/state/**"));
assert_eq!(
p("v1/h-aaaaaaaaaaaa/**").as_deref(),
Some("v1/h-aaaaaaaaaaaa/state/**")
);
assert_eq!(
p("v1/h-aaaaaaaaaaaa/*/demo/**").as_deref(),
Some("v1/h-aaaaaaaaaaaa/state/demo/**")
);
assert_eq!(
p("v1/h-aaaaaaaaaaaa/state/demo/health").as_deref(),
Some("v1/h-aaaaaaaaaaaa/state/demo/health")
);
assert_eq!(
p("v1/h-aaaaaaaaaaaa/state").as_deref(),
Some("v1/h-aaaaaaaaaaaa/state/**")
);
assert_eq!(p("v1/h-aaaaaaaaaaaa/telemetry/demo/**"), None);
assert_eq!(p("v1/h-aaaaaaaaaaaa"), None);
assert_eq!(p("v1"), None);
assert_eq!(
state_projection("acme", "acme/v1/**").as_deref(),
Some("acme/v1/*/state/**")
);
assert_eq!(state_projection("acme", "other/v1/**"), None);
}
}