use crate::{Store, StoreError};
use rusqlite::params;
use serde::{Deserialize, Serialize};
fn fingerprint_row<T: serde::Serialize>(row: &T) -> String {
scc_core::fnv1a64_hex(serde_json::to_string(row).unwrap_or_default().as_bytes())
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ContextSnapshot {
pub id: String,
pub created_at: String,
pub task: String,
pub epoch: String,
pub revision: i64,
pub artifact: String,
pub artifact_hash: String,
pub entity_ids: Vec<String>,
pub entity_fp: std::collections::BTreeMap<String, String>,
pub rel_fp: std::collections::BTreeMap<String, String>,
pub contract_fp: std::collections::BTreeMap<String, String>,
pub state_fp: std::collections::BTreeMap<String, String>,
pub flow_names: Vec<String>,
pub budget: usize,
pub warnings: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SnapshotDiff {
pub still_valid: Vec<String>,
pub invalidated: Vec<String>,
pub modified_entities: Vec<String>,
pub changed_relationships: Vec<String>,
pub changed_contracts: Vec<String>,
pub changed_state: Vec<String>,
pub changed_flows: Vec<String>,
pub artifact_changed: bool,
pub snapshot_revision: i64,
pub current_revision: i64,
}
pub struct SnapshotSave<'a> {
pub task: &'a str,
pub epoch: &'a str,
pub revision: i64,
pub artifact: &'a str,
pub entity_ids: &'a [String],
pub budget: usize,
pub warnings: &'a [String],
}
impl Store {
pub fn save_snapshot(&self, save: SnapshotSave<'_>) -> Result<ContextSnapshot, StoreError> {
let hex =
scc_core::fnv1a64_hex(format!("{}\0{}", save.task, save.epoch).as_bytes());
let visible: std::collections::BTreeSet<&str> =
save.entity_ids.iter().map(|s| s.as_str()).collect();
let mut entity_fp = std::collections::BTreeMap::new();
let mut contract_fp = std::collections::BTreeMap::new();
let mut state_fp = std::collections::BTreeMap::new();
for e in self.all_entities()? {
let fp = fingerprint_row(&e);
if visible.contains(e.id.as_str()) {
entity_fp.insert(e.id.clone(), fp.clone());
}
match e.kind.as_str() {
k if k == scc_core::kinds::CONTRACT
|| k == scc_core::kinds::ROUTE
|| k == scc_core::kinds::TOPIC
|| k == scc_core::kinds::CONFIGURATION =>
{
contract_fp.insert(e.id.clone(), fp);
}
k if k == scc_core::kinds::DATA_STORE || k == scc_core::kinds::DATA_ENTITY => {
state_fp.insert(e.id.clone(), fp);
}
_ => {}
}
}
let mut rel_fp = std::collections::BTreeMap::new();
for r in self.all_relationships()? {
if visible.contains(r.subject.as_str()) || visible.contains(r.object.as_str()) {
rel_fp.insert(r.id.clone(), fingerprint_row(&r));
}
}
let mut flow_names: Vec<String> = self
.flow_graphs()
.unwrap_or_default()
.into_iter()
.map(|g| g.name)
.collect();
flow_names.sort();
let snap = ContextSnapshot {
id: format!("snap-{}-{hex}", save.revision),
created_at: scc_core::now_rfc3339(),
task: save.task.to_string(),
epoch: save.epoch.to_string(),
revision: save.revision,
artifact: save.artifact.to_string(),
artifact_hash: scc_core::fnv1a64_hex(save.artifact.as_bytes()),
entity_ids: save.entity_ids.to_vec(),
entity_fp,
rel_fp,
contract_fp,
state_fp,
flow_names,
budget: save.budget,
warnings: save.warnings.to_vec(),
};
self.conn.execute(
"INSERT INTO context_snapshots
(id, created_at, task, epoch, revision, artifact, artifact_hash,
entity_ids, entity_fp, rel_fp, contract_fp, state_fp, flow_names,
budget, warnings)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)
ON CONFLICT(id) DO UPDATE SET
created_at = excluded.created_at, artifact = excluded.artifact,
artifact_hash = excluded.artifact_hash,
entity_ids = excluded.entity_ids, entity_fp = excluded.entity_fp,
rel_fp = excluded.rel_fp, contract_fp = excluded.contract_fp,
state_fp = excluded.state_fp, flow_names = excluded.flow_names,
budget = excluded.budget, warnings = excluded.warnings",
params![
snap.id,
snap.created_at,
snap.task,
snap.epoch,
snap.revision,
snap.artifact,
snap.artifact_hash,
serde_json::to_string(&snap.entity_ids).unwrap_or_default(),
serde_json::to_string(&snap.entity_fp).unwrap_or_default(),
serde_json::to_string(&snap.rel_fp).unwrap_or_default(),
serde_json::to_string(&snap.contract_fp).unwrap_or_default(),
serde_json::to_string(&snap.state_fp).unwrap_or_default(),
serde_json::to_string(&snap.flow_names).unwrap_or_default(),
snap.budget as i64,
serde_json::to_string(&snap.warnings).unwrap_or_default(),
],
)?;
Ok(snap)
}
pub fn load_snapshot(&self, id: &str) -> Result<Option<ContextSnapshot>, StoreError> {
let mut stmt = self.conn.prepare(
"SELECT id, created_at, task, epoch, revision, artifact, artifact_hash,
entity_ids, entity_fp, rel_fp, contract_fp, state_fp, flow_names,
budget, warnings
FROM context_snapshots WHERE id = ?1",
)?;
let js = |r: &rusqlite::Row, i: usize| -> String { r.get(i).unwrap_or_default() };
let mut rows = stmt.query_map(params![id], |r| {
Ok(ContextSnapshot {
id: r.get(0)?,
created_at: r.get(1)?,
task: r.get(2)?,
epoch: r.get(3)?,
revision: r.get(4)?,
artifact: r.get(5)?,
artifact_hash: r.get(6)?,
entity_ids: serde_json::from_str(&js(r, 7)).unwrap_or_default(),
entity_fp: serde_json::from_str(&js(r, 8)).unwrap_or_default(),
rel_fp: serde_json::from_str(&js(r, 9)).unwrap_or_default(),
contract_fp: serde_json::from_str(&js(r, 10)).unwrap_or_default(),
state_fp: serde_json::from_str(&js(r, 11)).unwrap_or_default(),
flow_names: serde_json::from_str(&js(r, 12)).unwrap_or_default(),
budget: r.get::<_, i64>(13)? as usize,
warnings: serde_json::from_str(&js(r, 14)).unwrap_or_default(),
})
})?;
match rows.next() {
Some(r) => Ok(Some(r?)),
None => Ok(None),
}
}
pub fn diff_snapshot(&self, id: &str) -> Result<Option<SnapshotDiff>, StoreError> {
let Some(snap) = self.load_snapshot(id)? else {
return Ok(None);
};
let mut live_fp = std::collections::BTreeMap::new();
for e in self.all_entities()? {
live_fp.insert(e.id.clone(), fingerprint_row(&e));
}
let mut still_valid = Vec::new();
let mut invalidated = Vec::new();
let mut modified_entities = Vec::new();
for eid in &snap.entity_ids {
match live_fp.get(eid) {
None => invalidated.push(eid.clone()),
Some(fp) => {
if snap.entity_fp.get(eid).map(|s| s == fp).unwrap_or(true) {
still_valid.push(eid.clone());
} else {
modified_entities.push(eid.clone());
}
}
}
}
let visible: std::collections::BTreeSet<&str> =
snap.entity_ids.iter().map(|s| s.as_str()).collect();
let mut live_rel_fp = std::collections::BTreeMap::new();
for r in self.all_relationships()? {
if visible.contains(r.subject.as_str()) || visible.contains(r.object.as_str()) {
live_rel_fp.insert(r.id.clone(), fingerprint_row(&r));
}
}
let mut changed_relationships: Vec<String> = live_rel_fp
.iter()
.filter(|(k, v)| snap.rel_fp.get(*k) != Some(*v))
.map(|(k, _)| k.clone())
.chain(
snap.rel_fp
.keys()
.filter(|k| !live_rel_fp.contains_key(*k))
.cloned(),
)
.collect();
changed_relationships.sort();
let fp_diff = |old: &std::collections::BTreeMap<String, String>| -> Vec<String> {
let mut live = std::collections::BTreeMap::new();
for e in self.all_entities().unwrap_or_default() {
live.insert(e.id.clone(), fingerprint_row(&e));
}
let mut out: Vec<String> = live
.iter()
.filter(|(k, v)| old.get(*k) != Some(*v))
.map(|(k, _)| k.clone())
.chain(old.keys().filter(|k| !live.contains_key(*k)).cloned())
.collect();
out.sort();
out
};
let mut live_kind: std::collections::BTreeMap<String, String> = std::collections::BTreeMap::new();
for e in self.all_entities().unwrap_or_default() {
live_kind.insert(e.id.clone(), e.kind.clone());
}
let changed_contracts: Vec<String> = fp_diff(&snap.contract_fp)
.into_iter()
.filter(|id| {
snap.contract_fp.contains_key(id)
|| live_kind.get(id).map(|k| {
k == scc_core::kinds::CONTRACT
|| k == scc_core::kinds::ROUTE
|| k == scc_core::kinds::TOPIC
|| k == scc_core::kinds::CONFIGURATION
}).unwrap_or(false)
})
.collect();
let changed_state: Vec<String> = fp_diff(&snap.state_fp)
.into_iter()
.filter(|id| {
snap.state_fp.contains_key(id)
|| live_kind.get(id).map(|k| {
k == scc_core::kinds::DATA_STORE || k == scc_core::kinds::DATA_ENTITY
}).unwrap_or(false)
})
.collect();
let mut live_flows: Vec<String> = self
.flow_graphs()
.unwrap_or_default()
.into_iter()
.map(|g| g.name)
.collect();
live_flows.sort();
let live_flow_set: std::collections::BTreeSet<&str> =
live_flows.iter().map(|s| s.as_str()).collect();
let saved_flow_set: std::collections::BTreeSet<&str> =
snap.flow_names.iter().map(|s| s.as_str()).collect();
let mut changed_flows: Vec<String> = live_flow_set
.symmetric_difference(&saved_flow_set)
.map(|s| s.to_string())
.collect();
changed_flows.sort();
let current_revision: i64 = self
.conn
.query_row("SELECT COALESCE(MAX(rev), 0) FROM graph_revisions", [], |r| {
r.get(0)
})
.unwrap_or(0);
let artifact_changed = !snap.artifact.is_empty()
&& (!invalidated.is_empty()
|| !modified_entities.is_empty()
|| !changed_relationships.is_empty());
Ok(Some(SnapshotDiff {
still_valid,
invalidated,
modified_entities,
changed_relationships,
changed_contracts,
changed_state,
changed_flows,
artifact_changed,
snapshot_revision: snap.revision,
current_revision,
}))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{tests::tmp_store, Entity};
#[test]
fn snapshots_persist_diff_and_survive() {
let (s, _d) = tmp_store();
s.upsert_file("a.py", "h1", "python", "source", 10).unwrap();
for id in ["repo://t/symbol/a.py/f", "repo://t/symbol/a.py/g"] {
s.insert_entity(&Entity::new(id, "symbol", "f"), &["a.py".into()])
.unwrap();
}
let r1 = s.record_current_revision().unwrap();
let epoch = "epoch:test";
let snap = s
.save_snapshot(SnapshotSave {
task: "do the thing",
epoch,
revision: r1.rev,
artifact: "pack showing f and g",
entity_ids: &[
"repo://t/symbol/a.py/f".into(),
"repo://t/symbol/a.py/g".into(),
],
budget: 100,
warnings: &[],
})
.unwrap();
assert_eq!(snap.revision, r1.rev);
assert_eq!(snap.epoch, epoch);
let snap2 = s
.save_snapshot(SnapshotSave {
task: "do the thing",
epoch,
revision: r1.rev,
artifact: "pack showing f and g",
entity_ids: &[
"repo://t/symbol/a.py/f".into(),
"repo://t/symbol/a.py/g".into(),
],
budget: 100,
warnings: &[],
})
.unwrap();
assert_eq!(snap.id, snap2.id);
let loaded = s.load_snapshot(&snap.id).unwrap().unwrap();
assert_eq!(loaded.artifact, "pack showing f and g");
let d0 = s.diff_snapshot(&snap.id).unwrap().unwrap();
assert!(d0.invalidated.is_empty());
assert_eq!(d0.still_valid.len(), 2);
s.delete_entity("repo://t/symbol/a.py/g").unwrap();
s.upsert_file("a.py", "h2", "python", "source", 10).unwrap();
let _r2 = s.record_current_revision().unwrap();
let d1 = s.diff_snapshot(&snap.id).unwrap().unwrap();
assert_eq!(d1.still_valid, vec!["repo://t/symbol/a.py/f".to_string()]);
assert_eq!(d1.invalidated, vec!["repo://t/symbol/a.py/g".to_string()]);
assert_eq!(d1.snapshot_revision, r1.rev);
assert_eq!(d1.current_revision, r1.rev + 1);
let mut f2 = Entity::new("repo://t/symbol/a.py/f", "symbol", "f");
f2.attr("confidence", serde_json::json!(0.99));
s.insert_entity(&f2, &["a.py".into()]).unwrap();
let d2 = s.diff_snapshot(&snap.id).unwrap().unwrap();
assert!(d2.still_valid.is_empty(), "{d2:?}");
assert_eq!(d2.modified_entities, vec!["repo://t/symbol/a.py/f".to_string()]);
assert!(d2.artifact_changed, "visible rows changed: re-render would differ");
assert!(s.diff_snapshot("snap-0-deadbeef").unwrap().is_none());
}
}