use crate::federation::{FederationPeer, FederationRule, ScopedEvent};
use crate::long_tier::{LongTierError, LongTierStore};
use crate::ring::{EventRing, RingConfig};
use crate::store::{EventStore, ScryerStoreError, ScopeFilter, ScopeInfo};
use observation::{Event, EventScope, HopRollup, Level};
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use task_runs::{EventFilter as TaskEventFilter, StoreError as TaskStoreError, TaskStore};
use thiserror::Error;
use tokio::sync::broadcast;
pub const MAX_HOP_BUCKETS: u64 = 240;
#[derive(Debug, Clone)]
pub struct ScryerConfig {
pub db_path: PathBuf,
pub ring: RingConfig,
pub subscriber_capacity: usize,
pub retention_ms: u64,
}
impl ScryerConfig {
pub fn new(db_path: impl AsRef<Path>) -> Self {
Self {
db_path: db_path.as_ref().to_owned(),
ring: RingConfig::default(),
subscriber_capacity: 256,
retention_ms: 7 * 24 * 3600 * 1000,
}
}
}
#[derive(Debug, Error)]
pub enum ScryerError {
#[error("store: {0}")]
Store(#[from] ScryerStoreError),
#[error("task-runs store: {0}")]
TaskStore(#[from] TaskStoreError),
#[error("long tier: {0}")]
LongTier(#[from] LongTierError),
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct EventFilter {
pub min_level: Option<Level>,
pub target: Option<String>,
pub offset_range: Option<(u32, u32)>,
pub seq_range: Option<(u32, u32)>,
pub limit: Option<u32>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct QueryCursor {
pub ring_cursor: u64,
pub last_seq: u32,
}
impl QueryCursor {
pub fn beginning() -> Self {
Self { ring_cursor: 0, last_seq: 0 }
}
}
pub struct TailResult {
pub events: Vec<Event>,
pub next_cursor: QueryCursor,
}
#[derive(Debug, Clone)]
pub struct AggregateBucket {
pub key: String,
pub count: u64,
}
pub struct Scryer {
ring: Arc<EventRing>,
store: Arc<EventStore>,
task_store: Option<Arc<TaskStore>>,
tx: broadcast::Sender<(EventScope, Event)>,
_cfg: ScryerConfig,
long_tier: Option<(Arc<LongTierStore>, u64)>,
}
impl Scryer {
pub fn new(cfg: ScryerConfig, task_store: Option<Arc<TaskStore>>) -> Result<Self, ScryerError> {
let ring = EventRing::new(cfg.ring.clone());
let store = Arc::new(EventStore::open(&cfg.db_path)?);
let (tx, _) = broadcast::channel(cfg.subscriber_capacity);
Ok(Self { ring, store, task_store, tx, _cfg: cfg, long_tier: None })
}
pub fn with_long_tier(
mut self,
store: Arc<LongTierStore>,
short_disk_retention_ms: u64,
) -> Self {
self.long_tier = Some((store, short_disk_retention_ms));
self
}
pub fn push(&self, scope: EventScope, event: Event) -> Result<(), ScryerError> {
let _ = self.tx.send((scope.clone(), event.clone()));
let (_cursor, should_flush) = self.ring.push(scope, event);
if should_flush {
self.flush_ring()?;
}
Ok(())
}
pub fn flush_ring(&self) -> Result<(), ScryerError> {
let pending = self.ring.take_pending();
if !pending.is_empty() {
self.store.insert_events(&pending)?;
}
Ok(())
}
pub async fn events(
&self,
scope: &EventScope,
filter: &EventFilter,
) -> Result<Vec<Event>, ScryerError> {
match (scope, &self.task_store) {
(EventScope::TaskRun(run_id), Some(ts)) => {
let tf = TaskEventFilter {
target: filter.target.clone(),
min_level: filter.min_level,
seq_range: filter.seq_range.map(|(lo, hi)| lo..hi),
offset_range: filter.offset_range,
field_filter: None,
limit: filter.limit,
};
let events = ts.query_events(run_id, &tf).await?;
Ok(events)
}
_ => {
let sf = ScopeFilter {
scope: scope.clone(),
min_level: filter.min_level,
target_prefix: filter.target.clone(),
offset_range: filter.offset_range,
seq_range: filter.seq_range,
limit: filter.limit.map(|l| l as usize),
};
let rows = self.store.query_events(&sf)?;
Ok(rows.into_iter().map(|(_, ev)| ev).collect())
}
}
}
pub async fn tail(
&self,
scope: &EventScope,
cursor: QueryCursor,
limit: usize,
) -> Result<TailResult, ScryerError> {
let limit = limit.min(1000);
let (events, next_ring) = self.ring.tail_since(scope, cursor.ring_cursor, limit);
let events = if events.is_empty() {
if let (EventScope::TaskRun(run_id), Some(ts)) = (scope, &self.task_store) {
let tf = TaskEventFilter {
seq_range: Some(cursor.last_seq..u32::MAX),
limit: Some(limit as u32),
..Default::default()
};
ts.query_events(run_id, &tf).await?
} else {
events
}
} else {
events
};
let next_seq = events.last().map(|e| e.seq + 1).unwrap_or(cursor.last_seq);
Ok(TailResult {
events,
next_cursor: QueryCursor { ring_cursor: next_ring, last_seq: next_seq },
})
}
pub fn subscribe(&self) -> broadcast::Receiver<(EventScope, Event)> {
self.tx.subscribe()
}
pub async fn federated_events(
&self,
scope: &EventScope,
filter: &EventFilter,
rule: &FederationRule,
peers: &[Arc<dyn FederationPeer>],
) -> Result<Vec<ScopedEvent>, ScryerError> {
let local = ScopedEvent::tag_all(scope, self.events(scope, filter).await?);
Ok(crate::federation::federated_events(local, peers, filter, rule).await)
}
pub fn list_scopes(&self, limit: usize) -> Result<Vec<ScopeInfo>, ScryerError> {
Ok(self.store.list_scopes(limit)?)
}
pub async fn events_all_scoped(
&self,
filter: &EventFilter,
) -> Result<Vec<ScopedEvent>, ScryerError> {
let scope_limit = filter.limit.map(|l| l as usize).unwrap_or(1000);
let mut merged: Vec<ScopedEvent> = Vec::new();
for info in self.store.list_scopes(1000)? {
let local = ScopedEvent::tag_all(&info.scope, self.events(&info.scope, filter).await?);
merged = crate::federation::merge_ordered(merged, local);
if merged.len() >= scope_limit {
merged.truncate(scope_limit);
break;
}
}
Ok(merged)
}
pub async fn events_all(
&self,
filter: &EventFilter,
) -> Result<Vec<Event>, ScryerError> {
Ok(ScopedEvent::into_events(self.events_all_scoped(filter).await?))
}
pub fn aggregate_all(
&self,
since_ms: u64,
group_by: &str,
) -> Result<Vec<AggregateBucket>, ScryerError> {
let mut counts: std::collections::HashMap<String, u64> =
std::collections::HashMap::new();
for info in self.store.list_scopes(1000)? {
for bucket in self.aggregate(&info.scope, since_ms, group_by)? {
*counts.entry(bucket.key).or_insert(0) += bucket.count;
}
}
let mut buckets: Vec<AggregateBucket> = counts
.into_iter()
.map(|(key, count)| AggregateBucket { key, count })
.collect();
buckets.sort_by(|a, b| b.count.cmp(&a.count).then(a.key.cmp(&b.key)));
Ok(buckets)
}
pub fn aggregate(
&self,
scope: &EventScope,
since_ms: u64,
group_by: &str,
) -> Result<Vec<AggregateBucket>, ScryerError> {
let mut all_events = Vec::new();
let boundary_ms = self.long_tier.as_ref().map_or(since_ms, |(_, ret)| {
since_ms.max(*ret)
});
let sf = ScopeFilter {
scope: scope.clone(),
offset_range: Some((boundary_ms as u32, u32::MAX)),
..ScopeFilter::for_scope(scope.clone())
};
let disk_rows = self.store.query_events(&sf)?;
all_events.extend(disk_rows.into_iter().map(|(_, e)| e));
if let Some((lt, ret_ms)) = &self.long_tier {
let until_ms = (*ret_ms).min(boundary_ms);
if since_ms < until_ms {
let lt_rows = lt.query_range(Some(scope), since_ms, until_ms)?;
all_events.extend(lt_rows.into_iter().map(|(_, e)| e));
}
}
let mut counts: std::collections::HashMap<String, u64> =
std::collections::HashMap::new();
for ev in &all_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<AggregateBucket> = counts
.into_iter()
.map(|(key, count)| AggregateBucket { key, count })
.collect();
buckets.sort_by(|a, b| b.count.cmp(&a.count).then(a.key.cmp(&b.key)));
Ok(buckets)
}
pub fn hop_rollups(
&self,
from_unix_nanos: u64,
to_unix_nanos: u64,
bucket_nanos: Option<u64>,
) -> Result<Vec<HopRollup>, ScryerError> {
let Some(requested) = bucket_nanos else {
return Ok(self
.store
.query_hop_rollups(None, from_unix_nanos, to_unix_nanos)?);
};
if to_unix_nanos <= from_unix_nanos {
return Ok(Vec::new());
}
let span = to_unix_nanos - from_unix_nanos;
let bucket = requested
.max(span.div_ceil(MAX_HOP_BUCKETS))
.max(observation::ROLLUP_WINDOW_NANOS)
.div_ceil(observation::ROLLUP_WINDOW_NANOS)
* observation::ROLLUP_WINDOW_NANOS;
let mut out = Vec::new();
let mut start = observation::rollup_window_start(from_unix_nanos);
while start < to_unix_nanos {
let end = start + bucket;
out.extend(self.store.query_hop_rollups(None, start, end)?);
start = end;
}
Ok(out)
}
pub fn ring(&self) -> &Arc<EventRing> {
&self.ring
}
pub fn store(&self) -> &Arc<EventStore> {
&self.store
}
}
#[cfg(test)]
mod tests {
use super::*;
use observation::{EventSource, TaskRunId};
use serde_json::json;
use tempfile::TempDir;
fn make_event(run_id: TaskRunId, seq: u32) -> Event {
Event {
run_id,
seq,
offset_ms: seq * 10,
level: Level::Info,
target: "test".to_string(),
msg: format!("event {seq}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
}
}
fn open_scryer(dir: &TempDir) -> Scryer {
let cfg = ScryerConfig::new(dir.path().join("events.db"));
Scryer::new(cfg, None).unwrap()
}
#[tokio::test]
async fn scryer_push_and_events_via_store() {
let dir = TempDir::new().unwrap();
let scryer = open_scryer(&dir);
let run_id = TaskRunId::new();
let scope = EventScope::TaskRun(run_id.clone());
for i in 0..5 {
let ev = make_event(run_id.clone(), i);
scryer.push(scope.clone(), ev).unwrap();
}
scryer.flush_ring().unwrap();
let events = scryer.events(&scope, &EventFilter::default()).await.unwrap();
assert_eq!(events.len(), 5);
}
#[tokio::test]
async fn scryer_tail_from_ring() {
let dir = TempDir::new().unwrap();
let scryer = open_scryer(&dir);
let run_id = TaskRunId::new();
let scope = EventScope::TaskRun(run_id.clone());
for i in 0..5 {
scryer.push(scope.clone(), make_event(run_id.clone(), i)).unwrap();
}
let result = scryer.tail(&scope, QueryCursor::beginning(), 100).await.unwrap();
assert_eq!(result.events.len(), 5);
assert_eq!(result.next_cursor.ring_cursor, 5);
}
#[test]
fn scryer_subscribe_receives_events() {
let dir = TempDir::new().unwrap();
let scryer = open_scryer(&dir);
let run_id = TaskRunId::new();
let scope = EventScope::TaskRun(run_id.clone());
let mut rx = scryer.subscribe();
scryer.push(scope.clone(), make_event(run_id.clone(), 0)).unwrap();
let (recv_scope, recv_ev) = rx.try_recv().unwrap();
assert_eq!(recv_scope, scope);
assert_eq!(recv_ev.seq, 0);
}
mod events {
use super::*;
use observation::{EventSource, ForgeId, Level};
use serde_json::json;
#[tokio::test]
async fn forge_scope() {
let dir = TempDir::new().unwrap();
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let scryer = Scryer::new(cfg, None).unwrap();
let forge_id = ForgeId::new();
let scope = EventScope::Forge(forge_id.clone());
for i in 0u32..4 {
let ev = Event {
run_id: forge_id.clone().into(),
seq: i,
offset_ms: i * 10,
level: Level::Info,
target: "forge::test".to_string(),
msg: format!("forge event {i}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
};
scryer.push(scope.clone(), ev).unwrap();
}
scryer.flush_ring().unwrap();
let events = scryer.events(&scope, &EventFilter::default()).await.unwrap();
assert_eq!(events.len(), 4, "expected 4 forge events");
assert_eq!(events[0].target, "forge::test");
assert_eq!(events[3].seq, 3);
let other_scope = EventScope::Forge(ForgeId::new());
let other = scryer.events(&other_scope, &EventFilter::default()).await.unwrap();
assert!(other.is_empty(), "different forge id must return no events");
}
}
#[tokio::test]
async fn scryer_events_taskrun_matches_task_events() {
use task_runs::{EventFilter as TF, Initiator, RunStatus, TaskRunMeta, TaskStore};
let dir = TempDir::new().unwrap();
let task_store = Arc::new(TaskStore::open(&dir.path().join("task-runs.db")).await.unwrap());
let run_id = TaskRunId::new();
let meta = TaskRunMeta {
id: run_id.clone(),
command: "echo hi".to_string(),
cwd: dir.path().to_path_buf(),
env: vec![],
started_at: 1000,
status: RunStatus::Pending,
label: None,
initiator: Initiator::Human { camp: "test".to_string() },
beholder_status: None,
pinned: false,
origin: None,
host_pid: None,
};
task_store.insert_run(&meta).await.unwrap();
for i in 0u32..5 {
use observation::{EventSource, Level};
use serde_json::json;
task_store
.append_event(
&run_id,
i * 10,
Level::Info,
"test",
&format!("event {i}"),
&json!({}),
None,
&EventSource::Synth,
)
.await
.unwrap();
}
let cfg = ScryerConfig::new(dir.path().join("scryer-events.db"));
let scryer = Scryer::new(cfg, Some(task_store.clone())).unwrap();
let scope = EventScope::TaskRun(run_id.clone());
let scryer_events = scryer.events(&scope, &EventFilter::default()).await.unwrap();
let task_events = task_store.query_events(&run_id, &TF::default()).await.unwrap();
assert_eq!(scryer_events.len(), task_events.len());
for (se, te) in scryer_events.iter().zip(task_events.iter()) {
assert_eq!(se.seq, te.seq);
assert_eq!(se.msg, te.msg);
}
}
mod aggregate {
use super::*;
use crate::long_tier::{InMemoryObjectStore, LongTierConfig, LongTierStore, MS_PER_DAY};
use observation::{EventSource, Level, TaskRunId};
use serde_json::json;
use workload_spec::MeshIdent;
fn make_svc_event(run_id: &TaskRunId, seq: u32, offset_ms: u32, level: Level) -> Event {
Event {
run_id: run_id.clone(),
seq,
offset_ms,
level,
target: "svc::db".to_string(),
msg: format!("event {seq}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
}
}
#[tokio::test]
async fn cross_boundary() {
let dir = TempDir::new().unwrap();
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let scryer = Scryer::new(cfg, None).unwrap();
let scope = EventScope::Service(MeshIdent("svc.prod".to_string()));
let run_id = TaskRunId::new();
let old_offset = MS_PER_DAY as u32; for i in 0u32..3 {
scryer
.push(scope.clone(), make_svc_event(&run_id, i, old_offset, Level::Warn))
.unwrap();
}
scryer.flush_ring().unwrap();
let cutoff_ms = 7 * MS_PER_DAY; let obj_store = Arc::new(InMemoryObjectStore::new());
let lt_cfg = LongTierConfig { machine_id: "m1".to_string(), retention_ms: cutoff_ms };
let lt = Arc::new(LongTierStore::new(
lt_cfg,
Arc::clone(&obj_store) as Arc<dyn crate::long_tier::ObjectStore>,
));
let promoted = lt.rollover(scryer.store(), cutoff_ms).unwrap();
assert_eq!(promoted, 3, "3 old events promoted to long tier");
let after_rollover = scryer.events(&scope, &EventFilter::default()).await.unwrap();
assert!(after_rollover.is_empty(), "old events removed from short-disk");
let recent_offset = (10 * MS_PER_DAY) as u32;
for i in 3u32..5 {
scryer
.push(scope.clone(), make_svc_event(&run_id, i, recent_offset, Level::Info))
.unwrap();
}
scryer.flush_ring().unwrap();
let dir2 = TempDir::new().unwrap();
let cfg2 = ScryerConfig::new(dir2.path().join("events2.db"));
let scryer2 = Scryer::new(cfg2, None)
.unwrap()
.with_long_tier(Arc::clone(<), cutoff_ms);
let recent_events: Vec<(EventScope, Event)> = (3u32..5)
.map(|i| {
(scope.clone(), make_svc_event(&run_id, i, recent_offset, Level::Info))
})
.collect();
scryer2.store().insert_events(&recent_events).unwrap();
let buckets = scryer2.aggregate(&scope, 0, "level").unwrap();
let total: u64 = buckets.iter().map(|b| b.count).sum();
assert_eq!(total, 5, "3 old (long-tier) + 2 recent (short-disk) = 5 total");
let warn_count = buckets
.iter()
.find(|b| b.key == "warn")
.map(|b| b.count)
.unwrap_or(0);
let info_count = buckets
.iter()
.find(|b| b.key == "info")
.map(|b| b.count)
.unwrap_or(0);
assert_eq!(warn_count, 3, "3 Warn events from long-tier");
assert_eq!(info_count, 2, "2 Info events from short-disk");
}
#[test]
fn short_disk_only() {
let dir = TempDir::new().unwrap();
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let scryer = Scryer::new(cfg, None).unwrap();
let scope = EventScope::Service(MeshIdent("svc.local".to_string()));
let run_id = TaskRunId::new();
for i in 0u32..4 {
let offset = if i < 2 { 1000u32 } else { 2000u32 };
let level = if i < 2 { Level::Error } else { Level::Info };
scryer.push(scope.clone(), make_svc_event(&run_id, i, offset, level)).unwrap();
}
scryer.flush_ring().unwrap();
let buckets = scryer.aggregate(&scope, 0, "level").unwrap();
let total: u64 = buckets.iter().map(|b| b.count).sum();
assert_eq!(total, 4);
}
}
}