use crate::store::{EventStore, ScryerStoreError};
use arrow_array::{
BooleanArray, Int64Array, StringArray, UInt32Array,
cast::AsArray,
};
use arrow_schema::{DataType, Field, Schema};
use bytes::Bytes;
use observation::{ChunkRef, Event, EventScope, EventSource, ForgeId, Level, TaskRunId};
use parquet::arrow::ArrowWriter;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use serde_json::Value;
use std::collections::BTreeMap;
use std::str::FromStr;
use std::sync::Arc;
use thiserror::Error;
use workload_spec::MeshIdent;
pub use yah_object_store::{Error as ObjectStoreError, InMemoryObjectStore, ObjectStore};
pub const MS_PER_DAY: u64 = 24 * 3_600 * 1_000;
#[derive(Debug, Error)]
pub enum LongTierError {
#[error("parquet: {0}")]
Parquet(#[from] parquet::errors::ParquetError),
#[error("arrow: {0}")]
Arrow(#[from] arrow_schema::ArrowError),
#[error("store: {0}")]
Store(#[from] ScryerStoreError),
#[error("object store: {0}")]
ObjectStore(#[from] ObjectStoreError),
#[error("not found: {0}")]
NotFound(String),
#[error("decode: {0}")]
Decode(String),
}
#[derive(Debug, Clone)]
pub struct LongTierConfig {
pub machine_id: String,
pub retention_ms: u64,
}
impl LongTierConfig {
pub fn new(machine_id: impl Into<String>) -> Self {
Self { machine_id: machine_id.into(), retention_ms: 7 * MS_PER_DAY }
}
}
pub struct LongTierStore {
cfg: LongTierConfig,
object_store: Arc<dyn ObjectStore>,
}
impl LongTierStore {
pub fn new(cfg: LongTierConfig, object_store: Arc<dyn ObjectStore>) -> Self {
Self { cfg, object_store }
}
fn shard_key(&self, day: u64) -> String {
format!("events/{}/{}.parquet", self.cfg.machine_id, day)
}
pub fn rollover(
&self,
event_store: &EventStore,
cutoff_ms: u64,
) -> Result<usize, LongTierError> {
let old_items = event_store.query_events_older_than(cutoff_ms)?;
if old_items.is_empty() {
return Ok(0);
}
let mut by_day: BTreeMap<u64, Vec<(EventScope, Event)>> = BTreeMap::new();
for item in &old_items {
let day = item.1.offset_ms as u64 / MS_PER_DAY;
by_day.entry(day).or_default().push(item.clone());
}
for (day, new_items) in &by_day {
let key = self.shard_key(*day);
let to_write = if let Some(existing_bytes) =
self.object_store.get(&key)?
{
let mut merged = read_shard_bytes(existing_bytes)?;
merged.extend_from_slice(new_items);
merged.sort_by(|(sa, ea), (sb, eb)| {
sa.kind_str()
.cmp(sb.kind_str())
.then(sa.id_str().cmp(&sb.id_str()))
.then(ea.seq.cmp(&eb.seq))
});
merged.dedup_by(|(sa, ea), (sb, eb)| {
sa.kind_str() == sb.kind_str()
&& sa.id_str() == sb.id_str()
&& ea.seq == eb.seq
});
merged
} else {
new_items.clone()
};
let shard_bytes = write_shard_bytes(&to_write)?;
self.object_store.put(&key, shard_bytes)?;
}
let pruned = event_store.prune_older_than(cutoff_ms)?;
Ok(pruned)
}
pub fn query_range(
&self,
scope: Option<&EventScope>,
since_ms: u64,
until_ms: u64,
) -> Result<Vec<(EventScope, Event)>, LongTierError> {
let since_day = since_ms / MS_PER_DAY;
let until_day = until_ms.saturating_sub(1) / MS_PER_DAY;
let mut results = Vec::new();
for day in since_day..=until_day {
let key = self.shard_key(day);
let data = match self.object_store.get(&key)? {
Some(d) => d,
None => continue,
};
for (ev_scope, ev) in read_shard_bytes(data)? {
let off = ev.offset_ms as u64;
let in_range = off >= since_ms && off < until_ms;
let scope_match = scope.map_or(true, |s| &ev_scope == s);
if in_range && scope_match {
results.push((ev_scope, ev));
}
}
}
Ok(results)
}
}
fn events_schema() -> Arc<Schema> {
Arc::new(Schema::new(vec![
Field::new("scope_kind", DataType::Utf8, false),
Field::new("scope_id", DataType::Utf8, false),
Field::new("seq", DataType::UInt32, false),
Field::new("offset_ms", DataType::UInt32, false),
Field::new("level", DataType::Utf8, false),
Field::new("target", DataType::Utf8, false),
Field::new("msg", DataType::Utf8, false),
Field::new("fields_json", DataType::Utf8, false),
Field::new("has_anchor", DataType::Boolean, false),
Field::new("anchor_seq", DataType::Int64, false), Field::new("source_kind", DataType::Utf8, false),
Field::new("source_name", DataType::Utf8, false),
]))
}
fn write_shard_bytes(items: &[(EventScope, Event)]) -> Result<Vec<u8>, LongTierError> {
if items.is_empty() {
let schema = events_schema();
let mut buf = Vec::new();
let writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
writer.close()?;
return Ok(buf);
}
let schema = events_schema();
let scope_kind_col: StringArray =
items.iter().map(|(s, _)| s.kind_str()).collect::<Vec<&str>>().into();
let scope_id_vals: Vec<String> = items.iter().map(|(s, _)| s.id_str()).collect();
let scope_id_col: StringArray =
scope_id_vals.iter().map(|s| s.as_str()).collect::<Vec<&str>>().into();
let seq_col: UInt32Array = items.iter().map(|(_, e)| e.seq).collect();
let offset_ms_col: UInt32Array = items.iter().map(|(_, e)| e.offset_ms).collect();
let level_col: StringArray =
items.iter().map(|(_, e)| e.level.as_str()).collect::<Vec<&str>>().into();
let target_vals: Vec<&str> = items.iter().map(|(_, e)| e.target.as_str()).collect();
let target_col: StringArray = target_vals.into();
let msg_vals: Vec<&str> = items.iter().map(|(_, e)| e.msg.as_str()).collect();
let msg_col: StringArray = msg_vals.into();
let fields_vals: Vec<String> =
items.iter().map(|(_, e)| e.fields.to_string()).collect();
let fields_col: StringArray =
fields_vals.iter().map(|s| s.as_str()).collect::<Vec<&str>>().into();
let has_anchor_col: BooleanArray = items.iter().map(|(_, e)| e.anchor.is_some()).collect();
let anchor_seq_col: Int64Array = items
.iter()
.map(|(_, e)| e.anchor.as_ref().map(|a| a.seq as i64).unwrap_or(0))
.collect();
let source_kind_col: StringArray =
items.iter().map(|(_, e)| e.source.kind_str()).collect::<Vec<&str>>().into();
let source_name_vals: Vec<&str> =
items.iter().map(|(_, e)| e.source.name_str()).collect();
let source_name_col: StringArray = source_name_vals.into();
let batch = arrow_array::RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(scope_kind_col),
Arc::new(scope_id_col),
Arc::new(seq_col),
Arc::new(offset_ms_col),
Arc::new(level_col),
Arc::new(target_col),
Arc::new(msg_col),
Arc::new(fields_col),
Arc::new(has_anchor_col),
Arc::new(anchor_seq_col),
Arc::new(source_kind_col),
Arc::new(source_name_col),
],
)?;
let mut buf = Vec::new();
let mut writer = ArrowWriter::try_new(&mut buf, schema, None)?;
writer.write(&batch)?;
writer.close()?;
Ok(buf)
}
fn read_shard_bytes(data: Vec<u8>) -> Result<Vec<(EventScope, Event)>, LongTierError> {
let bytes = Bytes::from(data);
let reader = ParquetRecordBatchReaderBuilder::try_new(bytes)
.map_err(LongTierError::Parquet)?
.with_batch_size(4096)
.build()
.map_err(LongTierError::Parquet)?;
let mut results = Vec::new();
for batch_result in reader {
let batch = batch_result.map_err(LongTierError::Arrow)?;
let n = batch.num_rows();
if n == 0 {
continue;
}
let scope_kind_col = batch.column(0).as_string::<i32>();
let scope_id_col = batch.column(1).as_string::<i32>();
let seq_col = batch.column(2).as_primitive::<arrow_array::types::UInt32Type>();
let offset_ms_col = batch.column(3).as_primitive::<arrow_array::types::UInt32Type>();
let level_col = batch.column(4).as_string::<i32>();
let target_col = batch.column(5).as_string::<i32>();
let msg_col = batch.column(6).as_string::<i32>();
let fields_col = batch.column(7).as_string::<i32>();
let has_anchor_col = batch.column(8).as_boolean();
let anchor_seq_col = batch.column(9).as_primitive::<arrow_array::types::Int64Type>();
let source_kind_col = batch.column(10).as_string::<i32>();
let source_name_col = batch.column(11).as_string::<i32>();
for i in 0..n {
let scope_kind = scope_kind_col.value(i);
let scope_id = scope_id_col.value(i);
let scope = decode_scope(scope_kind, scope_id)?;
let seq = seq_col.value(i);
let offset_ms = offset_ms_col.value(i);
let level_str = level_col.value(i);
let level =
Level::from_str(level_str).map_err(|e| LongTierError::Decode(e))?;
let target = target_col.value(i).to_string();
let msg = msg_col.value(i).to_string();
let fields_str = fields_col.value(i);
let fields: Value = serde_json::from_str(fields_str)
.unwrap_or(serde_json::json!({}));
let has_anchor = has_anchor_col.value(i);
let anchor = if has_anchor {
Some(ChunkRef { seq: anchor_seq_col.value(i) as u32 })
} else {
None
};
let source_kind_str = source_kind_col.value(i);
let source_name_str = source_name_col.value(i).to_string();
let source = decode_source(source_kind_str, source_name_str);
let run_id = match &scope {
EventScope::TaskRun(id) => id.clone(),
EventScope::Forge(id) => TaskRunId(id.0),
EventScope::Service(_) => scope_id
.parse::<TaskRunId>()
.unwrap_or_else(|_| TaskRunId::new()),
};
results.push((
scope,
Event { run_id, seq, offset_ms, level, target, msg, fields, anchor, source },
));
}
}
Ok(results)
}
fn decode_scope(kind: &str, id: &str) -> Result<EventScope, LongTierError> {
match kind {
"task_run" => id
.parse::<TaskRunId>()
.map(EventScope::TaskRun)
.map_err(|e| LongTierError::Decode(e.to_string())),
"service" => Ok(EventScope::Service(MeshIdent(id.to_string()))),
"forge" => id
.parse::<ForgeId>()
.map(EventScope::Forge)
.map_err(|e| LongTierError::Decode(e.to_string())),
other => Err(LongTierError::Decode(format!("unknown scope kind: {other}"))),
}
}
fn decode_source(kind: &str, name: String) -> EventSource {
match kind {
"beholder" => EventSource::Beholder { name, version: String::new() },
"shim" => EventSource::Shim { lib: name, version: String::new() },
_ => EventSource::Synth,
}
}
#[cfg(test)]
mod tests {
use super::*;
use observation::{EventSource, Level, TaskRunId};
use serde_json::json;
use tempfile::TempDir;
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 make_store(dir: &TempDir) -> EventStore {
EventStore::open(&dir.path().join("events.db")).unwrap()
}
fn make_lt(machine_id: &str) -> (LongTierStore, Arc<InMemoryObjectStore>) {
let obj_store = Arc::new(InMemoryObjectStore::new());
let cfg =
LongTierConfig { machine_id: machine_id.to_string(), retention_ms: 7 * MS_PER_DAY };
let lt = LongTierStore::new(cfg, Arc::clone(&obj_store) as Arc<dyn ObjectStore>);
(lt, obj_store)
}
#[test]
fn rollover() {
let dir = TempDir::new().unwrap();
let store = make_store(&dir);
let (lt, obj_store) = make_lt("m1");
let scope_a = EventScope::Service(MeshIdent("svc-a.prod".to_string()));
let run_id_a = TaskRunId::new();
let old_cutoff_ms = 7 * MS_PER_DAY; let day0_offset = MS_PER_DAY as u32;
let old_events: Vec<(EventScope, Event)> = (0u32..5)
.map(|i| (scope_a.clone(), make_event(&run_id_a, i, day0_offset, Level::Warn)))
.collect();
store.insert_events(&old_events).unwrap();
assert_eq!(store.count().unwrap(), 5, "5 events in short-disk before rollover");
let recent_events = vec![(
scope_a.clone(),
make_event(&run_id_a, 99, (10 * MS_PER_DAY) as u32, Level::Info),
)];
store.insert_events(&recent_events).unwrap();
assert_eq!(store.count().unwrap(), 6);
let promoted = lt.rollover(&store, old_cutoff_ms).unwrap();
assert_eq!(promoted, 5, "5 old events should have been promoted");
assert_eq!(store.count().unwrap(), 1, "1 recent event remains in short-disk");
let shard_key = format!("events/m1/1.parquet");
assert!(
obj_store.contains_key(&shard_key),
"Parquet shard events/m1/1.parquet should exist"
);
let shard_data = obj_store.get(&shard_key).unwrap().unwrap();
let shard_events = read_shard_bytes(shard_data).unwrap();
assert_eq!(shard_events.len(), 5, "shard should hold all 5 old events");
assert!(
shard_events.iter().all(|(_, e)| e.level == Level::Warn),
"all shard events should be Warn"
);
let recent_shard = format!("events/m1/10.parquet");
assert!(
!obj_store.contains_key(&recent_shard),
"recent event shard should not exist"
);
}
#[test]
fn rollover_idempotent() {
let dir = TempDir::new().unwrap();
let store = make_store(&dir);
let (lt, obj_store) = make_lt("m2");
let scope = EventScope::Service(MeshIdent("svc.local".to_string()));
let run_id = TaskRunId::new();
let items: Vec<(EventScope, Event)> =
(0u32..3).map(|i| (scope.clone(), make_event(&run_id, i, 1000, Level::Info))).collect();
store.insert_events(&items).unwrap();
lt.rollover(&store, 7 * MS_PER_DAY).unwrap();
store.insert_events(&items).unwrap();
lt.rollover(&store, 7 * MS_PER_DAY).unwrap();
let key = "events/m2/0.parquet".to_string();
let data = obj_store.get(&key).unwrap().unwrap();
let shard = read_shard_bytes(data).unwrap();
assert_eq!(shard.len(), 3, "merge must deduplicate — still 3 events");
}
#[test]
fn parquet_round_trip() {
let run_id = TaskRunId::new();
let scope = EventScope::Service(MeshIdent("svc.rt".to_string()));
let events: Vec<(EventScope, Event)> = vec![
(scope.clone(), make_event(&run_id, 0, 100, Level::Info)),
(
scope.clone(),
Event {
run_id: run_id.clone(),
seq: 1,
offset_ms: 200,
level: Level::Error,
target: "app::db".to_string(),
msg: "connection refused".to_string(),
fields: json!({"error": {"code": "ECONNREFUSED"}}),
anchor: Some(ChunkRef { seq: 42 }),
source: EventSource::Beholder {
name: "pino".to_string(),
version: "1.0".to_string(),
},
},
),
];
let bytes = write_shard_bytes(&events).unwrap();
let recovered = read_shard_bytes(bytes).unwrap();
assert_eq!(recovered.len(), 2);
assert_eq!(recovered[0].1.seq, 0);
assert_eq!(recovered[0].1.level, Level::Info);
assert_eq!(recovered[1].1.seq, 1);
assert_eq!(recovered[1].1.level, Level::Error);
assert_eq!(recovered[1].1.target, "app::db");
assert_eq!(recovered[1].1.anchor.as_ref().unwrap().seq, 42);
assert_eq!(recovered[1].1.source.kind_str(), "beholder");
}
#[test]
fn query_range_filters_by_range() {
let dir = TempDir::new().unwrap();
let store = make_store(&dir);
let (lt, _) = make_lt("m3");
let scope = EventScope::Service(MeshIdent("svc.range".to_string()));
let run_id = TaskRunId::new();
let day0: Vec<(EventScope, Event)> =
(0u32..3).map(|i| (scope.clone(), make_event(&run_id, i, 100, Level::Info))).collect();
let day2_ms = (2 * MS_PER_DAY + 1) as u32;
let day2: Vec<(EventScope, Event)> = (3u32..5)
.map(|i| (scope.clone(), make_event(&run_id, i, day2_ms, Level::Warn)))
.collect();
store.insert_events(&day0).unwrap();
store.insert_events(&day2).unwrap();
lt.rollover(&store, 7 * MS_PER_DAY).unwrap();
let results = lt.query_range(Some(&scope), 0, MS_PER_DAY).unwrap();
assert_eq!(results.len(), 3, "only day-0 events in range");
assert!(results.iter().all(|(_, e)| e.level == Level::Info));
let results = lt.query_range(Some(&scope), 2 * MS_PER_DAY, 3 * MS_PER_DAY).unwrap();
assert_eq!(results.len(), 2, "only day-2 events in range");
assert!(results.iter().all(|(_, e)| e.level == Level::Warn));
}
}