1use std::io::Write;
33use std::path::{Path, PathBuf};
34
35use graphforge_core::GfError;
36
37use crate::staging::RewriteBatch;
38
39const GENERATION_KEY: &str = "topology_generation";
41const 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#[must_use]
57pub fn generation_path(project_dir: &Path) -> PathBuf {
58 project_dir.join("topology").join("generation.json")
59}
60
61pub fn read_topology_generation(project_dir: &Path) -> Result<u64, GfError> {
73 Ok(read_generation_state(project_dir)?.topology)
74}
75
76pub 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
127pub fn bump_topology_generation(project_dir: &Path) -> Result<u64, GfError> {
134 Ok(bump_generations(project_dir, true, false)?.topology)
135}
136
137pub 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#[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#[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
223pub 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#[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 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 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 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 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 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}