use std::sync::Arc;
use std::time::Duration;
use pensieve_core::catalog::Catalog;
use pensieve_core::segment_format::SegmentFormat;
use tokio::sync::broadcast;
use tracing::{debug, info, warn};
use crate::graph_handler::stored_provider;
pub struct GraphSnapshotScheduler {
catalog: Arc<dyn Catalog>,
format: Arc<dyn SegmentFormat>,
pub poll_interval: Duration,
}
fn needs_rebuild(last_marker: Option<&str>, current_version: &str) -> bool {
last_marker != Some(current_version)
}
impl GraphSnapshotScheduler {
fn default_poll_interval() -> Duration {
Duration::from_secs(
std::env::var("PENSIEVE_GRAPH_SNAPSHOT_POLL_SECS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(120),
)
}
pub fn new(catalog: Arc<dyn Catalog>, format: Arc<dyn SegmentFormat>) -> Self {
Self {
catalog,
format,
poll_interval: Self::default_poll_interval(),
}
}
pub async fn tick(&self) -> anyhow::Result<usize> {
let Some(store) = self.format.object_store() else {
debug!("no object store: graph snapshot scheduler idle");
return Ok(0);
};
let mut built = 0usize;
let databases = self
.catalog
.list_databases()
.await
.map_err(|e| anyhow::anyhow!("list_databases: {e}"))?;
for db in databases {
let graphs = match self.catalog.list_graphs(&db).await {
Ok(g) => g,
Err(e) => {
warn!(database = %db, error = %e, "list_graphs failed; skipping db");
continue;
}
};
if graphs.is_empty() {
continue;
}
let tables = match self.catalog.list_tables_in_database(&db).await {
Ok(t) => t,
Err(e) => {
warn!(database = %db, error = %e, "list_tables failed; skipping db");
continue;
}
};
for reg in graphs {
let Some(edge_tbl) = tables.iter().find(|t| t.name == reg.edge_table) else {
debug!(database = %db, edge_table = %reg.edge_table,
"edge table not found; skipping graph");
continue;
};
let current_version = edge_tbl.current_snapshot_id.to_string();
let marker = pensieve_graph::snapshot::meta_path(&db, ®.edge_table);
let last = pensieve_graph::snapshot::load_snapshot_meta(&store, &marker)
.await
.unwrap_or(None);
if !needs_rebuild(last.as_deref(), ¤t_version) {
continue; }
let provider = stored_provider(&self.catalog, &self.format, reg.clone());
match provider.build_and_store_snapshot(&store).await {
Ok(()) => {
if let Err(e) =
pensieve_graph::snapshot::store_snapshot_meta(&store, &marker, ¤t_version)
.await
{
warn!(database = %db, graph = %reg.edge_table, error = %e,
"snapshot built but marker write failed; will rebuild next tick");
}
built += 1;
debug!(database = %db, edge_table = %reg.edge_table,
version = %current_version, "graph snapshot rebuilt");
}
Err(e) => {
warn!(database = %db, edge_table = %reg.edge_table, error = %e,
"graph snapshot build failed; will retry next tick");
}
}
}
}
if built > 0 {
info!(snapshots = built, "graph snapshot scheduler rebuilt snapshots");
}
Ok(built)
}
pub async fn run(self, mut shutdown: broadcast::Receiver<()>) {
info!(
poll_secs = self.poll_interval.as_secs(),
"graph snapshot scheduler started"
);
loop {
tokio::select! {
_ = shutdown.recv() => {
info!("graph snapshot scheduler shutting down");
break;
}
() = tokio::time::sleep(self.poll_interval) => {
if let Err(e) = self.tick().await {
warn!(error = %e, "graph snapshot scheduler tick failed");
}
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::needs_rebuild;
#[test]
fn rebuild_decision_debounces_on_unchanged_version() {
assert!(needs_rebuild(None, "snap-1"));
assert!(!needs_rebuild(Some("snap-1"), "snap-1"));
assert!(needs_rebuild(Some("snap-1"), "snap-2"));
}
}