use std::sync::Arc;
use std::time::Duration;
use crate::long_tier::{LongTierError, LongTierStore};
use crate::service::Scryer;
pub const DEFAULT_PROMOTE_INTERVAL: Duration = Duration::from_secs(3600);
#[derive(Debug, Clone)]
pub struct PromotionConfig {
pub interval: Duration,
pub retention_ms: u64,
}
impl PromotionConfig {
pub fn new(retention_ms: u64) -> Self {
Self { interval: DEFAULT_PROMOTE_INTERVAL, retention_ms }
}
pub fn with_interval(mut self, interval: Duration) -> Self {
self.interval = interval;
self
}
}
pub struct PromotionConsumer {
scryer: Arc<Scryer>,
long_tier: Arc<LongTierStore>,
cfg: PromotionConfig,
}
impl PromotionConsumer {
pub fn new(
scryer: Arc<Scryer>,
long_tier: Arc<LongTierStore>,
cfg: PromotionConfig,
) -> Self {
Self { scryer, long_tier, cfg }
}
pub fn run_once(&self) -> Result<usize, LongTierError> {
self.long_tier.rollover(self.scryer.store(), self.cfg.retention_ms)
}
pub fn spawn(self) -> tokio::task::JoinHandle<()> {
let consumer = Arc::new(self);
tokio::spawn(async move {
let mut ticker = tokio::time::interval(consumer.cfg.interval);
loop {
ticker.tick().await;
let c = Arc::clone(&consumer);
match tokio::task::spawn_blocking(move || c.run_once()).await {
Ok(Ok(0)) => {}
Ok(Ok(n)) => {
eprintln!("scryer promotion: {n} events promoted to long tier");
}
Ok(Err(e)) => eprintln!("scryer promotion: rollover error: {e}"),
Err(e) => eprintln!("scryer promotion: pass panicked: {e}"),
}
}
})
}
}
#[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;
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!("test::{seq}"),
msg: format!("msg {seq}"),
fields: json!({ "seq": seq }),
anchor: None,
source: EventSource::Synth,
}
}
fn setup(dir: &TempDir) -> (Arc<Scryer>, Arc<LongTierStore>, Arc<InMemoryObjectStore>) {
let cfg = ScryerConfig::new(dir.path().join("events.db"));
let retention_ms = 7 * MS_PER_DAY;
let scryer = Arc::new(Scryer::new(cfg, None).unwrap());
let obj_store = Arc::new(InMemoryObjectStore::new());
let lt_cfg = LongTierConfig { machine_id: "m1".to_string(), retention_ms };
let lt = Arc::new(LongTierStore::new(
lt_cfg,
Arc::clone(&obj_store) as Arc<dyn ObjectStore>,
));
(scryer, lt, obj_store)
}
#[test]
fn run_once_promotes_aged_events() {
let dir = TempDir::new().unwrap();
let (scryer, lt, obj_store) = setup(&dir);
let retention_ms = 7 * MS_PER_DAY;
let scope = EventScope::Service(MeshIdent("svc.prod".to_string()));
let run_id = TaskRunId::new();
let aged: Vec<(EventScope, Event)> = (0u32..4)
.map(|i| (scope.clone(), make_event(&run_id, i, MS_PER_DAY as u32, Level::Warn)))
.collect();
let recent = vec![(
scope.clone(),
make_event(&run_id, 99, (10 * MS_PER_DAY) as u32, Level::Info),
)];
scryer.store().insert_events(&aged).unwrap();
scryer.store().insert_events(&recent).unwrap();
assert_eq!(scryer.store().count().unwrap(), 5);
let consumer = PromotionConsumer::new(
Arc::clone(&scryer),
Arc::clone(<),
PromotionConfig::new(retention_ms),
);
let promoted = consumer.run_once().unwrap();
assert_eq!(promoted, 4, "4 aged events promoted");
assert_eq!(scryer.store().count().unwrap(), 1, "recent event remains");
assert!(
obj_store.contains_key("events/m1/1.parquet"),
"day-1 Parquet shard written to the object store"
);
}
#[test]
fn run_once_noop_when_nothing_aged() {
let dir = TempDir::new().unwrap();
let (scryer, lt, obj_store) = setup(&dir);
let retention_ms = 7 * MS_PER_DAY;
let scope = EventScope::Service(MeshIdent("svc.fresh".to_string()));
let run_id = TaskRunId::new();
let fresh = vec![(
scope.clone(),
make_event(&run_id, 0, (10 * MS_PER_DAY) as u32, Level::Info),
)];
scryer.store().insert_events(&fresh).unwrap();
let consumer = PromotionConsumer::new(scryer.clone(), lt, PromotionConfig::new(retention_ms));
assert_eq!(consumer.run_once().unwrap(), 0, "nothing past the cutoff");
assert_eq!(scryer.store().count().unwrap(), 1);
assert!(obj_store.keys().is_empty(), "no shards written");
}
#[tokio::test]
async fn spawned_loop_runs_a_pass() {
let dir = TempDir::new().unwrap();
let (scryer, lt, obj_store) = setup(&dir);
let retention_ms = 7 * MS_PER_DAY;
let scope = EventScope::Service(MeshIdent("svc.loop".to_string()));
let run_id = TaskRunId::new();
let aged: Vec<(EventScope, Event)> = (0u32..3)
.map(|i| (scope.clone(), make_event(&run_id, i, MS_PER_DAY as u32, Level::Error)))
.collect();
scryer.store().insert_events(&aged).unwrap();
let cfg = PromotionConfig::new(retention_ms).with_interval(Duration::from_millis(20));
let handle = PromotionConsumer::new(scryer.clone(), lt, cfg).spawn();
let mut promoted = false;
for _ in 0..50 {
if obj_store.contains_key("events/m1/1.parquet") {
promoted = true;
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
handle.abort();
assert!(promoted, "spawned loop promoted aged events to a shard");
assert_eq!(scryer.store().count().unwrap(), 0, "short-disk drained");
}
}