use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use observation::{Event, EventScope};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::long_tier::{LongTierConfig, LongTierError, LongTierStore, ObjectStore, ObjectStoreError};
type ScopedRow = (EventScope, Event);
pub const SNAPSHOT_PREFIX: &str = "analytics/snapshots/";
pub const POINTER_KEY: &str = "analytics/current.json";
pub const DEFAULT_SNAPSHOT_INTERVAL: Duration = Duration::from_secs(300);
pub const DEFAULT_EVENT_LIMIT: usize = 200;
#[derive(Debug, thiserror::Error)]
pub enum SnapshotError {
#[error("long tier: {0}")]
LongTier(#[from] LongTierError),
#[error("object store: {0}")]
ObjectStore(#[from] ObjectStoreError),
#[error("encode: {0}")]
Encode(#[from] serde_json::Error),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SnapshotLevelCount {
pub level: String,
pub count: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SnapshotBucket {
pub key: String,
pub count: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SnapshotEvent {
pub offset_ms: u32,
pub level: String,
pub target: String,
pub msg: String,
pub scope_kind: String,
pub scope_id: String,
pub fields: serde_json::Value,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AnalyticsSnapshot {
pub generated_at_ms: u64,
pub window_start_ms: u64,
pub window_end_ms: u64,
pub total_events: u64,
pub by_level: Vec<SnapshotLevelCount>,
pub group_by: String,
pub timeseries: Vec<SnapshotBucket>,
pub events: Vec<SnapshotEvent>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SnapshotPointer {
pub hash: String,
pub blob_key: String,
pub generated_at_ms: u64,
pub total_events: u64,
}
#[derive(Debug, Clone)]
pub struct SnapshotConfig {
pub machines: Vec<String>,
pub retention_ms: u64,
pub group_by: String,
pub event_limit: usize,
pub min_level: String,
pub interval: Duration,
}
impl SnapshotConfig {
pub fn new(machines: Vec<String>, retention_ms: u64) -> Self {
Self {
machines,
retention_ms,
group_by: "hour".to_string(),
event_limit: DEFAULT_EVENT_LIMIT,
min_level: "info".to_string(),
interval: DEFAULT_SNAPSHOT_INTERVAL,
}
}
pub fn with_interval(mut self, interval: Duration) -> Self {
self.interval = interval;
self
}
pub fn with_group_by(mut self, group_by: impl Into<String>) -> Self {
self.group_by = group_by.into();
self
}
pub fn with_min_level(mut self, min_level: impl Into<String>) -> Self {
self.min_level = min_level.into();
self
}
pub fn with_event_limit(mut self, limit: usize) -> Self {
self.event_limit = limit;
self
}
}
pub struct SnapshotProducer {
object_store: Arc<dyn ObjectStore>,
cfg: SnapshotConfig,
}
impl SnapshotProducer {
pub fn new(object_store: Arc<dyn ObjectStore>, cfg: SnapshotConfig) -> Self {
Self { object_store, cfg }
}
pub fn build_snapshot(&self) -> Result<AnalyticsSnapshot, SnapshotError> {
let until_ms = self.cfg.retention_ms;
let mut events: Vec<ScopedRow> = Vec::new();
for machine in &self.cfg.machines {
let lt = LongTierStore::new(
LongTierConfig { machine_id: machine.clone(), retention_ms: self.cfg.retention_ms },
Arc::clone(&self.object_store),
);
events.extend(lt.query_range(None, 0, until_ms)?);
}
let total_events = events.len() as u64;
let by_level = aggregate_by_level(&events);
let timeseries = aggregate_buckets(&events, &self.cfg.group_by);
let recent = recent_events(&events, &self.cfg.min_level, self.cfg.event_limit);
Ok(AnalyticsSnapshot {
generated_at_ms: now_epoch_ms(),
window_start_ms: 0,
window_end_ms: until_ms,
total_events,
by_level,
group_by: self.cfg.group_by.clone(),
timeseries,
events: recent,
})
}
pub fn publish(&self, snapshot: &AnalyticsSnapshot) -> Result<String, SnapshotError> {
let body = serde_json::to_vec(snapshot)?;
let hash = sha256_hex(&body);
let blob_key = format!("{SNAPSHOT_PREFIX}{hash}.json");
self.object_store.put(&blob_key, body)?;
let pointer = SnapshotPointer {
hash: hash.clone(),
blob_key,
generated_at_ms: snapshot.generated_at_ms,
total_events: snapshot.total_events,
};
self.object_store.put(POINTER_KEY, serde_json::to_vec(&pointer)?)?;
Ok(hash)
}
pub fn run_once(&self) -> Result<String, SnapshotError> {
let snapshot = self.build_snapshot()?;
self.publish(&snapshot)
}
pub fn spawn(self) -> tokio::task::JoinHandle<()> {
let producer = Arc::new(self);
tokio::spawn(async move {
let mut ticker = tokio::time::interval(producer.cfg.interval);
loop {
ticker.tick().await;
let p = Arc::clone(&producer);
match tokio::task::spawn_blocking(move || p.run_once()).await {
Ok(Ok(hash)) => {
eprintln!("scryer snapshot: published analytics snapshot {hash}");
}
Ok(Err(e)) => eprintln!("scryer snapshot: publish error: {e}"),
Err(e) => eprintln!("scryer snapshot: pass panicked: {e}"),
}
}
})
}
}
fn level_rank(level: &str) -> u8 {
match level {
"trace" => 0,
"debug" => 1,
"info" => 2,
"warn" => 3,
"error" => 4,
"fatal" => 5,
_ => 2,
}
}
fn aggregate_by_level(events: &[ScopedRow]) -> Vec<SnapshotLevelCount> {
let mut counts: std::collections::HashMap<String, u64> = std::collections::HashMap::new();
for (_scope, ev) in events {
*counts.entry(ev.level.as_str().to_string()).or_insert(0) += 1;
}
let mut rows: Vec<SnapshotLevelCount> = counts
.into_iter()
.map(|(level, count)| SnapshotLevelCount { level, count })
.collect();
rows.sort_by(|a, b| b.count.cmp(&a.count).then(a.level.cmp(&b.level)));
rows
}
fn aggregate_buckets(events: &[ScopedRow], group_by: &str) -> Vec<SnapshotBucket> {
let mut counts: std::collections::HashMap<String, u64> = std::collections::HashMap::new();
for (_scope, ev) in events {
let key = match group_by {
"level" => ev.level.as_str().to_string(),
"target" => ev.target.splitn(2, "::").next().unwrap_or(&ev.target).to_string(),
"hour" => format!("h{}", ev.offset_ms / 3_600_000),
_ => ev.level.as_str().to_string(),
};
*counts.entry(key).or_insert(0) += 1;
}
let mut buckets: Vec<SnapshotBucket> = counts
.into_iter()
.map(|(key, count)| SnapshotBucket { key, count })
.collect();
buckets.sort_by(|a, b| b.count.cmp(&a.count).then(a.key.cmp(&b.key)));
buckets
}
fn recent_events(events: &[ScopedRow], min_level: &str, limit: usize) -> Vec<SnapshotEvent> {
let floor = level_rank(min_level);
let mut rows: Vec<&ScopedRow> = events
.iter()
.filter(|(_scope, ev)| level_rank(ev.level.as_str()) >= floor)
.collect();
rows.sort_by(|(_, a), (_, b)| b.offset_ms.cmp(&a.offset_ms));
rows.into_iter()
.take(limit)
.map(|(scope, ev)| SnapshotEvent {
offset_ms: ev.offset_ms,
level: ev.level.as_str().to_string(),
target: ev.target.clone(),
msg: ev.msg.clone(),
scope_kind: scope.kind_str().to_string(),
scope_id: scope.id_str(),
fields: ev.fields.clone(),
})
.collect()
}
fn sha256_hex(data: &[u8]) -> String {
let mut h = Sha256::new();
h.update(data);
hex::encode(h.finalize())
}
fn now_epoch_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::long_tier::{InMemoryObjectStore, LongTierConfig, LongTierStore, MS_PER_DAY, ObjectStore};
use crate::service::{Scryer, ScryerConfig};
use observation::{Event, EventScope, EventSource, Level, TaskRunId};
use serde_json::json;
use tempfile::TempDir;
use workload_spec::MeshIdent;
const RETENTION_MS: u64 = 7 * MS_PER_DAY;
fn make_event(run_id: &TaskRunId, seq: u32, offset_ms: u32, level: Level) -> Event {
Event {
run_id: run_id.clone(),
seq,
offset_ms,
level,
target: format!("cargo::stage{}", seq % 3),
msg: format!("msg {seq}"),
fields: json!({ "seq": seq }),
anchor: None,
source: EventSource::Synth,
}
}
fn promote_into(
store: &Arc<InMemoryObjectStore>,
machine: &str,
scope: &EventScope,
events: Vec<Event>,
) {
let dir = TempDir::new().unwrap();
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let scryer = Scryer::new(cfg, None).unwrap();
let items: Vec<(EventScope, Event)> =
events.into_iter().map(|e| (scope.clone(), e)).collect();
scryer.store().insert_events(&items).unwrap();
let lt = LongTierStore::new(
LongTierConfig { machine_id: machine.to_string(), retention_ms: RETENTION_MS },
Arc::clone(store) as Arc<dyn ObjectStore>,
);
lt.rollover(scryer.store(), RETENTION_MS).unwrap();
}
#[test]
fn build_and_publish_single_machine() {
let store = Arc::new(InMemoryObjectStore::new());
let scope = EventScope::Service(MeshIdent("svc.prod".to_string()));
let run_id = TaskRunId::new();
let one_day = MS_PER_DAY as u32;
let events = vec![
make_event(&run_id, 0, one_day, Level::Warn),
make_event(&run_id, 1, one_day, Level::Warn),
make_event(&run_id, 2, one_day, Level::Warn),
make_event(&run_id, 3, one_day, Level::Error),
make_event(&run_id, 4, one_day, Level::Error),
make_event(&run_id, 5, one_day, Level::Info),
];
promote_into(&store, "m1", &scope, events);
let producer = SnapshotProducer::new(
Arc::clone(&store) as Arc<dyn ObjectStore>,
SnapshotConfig::new(vec!["m1".to_string()], RETENTION_MS),
);
let snap = producer.build_snapshot().unwrap();
assert_eq!(snap.total_events, 6);
assert_eq!(
snap.by_level,
vec![
SnapshotLevelCount { level: "warn".into(), count: 3 },
SnapshotLevelCount { level: "error".into(), count: 2 },
SnapshotLevelCount { level: "info".into(), count: 1 },
]
);
assert_eq!(snap.group_by, "hour");
assert_eq!(snap.timeseries, vec![SnapshotBucket { key: "h24".into(), count: 6 }]);
assert_eq!(snap.events.len(), 6);
assert_eq!(snap.window_end_ms, RETENTION_MS);
let hash = producer.publish(&snap).unwrap();
let blob_key = format!("{SNAPSHOT_PREFIX}{hash}.json");
assert!(store.contains_key(&blob_key), "content-addressed blob present");
assert!(store.contains_key(POINTER_KEY), "pointer present");
let ptr: SnapshotPointer =
serde_json::from_slice(&store.get(POINTER_KEY).unwrap().unwrap()).unwrap();
assert_eq!(ptr.hash, hash);
assert_eq!(ptr.blob_key, blob_key);
assert_eq!(ptr.total_events, 6);
let round_trip: AnalyticsSnapshot =
serde_json::from_slice(&store.get(&blob_key).unwrap().unwrap()).unwrap();
assert_eq!(round_trip, snap);
}
#[test]
fn aggregates_across_machines() {
let store = Arc::new(InMemoryObjectStore::new());
let one_day = MS_PER_DAY as u32;
let run_a = TaskRunId::new();
let run_b = TaskRunId::new();
promote_into(
&store,
"m1",
&EventScope::Service(MeshIdent("svc.a".to_string())),
vec![make_event(&run_a, 0, one_day, Level::Error)],
);
promote_into(
&store,
"m2",
&EventScope::Service(MeshIdent("svc.b".to_string())),
vec![
make_event(&run_b, 0, one_day, Level::Error),
make_event(&run_b, 1, one_day, Level::Info),
],
);
let producer = SnapshotProducer::new(
Arc::clone(&store) as Arc<dyn ObjectStore>,
SnapshotConfig::new(vec!["m1".to_string(), "m2".to_string()], RETENTION_MS),
);
let snap = producer.build_snapshot().unwrap();
assert_eq!(snap.total_events, 3);
assert_eq!(
snap.by_level,
vec![
SnapshotLevelCount { level: "error".into(), count: 2 },
SnapshotLevelCount { level: "info".into(), count: 1 },
]
);
let mut scope_ids: Vec<&str> = snap.events.iter().map(|e| e.scope_id.as_str()).collect();
scope_ids.sort();
scope_ids.dedup();
assert_eq!(scope_ids, vec!["svc.a", "svc.b"]);
}
#[test]
fn min_level_filters_events_not_rollup() {
let store = Arc::new(InMemoryObjectStore::new());
let scope = EventScope::Service(MeshIdent("svc.filter".to_string()));
let run_id = TaskRunId::new();
let one_day = MS_PER_DAY as u32;
promote_into(
&store,
"m1",
&scope,
vec![
make_event(&run_id, 0, one_day, Level::Debug),
make_event(&run_id, 1, one_day, Level::Info),
make_event(&run_id, 2, one_day, Level::Error),
],
);
let producer = SnapshotProducer::new(
Arc::clone(&store) as Arc<dyn ObjectStore>,
SnapshotConfig::new(vec!["m1".to_string()], RETENTION_MS).with_min_level("warn"),
);
let snap = producer.build_snapshot().unwrap();
assert_eq!(snap.total_events, 3);
assert_eq!(snap.events.len(), 1);
assert_eq!(snap.events[0].level, "error");
}
#[test]
fn empty_corpus_publishes_empty_snapshot() {
let store = Arc::new(InMemoryObjectStore::new());
let producer = SnapshotProducer::new(
Arc::clone(&store) as Arc<dyn ObjectStore>,
SnapshotConfig::new(vec!["m1".to_string()], RETENTION_MS),
);
let snap = producer.build_snapshot().unwrap();
assert_eq!(snap.total_events, 0);
assert!(snap.by_level.is_empty());
assert!(snap.timeseries.is_empty());
assert!(snap.events.is_empty());
let hash = producer.publish(&snap).unwrap();
assert!(store.contains_key(POINTER_KEY));
assert!(store.contains_key(&format!("{SNAPSHOT_PREFIX}{hash}.json")));
}
#[test]
fn run_once_builds_and_publishes() {
let store = Arc::new(InMemoryObjectStore::new());
let scope = EventScope::Service(MeshIdent("svc.once".to_string()));
let run_id = TaskRunId::new();
promote_into(
&store,
"m1",
&scope,
vec![make_event(&run_id, 0, MS_PER_DAY as u32, Level::Info)],
);
let producer = SnapshotProducer::new(
Arc::clone(&store) as Arc<dyn ObjectStore>,
SnapshotConfig::new(vec!["m1".to_string()], RETENTION_MS),
);
let hash = producer.run_once().unwrap();
let ptr: SnapshotPointer =
serde_json::from_slice(&store.get(POINTER_KEY).unwrap().unwrap()).unwrap();
assert_eq!(ptr.hash, hash);
assert_eq!(ptr.total_events, 1);
}
}