Skip to main content

graphforge_storage/
generation.rs

1//! Project mutation generations — staleness signals for derived indexes.
2//!
3//! Every committed batch that rewrites topology (`topology/nodes.parquet` or
4//! any file under `topology/edges/`) bumps a monotonically increasing counter
5//! persisted at `topology/generation.json`. Derived indexes record the counter
6//! they were built from (see the [`adjacency`](crate::adjacency) manifest); a
7//! mismatch against the current value marks the index **stale**, and a stale
8//! or missing index falls back to scan-and-build — identical results, only
9//! slower. Property-only writes (`properties/`, `edge_properties/`) never bump
10//! the topology counter because properties cannot change adjacency.
11//!
12//! Search artifacts use the sibling `search_generation`. It advances for node
13//! topology and node-property commits, but not for edge-only, provenance, or
14//! knowledge-layer writes. Older projects without this key inherit the current
15//! topology generation until their first search-relevant mutation.
16//!
17//! # Crash-safety invariant
18//!
19//! [`commit_topology_aware`] bumps the counter **strictly before** the first
20//! rename of the staged batch. A crash after the bump but before (or during)
21//! the commit leaves the counter advanced over an unchanged or partially
22//! renamed topology — any existing index now merely *looks* stale and is
23//! rebuilt, costing one spurious rebuild. The reverse order would be unsound:
24//! a crash between commit and bump would leave new topology under the old
25//! counter, making a stale index look **fresh** and silently serving wrong
26//! traversals. Spurious bumps are safe; missed bumps are not.
27//!
28//! Multi-process writers can lose a bump (read-increment-rename is not
29//! cross-process atomic); this matches the consistency envelope of every
30//! Parquet rewrite in this embedded engine (see [`crate::staging`]).
31
32use std::io::Write;
33use std::path::{Path, PathBuf};
34
35use graphforge_core::GfError;
36
37use crate::staging::RewriteBatch;
38
39/// JSON key holding the counter inside `topology/generation.json`.
40const GENERATION_KEY: &str = "topology_generation";
41/// JSON key holding the graph-native search source counter.
42const SEARCH_GENERATION_KEY: &str = "search_generation";
43
44#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
45struct GenerationState {
46    topology: u64,
47    search: u64,
48}
49
50fn storage_err(e: impl std::fmt::Display) -> GfError {
51    GfError::Storage(e.to_string())
52}
53
54/// Path of the generation counter file within `project_dir`:
55/// `topology/generation.json`.
56#[must_use]
57pub fn generation_path(project_dir: &Path) -> PathBuf {
58    project_dir.join("topology").join("generation.json")
59}
60
61/// The project's current topology generation.
62///
63/// A missing file is generation **0** (a project that has never written
64/// topology with a counter-aware binary), mirroring the absent-file semantics
65/// of the catalog readers. This is sound because no index manifest can
66/// predate the counter: every binary that writes `indexes/` also bumps.
67///
68/// # Errors
69/// Returns [`GfError::Storage`] if the file exists but cannot be read or is
70/// not of the form `{"topology_generation": N}` — callers treating the index
71/// as a capability must then consider it always-stale, never fresh.
72pub fn read_topology_generation(project_dir: &Path) -> Result<u64, GfError> {
73    Ok(read_generation_state(project_dir)?.topology)
74}
75
76/// The generation of the committed node topology and node properties consumed
77/// by graph-native search.
78///
79/// A legacy counter without `search_generation` inherits
80/// `topology_generation`, preserving reopen compatibility. A missing counter
81/// file is generation zero.
82///
83/// # Errors
84/// Returns [`GfError::Storage`] when the counter exists but is corrupt or
85/// unreadable.
86pub fn read_search_generation(project_dir: &Path) -> Result<u64, GfError> {
87    Ok(read_generation_state(project_dir)?.search)
88}
89
90fn read_generation_state(project_dir: &Path) -> Result<GenerationState, GfError> {
91    let path = generation_path(project_dir);
92    let contents = match std::fs::read_to_string(&path) {
93        Ok(c) => c,
94        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
95            return Ok(GenerationState::default());
96        }
97        Err(e) => {
98            return Err(GfError::Storage(format!(
99                "cannot read {}: {e}",
100                path.display()
101            )));
102        }
103    };
104    let value: serde_json::Value = serde_json::from_str(&contents)
105        .map_err(|e| GfError::Storage(format!("corrupt {}: {e}", path.display())))?;
106    let topology = value
107        .get(GENERATION_KEY)
108        .and_then(serde_json::Value::as_u64)
109        .ok_or_else(|| {
110            GfError::Storage(format!(
111                "corrupt {}: expected {{\"{GENERATION_KEY}\": <u64>}}",
112                path.display()
113            ))
114        })?;
115    let search = match value.get(SEARCH_GENERATION_KEY) {
116        Some(value) => value.as_u64().ok_or_else(|| {
117            GfError::Storage(format!(
118                "corrupt {}: expected \"{SEARCH_GENERATION_KEY}\" to be a u64",
119                path.display()
120            ))
121        })?,
122        None => topology,
123    };
124    Ok(GenerationState { topology, search })
125}
126
127/// Atomically persist `current + 1` (sibling temp + rename) and return the
128/// new value. Creates `topology/` if needed.
129///
130/// # Errors
131/// Returns [`GfError::Storage`] if the current value cannot be read (corrupt
132/// file) or on I/O failure; on failure the prior file is untouched.
133pub fn bump_topology_generation(project_dir: &Path) -> Result<u64, GfError> {
134    Ok(bump_generations(project_dir, true, false)?.topology)
135}
136
137/// Atomically advance and persist the graph-native search generation.
138///
139/// # Errors
140/// Returns [`GfError::Storage`] if the existing generation is corrupt or the
141/// replacement cannot be persisted.
142pub fn bump_search_generation(project_dir: &Path) -> Result<u64, GfError> {
143    Ok(bump_generations(project_dir, false, true)?.search)
144}
145
146fn bump_generations(
147    project_dir: &Path,
148    bump_topology: bool,
149    bump_search: bool,
150) -> Result<GenerationState, GfError> {
151    let mut next = read_generation_state(project_dir)?;
152    if bump_topology {
153        next.topology = next
154            .topology
155            .checked_add(1)
156            .ok_or_else(|| GfError::Storage("topology generation counter overflow".to_owned()))?;
157    }
158    if bump_search {
159        next.search = next
160            .search
161            .checked_add(1)
162            .ok_or_else(|| GfError::Storage("search generation counter overflow".to_owned()))?;
163    }
164    let path = generation_path(project_dir);
165    let parent = path.parent().expect("generation path always has a parent");
166    std::fs::create_dir_all(parent).map_err(storage_err)?;
167    let mut tmp = tempfile::Builder::new()
168        .prefix("generation.json.")
169        .suffix(".tmp")
170        .tempfile_in(parent)
171        .map_err(storage_err)?;
172    let body = serde_json::json!({
173        GENERATION_KEY: next.topology,
174        SEARCH_GENERATION_KEY: next.search,
175    })
176    .to_string();
177    tmp.write_all(body.as_bytes()).map_err(storage_err)?;
178    tmp.as_file().sync_all().map_err(storage_err)?;
179    tmp.persist(&path).map_err(|e| storage_err(e.error))?;
180    sync_directory(parent)?;
181    Ok(next)
182}
183
184#[cfg(unix)]
185fn sync_directory(path: &Path) -> Result<(), GfError> {
186    std::fs::File::open(path)
187        .and_then(|directory| directory.sync_all())
188        .map_err(storage_err)
189}
190
191#[cfg(not(unix))]
192fn sync_directory(_path: &Path) -> Result<(), GfError> {
193    Ok(())
194}
195
196/// Whether any staged destination in `staged` rewrites topology:
197/// `topology/nodes.parquet` or any file under `topology/edges/` (including
198/// `_exploratory.parquet`). Paths elsewhere (`properties/`,
199/// `edge_properties/`, `provenance/`, …) do not count.
200#[must_use]
201pub fn touches_topology(staged: &RewriteBatch, project_dir: &Path) -> bool {
202    let topology = project_dir.join("topology");
203    let nodes = topology.join("nodes.parquet");
204    let edges = topology.join("edges");
205    staged
206        .staged_paths()
207        .any(|path| path == nodes || path.starts_with(&edges))
208}
209
210/// Whether a staged batch changes graph-native search inputs: node identity or
211/// label membership (`topology/nodes.parquet`) or node properties
212/// (`properties/`). Edge-only and knowledge-layer writes are intentionally
213/// excluded.
214#[must_use]
215pub fn touches_search_source(staged: &RewriteBatch, project_dir: &Path) -> bool {
216    let nodes = project_dir.join("topology").join("nodes.parquet");
217    let properties = project_dir.join("properties");
218    staged
219        .staged_paths()
220        .any(|path| path == nodes || path.starts_with(&properties))
221}
222
223/// Commit `staged`, bumping each affected generation **first** (see the module
224/// docs for why bump-before-commit is the only sound order). Edge topology
225/// advances only topology; node topology advances topology and search; node
226/// properties advance only search.
227///
228/// Returns `Some(new_generation)` when the batch bumped, `None` otherwise — the
229/// caller tags an adjacency delta segment (#765) with the returned value rather
230/// than re-reading the counter (which a concurrent bump could have advanced).
231///
232/// # Errors
233/// Returns [`GfError::Storage`] on bump or rename failure. A bump followed by
234/// a failed commit leaves the counter advanced — safe (the index reads as
235/// stale), see the crash-safety invariant.
236pub fn commit_topology_aware(
237    staged: RewriteBatch,
238    project_dir: &Path,
239) -> Result<Option<u64>, GfError> {
240    let topology = touches_topology(&staged, project_dir);
241    let search = touches_search_source(&staged, project_dir);
242    let bumped = if topology || search {
243        let generations = bump_generations(project_dir, topology, search)?;
244        topology.then_some(generations.topology)
245    } else {
246        None
247    };
248    staged.commit()?;
249    Ok(bumped)
250}
251
252// ---------------------------------------------------------------------------
253// Tests
254// ---------------------------------------------------------------------------
255
256#[cfg(test)]
257mod tests {
258    use std::sync::Arc;
259
260    use arrow::array::{Int64Array, RecordBatch};
261    use arrow::datatypes::{DataType, Field, Schema, SchemaRef};
262    use tempfile::TempDir;
263
264    use super::*;
265
266    fn int_batch() -> (SchemaRef, RecordBatch) {
267        let schema = Arc::new(Schema::new(vec![Field::new("v", DataType::Int64, false)]));
268        let batch = RecordBatch::try_new(
269            Arc::clone(&schema),
270            vec![Arc::new(Int64Array::from(vec![1]))],
271        )
272        .unwrap();
273        (schema, batch)
274    }
275
276    fn staged_for(dir: &Path, rel_paths: &[&str]) -> RewriteBatch {
277        let mut staged = RewriteBatch::new();
278        for rel in rel_paths {
279            let (schema, batch) = int_batch();
280            staged.stage(&dir.join(rel), schema, &batch).unwrap();
281        }
282        staged
283    }
284
285    #[test]
286    fn missing_file_reads_as_generation_zero() {
287        let dir = TempDir::new().unwrap();
288        assert_eq!(read_topology_generation(dir.path()).unwrap(), 0);
289        assert_eq!(read_search_generation(dir.path()).unwrap(), 0);
290    }
291
292    #[test]
293    fn bump_increments_and_persists() {
294        let dir = TempDir::new().unwrap();
295        assert_eq!(bump_topology_generation(dir.path()).unwrap(), 1);
296        assert_eq!(bump_topology_generation(dir.path()).unwrap(), 2);
297        assert_eq!(bump_topology_generation(dir.path()).unwrap(), 3);
298        assert_eq!(read_topology_generation(dir.path()).unwrap(), 3);
299        assert_eq!(read_search_generation(dir.path()).unwrap(), 0);
300        // No temp residue next to the counter.
301        let temps = std::fs::read_dir(dir.path().join("topology"))
302            .unwrap()
303            .filter_map(Result::ok)
304            .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
305            .count();
306        assert_eq!(temps, 0);
307    }
308
309    #[test]
310    fn corrupt_file_is_an_error_not_zero() {
311        let dir = TempDir::new().unwrap();
312        let path = generation_path(dir.path());
313        std::fs::create_dir_all(path.parent().unwrap()).unwrap();
314
315        for bad in ["not json", "{}", "{\"topology_generation\": -1}", "[3]"] {
316            std::fs::write(&path, bad).unwrap();
317            assert!(
318                matches!(
319                    read_topology_generation(dir.path()),
320                    Err(GfError::Storage(_))
321                ),
322                "{bad:?} must not parse"
323            );
324            // A corrupt counter must also fail the bump, not silently reset.
325            assert!(matches!(
326                bump_topology_generation(dir.path()),
327                Err(GfError::Storage(_))
328            ));
329        }
330    }
331
332    #[test]
333    fn touches_topology_matrix() {
334        let dir = TempDir::new().unwrap();
335        for (rel, topology, search) in [
336            ("topology/nodes.parquet", true, true),
337            ("topology/edges/KNOWS.parquet", true, false),
338            ("topology/edges/_exploratory.parquet", true, false),
339            ("topology/runtime_catalog.parquet", false, false),
340            ("properties/Person.parquet", false, true),
341            ("edge_properties/KNOWS.parquet", false, false),
342            ("auxiliary/records.parquet", false, false),
343        ] {
344            let staged = staged_for(dir.path(), &[rel]);
345            assert_eq!(
346                touches_topology(&staged, dir.path()),
347                topology,
348                "{rel} should {}count as topology",
349                if topology { "" } else { "not " }
350            );
351            assert_eq!(
352                touches_search_source(&staged, dir.path()),
353                search,
354                "{rel} should {}count as a search source",
355                if search { "" } else { "not " }
356            );
357        }
358    }
359
360    #[test]
361    fn commit_topology_aware_bumps_only_for_topology() {
362        let dir = TempDir::new().unwrap();
363
364        // Property-only batch: commit, no bump.
365        let staged = staged_for(dir.path(), &["properties/Person.parquet"]);
366        commit_topology_aware(staged, dir.path()).unwrap();
367        assert_eq!(read_topology_generation(dir.path()).unwrap(), 0);
368        assert_eq!(read_search_generation(dir.path()).unwrap(), 1);
369
370        // Mixed batch staging topology: exactly one bump.
371        let staged = staged_for(
372            dir.path(),
373            &["topology/nodes.parquet", "properties/Person.parquet"],
374        );
375        commit_topology_aware(staged, dir.path()).unwrap();
376        assert_eq!(read_topology_generation(dir.path()).unwrap(), 1);
377        assert_eq!(read_search_generation(dir.path()).unwrap(), 2);
378        assert!(dir.path().join("topology/nodes.parquet").exists());
379
380        // Edge-only topology advances adjacency without invalidating search.
381        let staged = staged_for(dir.path(), &["topology/edges/KNOWS.parquet"]);
382        commit_topology_aware(staged, dir.path()).unwrap();
383        assert_eq!(read_topology_generation(dir.path()).unwrap(), 2);
384        assert_eq!(read_search_generation(dir.path()).unwrap(), 2);
385    }
386
387    #[test]
388    fn legacy_counter_seeds_search_generation() {
389        let dir = TempDir::new().unwrap();
390        let path = generation_path(dir.path());
391        std::fs::create_dir_all(path.parent().unwrap()).unwrap();
392        std::fs::write(&path, r#"{"topology_generation":7}"#).unwrap();
393
394        assert_eq!(read_search_generation(dir.path()).unwrap(), 7);
395        assert_eq!(bump_search_generation(dir.path()).unwrap(), 8);
396        assert_eq!(read_topology_generation(dir.path()).unwrap(), 7);
397    }
398}