Skip to main content

graphforge_storage/
mutator.rs

1//! In-place mutation primitives for `DELETE` / `DETACH DELETE` (#740).
2//!
3//! Unlike [`GraphWriter`](crate::GraphWriter), which buffers and **appends**,
4//! these functions **rewrite** committed Parquet files: they read the current
5//! on-disk rows, drop the targeted ones (by `node_uuid` / `edge_uuid`), and
6//! write the survivors back. They operate directly on a project directory and
7//! take effect immediately — there is no buffering.
8//!
9//! Each file is read, filtered against a keep-mask, and its replacement staged
10//! into a [`RewriteBatch`]; a delete's rewrites then commit **all-or-nothing**
11//! (#790), with `topology/nodes.parquet` renamed last so even a (rare)
12//! rename-phase failure leaves at worst orphaned-but-unreferenced rows, never
13//! a deleted node with surviving edges. A failure while building any
14//! replacement file leaves the prior state fully intact. Files not touched by
15//! a given delete are never opened.
16
17use std::collections::{HashMap, HashSet};
18use std::fs;
19use std::hash::BuildHasher;
20use std::path::Path;
21
22use arrow::array::{Array, FixedSizeBinaryArray, ListArray, RecordBatch, UInt32Array, UInt64Array};
23use arrow::compute::filter_record_batch;
24use arrow::datatypes::SchemaRef;
25
26use graphforge_core::GfError;
27
28use crate::catalog::{
29    discover_parquet_schema, normalize_topology_nodes, read_nodes, read_parquet_or_empty,
30};
31use crate::schemas::TOPOLOGY_NODES_SCHEMA;
32use crate::staging::RewriteBatch;
33
34/// Stage label additions for persisted nodes into the statement rewrite batch.
35/// Existing labels and the immutable primary `type_id` are preserved.
36pub fn stage_add_node_labels<S1: BuildHasher, S2: BuildHasher>(
37    staged: &mut RewriteBatch,
38    dir: &Path,
39    additions: &HashMap<[u8; 16], HashSet<u32, S2>, S1>,
40) -> Result<u64, GfError> {
41    let removals: HashMap<[u8; 16], HashSet<u32>> = HashMap::new();
42    stage_mutate_node_labels(staged, dir, additions, &removals).map(|(added, _)| added)
43}
44
45/// Stage label additions and removals for persisted nodes in one topology
46/// rewrite. The scalar primary `type_id` remains unchanged as a routing key;
47/// `type_ids` is the authoritative label-membership set.
48pub fn stage_mutate_node_labels<AS1, AS2, RS1, RS2>(
49    staged: &mut RewriteBatch,
50    dir: &Path,
51    additions: &HashMap<[u8; 16], HashSet<u32, AS2>, AS1>,
52    removals: &HashMap<[u8; 16], HashSet<u32, RS2>, RS1>,
53) -> Result<(u64, u64), GfError>
54where
55    AS1: BuildHasher,
56    AS2: BuildHasher,
57    RS1: BuildHasher,
58    RS2: BuildHasher,
59{
60    if additions.is_empty() && removals.is_empty() {
61        return Ok((0, 0));
62    }
63    let path = dir.join("topology").join("nodes.parquet");
64    let read_path = staged
65        .staged_temp(&path)
66        .map_or_else(|| path.clone(), Path::to_path_buf);
67    if !read_path.exists() {
68        return Ok((0, 0));
69    }
70    let batches = normalize_topology_nodes(
71        read_parquet_or_empty(&read_path, TOPOLOGY_NODES_SCHEMA.clone()).map_err(pq_err)?,
72    )
73    .map_err(pq_err)?;
74    let mut changed = 0u64;
75    let mut removed = 0u64;
76    let mut rebuilt = Vec::with_capacity(batches.len());
77    for batch in batches {
78        let uuids = batch
79            .column_by_name("node_uuid")
80            .and_then(|a| a.as_any().downcast_ref::<FixedSizeBinaryArray>())
81            .ok_or_else(|| GfError::Storage("node topology missing node_uuid".into()))?;
82        let labels = batch
83            .column_by_name("type_ids")
84            .and_then(|a| a.as_any().downcast_ref::<ListArray>())
85            .ok_or_else(|| GfError::Storage("node topology missing type_ids".into()))?;
86        let mut rows = Vec::with_capacity(batch.num_rows());
87        for row in 0..batch.num_rows() {
88            let values = labels.value(row);
89            let values = values
90                .as_any()
91                .downcast_ref::<UInt32Array>()
92                .ok_or_else(|| GfError::Storage("node type_ids are not UInt32".into()))?;
93            let mut merged = values.values().to_vec();
94            if let Some(drop) = removals.get(&uuid_at(uuids, row)) {
95                let before = merged.len();
96                merged.retain(|label| !drop.contains(label));
97                removed += before.saturating_sub(merged.len()) as u64;
98            }
99            if let Some(extra) = additions.get(&uuid_at(uuids, row)) {
100                let before = merged.len();
101                merged.extend(extra.iter().copied());
102                merged.sort_unstable();
103                merged.dedup();
104                changed += merged.len().saturating_sub(before) as u64;
105            }
106            rows.push(Some(merged.into_iter().map(Some).collect::<Vec<_>>()));
107        }
108        let nullable = ListArray::from_iter_primitive::<arrow::datatypes::UInt32Type, _, _>(rows);
109        let new_labels = ListArray::new(
110            std::sync::Arc::new(arrow::datatypes::Field::new(
111                "item",
112                arrow::datatypes::DataType::UInt32,
113                false,
114            )),
115            nullable.offsets().clone(),
116            nullable.values().clone(),
117            nullable.nulls().cloned(),
118        );
119        let index = batch.schema().index_of("type_ids").map_err(pq_err)?;
120        let mut columns = batch.columns().to_vec();
121        columns[index] = std::sync::Arc::new(new_labels);
122        rebuilt.push(RecordBatch::try_new(batch.schema(), columns).map_err(pq_err)?);
123    }
124    if changed > 0 || removed > 0 {
125        let merged =
126            arrow::compute::concat_batches(&TOPOLOGY_NODES_SCHEMA, &rebuilt).map_err(pq_err)?;
127        staged.restage(&path, TOPOLOGY_NODES_SCHEMA.clone(), &merged)?;
128    }
129    Ok((changed, removed))
130}
131
132fn pq_err(e: impl std::fmt::Display) -> GfError {
133    GfError::Storage(e.to_string())
134}
135
136fn io_err(e: &std::io::Error) -> GfError {
137    GfError::Storage(e.to_string())
138}
139
140/// Read a row's `FixedSizeBinary(16)` cell into a `[u8; 16]`.
141fn uuid_at(col: &FixedSizeBinaryArray, row: usize) -> [u8; 16] {
142    let mut out = [0u8; 16];
143    out.copy_from_slice(col.value(row));
144    out
145}
146
147/// Build a keep-mask: `true` for rows whose `key_col` UUID is **not** in
148/// `targets`. Returns `None` (caller skips the rewrite) when no row matches a
149/// target, so an untouched file is never rewritten.
150fn keep_mask<S: BuildHasher>(
151    batch: &RecordBatch,
152    key_col: &str,
153    targets: &HashSet<[u8; 16], S>,
154) -> Result<Option<arrow::array::BooleanArray>, GfError> {
155    let col = batch
156        .column_by_name(key_col)
157        .and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
158        .ok_or_else(|| GfError::Storage(format!("file missing {key_col} column")))?;
159    let mut any_dropped = false;
160    let mask: arrow::array::BooleanArray = (0..col.len())
161        .map(|r| {
162            let keep = !targets.contains(&uuid_at(col, r));
163            if !keep {
164                any_dropped = true;
165            }
166            Some(keep)
167        })
168        .collect();
169    Ok(any_dropped.then_some(mask))
170}
171
172/// Stage a rewrite of one Parquet file into `staged`, dropping rows whose
173/// `key_col` UUID is in `targets`. Returns the number of rows that will be
174/// removed once the batch commits. A missing file or no-match stages nothing.
175///
176/// Reads **through** `staged`: content already staged for this file in the
177/// same statement (an earlier SET/REMOVE or append) is the base, so the
178/// restaged result is the net of all the statement's effects (#792).
179fn stage_rewrite_dropping<S: BuildHasher>(
180    staged: &mut RewriteBatch,
181    path: &Path,
182    schema: SchemaRef,
183    key_col: &str,
184    targets: &HashSet<[u8; 16], S>,
185) -> Result<u64, GfError> {
186    let read_path = match staged.staged_temp(path) {
187        Some(tmp) => tmp.to_path_buf(),
188        None if !path.exists() => return Ok(0),
189        None => path.to_path_buf(),
190    };
191    let batches = read_parquet_or_empty(&read_path, schema.clone()).map_err(pq_err)?;
192    let mut removed = 0u64;
193    let mut kept: Vec<RecordBatch> = Vec::with_capacity(batches.len());
194    for batch in &batches {
195        match keep_mask(batch, key_col, targets)? {
196            Some(mask) => {
197                let before = batch.num_rows() as u64;
198                let filtered = filter_record_batch(batch, &mask).map_err(pq_err)?;
199                removed += before - filtered.num_rows() as u64;
200                kept.push(filtered);
201            }
202            None => kept.push(batch.clone()),
203        }
204    }
205    if removed == 0 {
206        return Ok(0); // nothing in this file matched — leave it untouched
207    }
208    let merged = arrow::compute::concat_batches(&schema, &kept).map_err(pq_err)?;
209    staged.restage(path, schema, &merged)?;
210    Ok(removed)
211}
212
213fn stage_rewrite_nodes_dropping<S: BuildHasher>(
214    staged: &mut RewriteBatch,
215    path: &Path,
216    targets: &HashSet<[u8; 16], S>,
217) -> Result<u64, GfError> {
218    let read_path = match staged.staged_temp(path) {
219        Some(tmp) => tmp.to_path_buf(),
220        None if !path.exists() => return Ok(0),
221        None => path.to_path_buf(),
222    };
223    let Some(stored_schema) = discover_parquet_schema(&read_path) else {
224        return Ok(0);
225    };
226    let batches = read_parquet_or_empty(&read_path, stored_schema).map_err(pq_err)?;
227    let batches = normalize_topology_nodes(batches).map_err(pq_err)?;
228    let mut removed = 0u64;
229    let mut kept = Vec::with_capacity(batches.len());
230    for batch in &batches {
231        match keep_mask(batch, "node_uuid", targets)? {
232            Some(mask) => {
233                let before = batch.num_rows() as u64;
234                let filtered = filter_record_batch(batch, &mask).map_err(pq_err)?;
235                removed += before - filtered.num_rows() as u64;
236                kept.push(filtered);
237            }
238            None => kept.push(batch.clone()),
239        }
240    }
241    if removed == 0 {
242        return Ok(0);
243    }
244    let merged = arrow::compute::concat_batches(&TOPOLOGY_NODES_SCHEMA, &kept).map_err(pq_err)?;
245    staged.restage(path, TOPOLOGY_NODES_SCHEMA.clone(), &merged)?;
246    Ok(removed)
247}
248
249/// Every `*.parquet` path directly under `dir/<subdir>`, or an empty list when
250/// the directory is absent. Enumeration order is filesystem-dependent; callers
251/// needing determinism must sort or group the results themselves.
252pub(crate) fn parquet_files_in(
253    dir: &Path,
254    subdir: &str,
255) -> Result<Vec<std::path::PathBuf>, GfError> {
256    let d = dir.join(subdir);
257    let entries = match fs::read_dir(&d) {
258        Ok(rd) => rd,
259        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
260        Err(e) => return Err(io_err(&e)),
261    };
262    let mut out = Vec::new();
263    for entry in entries {
264        let path = entry.map_err(|e| io_err(&e))?.path();
265        if path.extension().and_then(|s| s.to_str()) == Some("parquet") {
266            out.push(path);
267        }
268    }
269    Ok(out)
270}
271
272/// Stage the deletion of the given nodes into `staged`: every
273/// `properties/*.parquet` first, then `topology/nodes.parquet` **last** — the
274/// authoritative existence record commits only after everything that refers
275/// to those nodes (#790). Returns the node rows that will be removed.
276///
277/// # Errors
278/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
279#[allow(clippy::implicit_hasher)]
280pub fn stage_delete_nodes<S: BuildHasher>(
281    staged: &mut RewriteBatch,
282    dir: &Path,
283    node_uuids: &HashSet<[u8; 16], S>,
284) -> Result<u64, GfError> {
285    if node_uuids.is_empty() {
286        return Ok(0);
287    }
288    // Drop the deleted nodes' property rows so they don't dangle.
289    for path in parquet_files_in(dir, "properties")? {
290        if let Some(schema) = discover_parquet_schema(&path) {
291            stage_rewrite_dropping(staged, &path, schema, "node_uuid", node_uuids)?;
292        }
293    }
294    stage_rewrite_nodes_dropping(
295        staged,
296        &dir.join("topology").join("nodes.parquet"),
297        node_uuids,
298    )
299}
300
301/// Stage the deletion of the given edges into `staged`: every
302/// `topology/edges/*.parquet`, then any `edge_properties/*.parquet`. Returns
303/// the edge rows that will be removed.
304///
305/// # Errors
306/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
307#[allow(clippy::implicit_hasher)]
308pub fn stage_delete_edges<S: BuildHasher>(
309    staged: &mut RewriteBatch,
310    dir: &Path,
311    edge_uuids: &HashSet<[u8; 16], S>,
312) -> Result<u64, GfError> {
313    if edge_uuids.is_empty() {
314        return Ok(0);
315    }
316    let mut removed = 0u64;
317    for path in parquet_files_in(dir, "topology/edges")? {
318        if let Some(schema) = discover_parquet_schema(&path) {
319            removed += stage_rewrite_dropping(staged, &path, schema, "edge_uuid", edge_uuids)?;
320        }
321    }
322    // Drop edge-property rows (the `edge_properties/` dir exists once edge
323    // properties have been written, #784).
324    for path in parquet_files_in(dir, "edge_properties")? {
325        if let Some(schema) = discover_parquet_schema(&path) {
326            stage_rewrite_dropping(staged, &path, schema, "edge_uuid", edge_uuids)?;
327        }
328    }
329    Ok(removed)
330}
331
332/// Delete the given nodes by `node_uuid`, rewriting `topology/nodes.parquet` and
333/// dropping the same nodes' rows from every `properties/*.parquet` file.
334///
335/// All rewrites stage and commit as one batch, `topology/nodes.parquet` last
336/// (#790): a failure while building the replacement files leaves the prior
337/// state fully intact.
338///
339/// Returns the number of node rows removed. Does **not** touch edges — callers
340/// enforce openCypher's "no relationships without DETACH" rule (see
341/// [`incident_edge_uuids`]) and delete incident edges via [`delete_edges`].
342///
343/// # Errors
344/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
345pub fn delete_nodes<S: BuildHasher>(
346    dir: &Path,
347    node_uuids: &HashSet<[u8; 16], S>,
348) -> Result<u64, GfError> {
349    let mut staged = RewriteBatch::new();
350    let removed = stage_delete_nodes(&mut staged, dir, node_uuids)?;
351    if let Some(g) = crate::generation::commit_topology_aware(staged, dir)? {
352        crate::adjacency_delta::discard_segment(dir, g); // delete writes no segment
353    }
354    Ok(removed)
355}
356
357/// Delete the given edges by `edge_uuid`, rewriting every `topology/edges/*.parquet`
358/// file and dropping the same edges' rows from any `edge_properties/*.parquet`.
359///
360/// All rewrites stage and commit as one batch (#790).
361///
362/// Returns the number of edge rows removed.
363///
364/// # Errors
365/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
366pub fn delete_edges<S: BuildHasher>(
367    dir: &Path,
368    edge_uuids: &HashSet<[u8; 16], S>,
369) -> Result<u64, GfError> {
370    let mut staged = RewriteBatch::new();
371    let removed = stage_delete_edges(&mut staged, dir, edge_uuids)?;
372    if let Some(g) = crate::generation::commit_topology_aware(staged, dir)? {
373        crate::adjacency_delta::discard_segment(dir, g); // delete writes no segment
374    }
375    Ok(removed)
376}
377
378/// Delete nodes and edges as **one all-or-nothing statement** (#790): stages
379/// every rewrite — edge files, edge properties, node properties, and
380/// `topology/nodes.parquet` strictly last — then commits once. A failure while
381/// building any replacement file leaves the prior on-disk state fully intact;
382/// a (rare) rename-phase failure is bounded by the ordering to consistent
383/// states — at worst orphaned-but-unreferenced rows, never a deleted node with
384/// surviving edges.
385///
386/// Returns `(node_rows_removed, edge_rows_removed)`.
387///
388/// # Errors
389/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
390#[allow(clippy::implicit_hasher)]
391pub fn delete_nodes_and_edges<S: BuildHasher>(
392    dir: &Path,
393    node_uuids: &HashSet<[u8; 16], S>,
394    edge_uuids: &HashSet<[u8; 16], S>,
395) -> Result<(u64, u64), GfError> {
396    let mut staged = RewriteBatch::new();
397    let edges_removed = stage_delete_edges(&mut staged, dir, edge_uuids)?;
398    let nodes_removed = stage_delete_nodes(&mut staged, dir, node_uuids)?;
399    if let Some(g) = crate::generation::commit_topology_aware(staged, dir)? {
400        crate::adjacency_delta::discard_segment(dir, g); // delete writes no segment
401    }
402    Ok((nodes_removed, edges_removed))
403}
404
405/// Return the `edge_uuid`s of every edge incident to any of `node_uuids`
406/// (as `src` or `dst`), across all edge files.
407///
408/// Used to enforce openCypher's `DELETE` semantics: deleting a node that still
409/// has relationships **without** `DETACH` is an error; `DETACH DELETE` deletes
410/// these incident edges alongside the node.
411///
412/// Endpoints in the edge files are keyed by the surrogate `src_id`/`dst_id`, so
413/// this first maps the target `node_uuid`s to their `node_id`s via
414/// `topology/nodes.parquet`, then scans the edge files for those ids.
415///
416/// # Errors
417/// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
418pub fn incident_edge_uuids<S: BuildHasher>(
419    dir: &Path,
420    node_uuids: &HashSet<[u8; 16], S>,
421) -> Result<Vec<[u8; 16]>, GfError> {
422    if node_uuids.is_empty() {
423        return Ok(Vec::new());
424    }
425    // node_uuid → node_id for the targets.
426    let target_ids = node_ids_for(dir, node_uuids)?;
427    if target_ids.is_empty() {
428        return Ok(Vec::new());
429    }
430
431    let mut out = Vec::new();
432    for path in parquet_files_in(dir, "topology/edges")? {
433        let Some(schema) = discover_parquet_schema(&path) else {
434            continue;
435        };
436        for batch in read_parquet_or_empty(&path, schema).map_err(pq_err)? {
437            let edge_uuid = batch
438                .column_by_name("edge_uuid")
439                .and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
440                .ok_or_else(|| GfError::Storage("edge file missing edge_uuid".to_owned()))?;
441            let src_id = u64_col(&batch, "src_id")?;
442            let dst_id = u64_col(&batch, "dst_id")?;
443            for r in 0..batch.num_rows() {
444                if target_ids.contains(&src_id.value(r)) || target_ids.contains(&dst_id.value(r)) {
445                    out.push(uuid_at(edge_uuid, r));
446                }
447            }
448        }
449    }
450    Ok(out)
451}
452
453/// Resolve the `node_id` surrogates of the given `node_uuid`s from
454/// `topology/nodes.parquet`.
455fn node_ids_for<S: BuildHasher>(
456    dir: &Path,
457    node_uuids: &HashSet<[u8; 16], S>,
458) -> Result<HashSet<u64>, GfError> {
459    let mut ids = HashSet::new();
460    for batch in read_nodes(dir).map_err(pq_err)? {
461        let uuid = batch
462            .column_by_name("node_uuid")
463            .and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
464            .ok_or_else(|| GfError::Storage("nodes file missing node_uuid".to_owned()))?;
465        let id = u64_col(&batch, "node_id")?;
466        for r in 0..batch.num_rows() {
467            if node_uuids.contains(&uuid_at(uuid, r)) {
468                ids.insert(id.value(r));
469            }
470        }
471    }
472    Ok(ids)
473}
474
475/// Borrow a `UInt64` column by name, erroring if absent or mistyped.
476fn u64_col<'a>(batch: &'a RecordBatch, name: &str) -> Result<&'a UInt64Array, GfError> {
477    batch
478        .column_by_name(name)
479        .and_then(|c| c.as_any().downcast_ref::<UInt64Array>())
480        .ok_or_else(|| GfError::Storage(format!("column {name} missing or not UInt64")))
481}
482
483// ---------------------------------------------------------------------------
484// Tests
485// ---------------------------------------------------------------------------
486
487#[cfg(test)]
488mod tests {
489    use std::collections::HashMap;
490
491    use graphforge_core::OntologyMode;
492    use graphforge_core::TypeId;
493    use graphforge_core::uuid::{Uuid, new_v7, to_bytes};
494    use graphforge_ir::IrLiteral;
495    use tempfile::TempDir;
496
497    use super::*;
498    use crate::GraphWriter;
499
500    const TS: i64 = 1_700_000_000_000_000;
501
502    /// Total rows across the batches of a (possibly absent) parquet file.
503    fn row_count(dir: &Path, rel: &str) -> usize {
504        let path = dir.join(rel);
505        let Some(schema) = discover_parquet_schema(&path) else {
506            return 0;
507        };
508        read_parquet_or_empty(&path, schema)
509            .unwrap()
510            .iter()
511            .map(RecordBatch::num_rows)
512            .sum()
513    }
514
515    fn set(uuids: &[Uuid]) -> HashSet<[u8; 16]> {
516        uuids.iter().map(to_bytes).collect()
517    }
518
519    /// Build A--KNOWS-->B--KNOWS-->C in Strict mode; return (dir, a, b, c).
520    fn chain() -> (TempDir, Uuid, Uuid, Uuid) {
521        let dir = TempDir::new().unwrap();
522        let (a, b, c) = (new_v7(), new_v7(), new_v7());
523        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
524        w.create_node(a, TypeId(0)).unwrap();
525        w.create_node(b, TypeId(0)).unwrap();
526        w.create_node(c, TypeId(0)).unwrap();
527        w.create_edge(new_v7(), "KNOWS", &a, &b).unwrap();
528        w.create_edge(new_v7(), "KNOWS", &b, &c).unwrap();
529        w.flush().unwrap();
530        (dir, a, b, c)
531    }
532
533    #[test]
534    fn delete_edges_removes_only_targeted_rows() {
535        let dir = TempDir::new().unwrap();
536        let (a, b) = (new_v7(), new_v7());
537        let (e1, e2) = (new_v7(), new_v7());
538        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
539        w.create_node(a, TypeId(0)).unwrap();
540        w.create_node(b, TypeId(0)).unwrap();
541        w.create_edge(e1, "KNOWS", &a, &b).unwrap();
542        w.create_edge(e2, "KNOWS", &b, &a).unwrap();
543        w.flush().unwrap();
544        assert_eq!(row_count(dir.path(), "topology/edges/KNOWS.parquet"), 2);
545
546        let removed = delete_edges(dir.path(), &set(&[e1])).unwrap();
547        assert_eq!(removed, 1);
548        assert_eq!(row_count(dir.path(), "topology/edges/KNOWS.parquet"), 1);
549    }
550
551    #[test]
552    fn delete_nodes_removes_node_rows_and_leaves_edges() {
553        let (dir, a, _b, _c) = chain();
554        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 3);
555
556        let removed = delete_nodes(dir.path(), &set(&[a])).unwrap();
557        assert_eq!(removed, 1);
558        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 2);
559        // delete_nodes does not touch edges — that's the caller's job (DETACH).
560        assert_eq!(row_count(dir.path(), "topology/edges/KNOWS.parquet"), 2);
561    }
562
563    #[test]
564    fn delete_nodes_drops_their_property_rows() {
565        let dir = TempDir::new().unwrap();
566        let (a, b) = (new_v7(), new_v7());
567        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
568        w.create_node(a, TypeId(0)).unwrap();
569        w.create_node(b, TypeId(0)).unwrap();
570        w.set_properties(
571            &a,
572            None,
573            HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
574        )
575        .unwrap();
576        w.set_properties(
577            &b,
578            None,
579            HashMap::from([("name".to_owned(), IrLiteral::Str("B".into()))]),
580        )
581        .unwrap();
582        w.flush().unwrap();
583        assert_eq!(row_count(dir.path(), "properties/_untyped.parquet"), 2);
584
585        delete_nodes(dir.path(), &set(&[a])).unwrap();
586        assert_eq!(
587            row_count(dir.path(), "properties/_untyped.parquet"),
588            1,
589            "the deleted node's property row is dropped too"
590        );
591    }
592
593    #[test]
594    fn incident_edge_uuids_finds_edges_on_either_endpoint() {
595        // B is the middle of A->B->C, so both edges are incident to B.
596        let (dir, _a, b, _c) = chain();
597        let incident = incident_edge_uuids(dir.path(), &set(&[b])).unwrap();
598        assert_eq!(incident.len(), 2, "both chain edges touch B");
599    }
600
601    #[test]
602    fn incident_edge_uuids_for_leaf_finds_one_edge() {
603        // A is only the source of A->B.
604        let (dir, a, _b, _c) = chain();
605        let incident = incident_edge_uuids(dir.path(), &set(&[a])).unwrap();
606        assert_eq!(incident.len(), 1);
607    }
608
609    #[test]
610    fn detach_delete_flow_removes_node_and_incident_edges() {
611        // Emulate DETACH DELETE B: collect incident edges, delete them, delete B.
612        let (dir, _a, b, _c) = chain();
613        let incident: HashSet<[u8; 16]> = incident_edge_uuids(dir.path(), &set(&[b]))
614            .unwrap()
615            .into_iter()
616            .collect();
617        assert_eq!(delete_edges(dir.path(), &incident).unwrap(), 2);
618        assert_eq!(delete_nodes(dir.path(), &set(&[b])).unwrap(), 1);
619        assert_eq!(row_count(dir.path(), "topology/edges/KNOWS.parquet"), 0);
620        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 2);
621    }
622
623    #[test]
624    fn empty_target_sets_are_noops() {
625        let (dir, _a, _b, _c) = chain();
626        assert_eq!(delete_nodes(dir.path(), &HashSet::new()).unwrap(), 0);
627        assert_eq!(delete_edges(dir.path(), &HashSet::new()).unwrap(), 0);
628        assert!(
629            incident_edge_uuids(dir.path(), &HashSet::new())
630                .unwrap()
631                .is_empty()
632        );
633        // Nothing changed.
634        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 3);
635        assert_eq!(row_count(dir.path(), "topology/edges/KNOWS.parquet"), 2);
636    }
637
638    #[test]
639    fn deletes_bump_topology_generation() {
640        use crate::generation::read_topology_generation;
641
642        let (dir, _a, b, _c) = chain(); // the fixture's flush is bump #1
643        assert_eq!(read_topology_generation(dir.path()).unwrap(), 1);
644
645        let incident: HashSet<[u8; 16]> = incident_edge_uuids(dir.path(), &set(&[b]))
646            .unwrap()
647            .into_iter()
648            .collect();
649        assert_eq!(delete_edges(dir.path(), &incident).unwrap(), 2);
650        assert_eq!(read_topology_generation(dir.path()).unwrap(), 2);
651
652        assert_eq!(delete_nodes(dir.path(), &set(&[b])).unwrap(), 1);
653        assert_eq!(read_topology_generation(dir.path()).unwrap(), 3);
654    }
655
656    #[test]
657    fn zero_match_delete_does_not_bump_topology_generation() {
658        use crate::generation::read_topology_generation;
659
660        let (dir, _a, _b, _c) = chain();
661        assert_eq!(read_topology_generation(dir.path()).unwrap(), 1);
662
663        // No matching rows ⇒ nothing staged under topology/ ⇒ no bump.
664        assert_eq!(delete_edges(dir.path(), &set(&[new_v7()])).unwrap(), 0);
665        assert_eq!(delete_nodes(dir.path(), &set(&[new_v7()])).unwrap(), 0);
666        assert_eq!(read_topology_generation(dir.path()).unwrap(), 1);
667    }
668
669    // -----------------------------------------------------------------------
670    // Staged all-or-nothing commit (#790)
671    // -----------------------------------------------------------------------
672
673    /// Chain fixture plus a node property on `a` and an edge property on the
674    /// A→B edge, so a full delete touches all four file groups.
675    fn chain_with_properties() -> (TempDir, Uuid, Uuid, Uuid, Uuid) {
676        let dir = TempDir::new().unwrap();
677        let (a, b, c) = (new_v7(), new_v7(), new_v7());
678        let e_ab = new_v7();
679        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
680        w.create_node(a, TypeId(0)).unwrap();
681        w.create_node(b, TypeId(0)).unwrap();
682        w.create_node(c, TypeId(0)).unwrap();
683        w.create_edge(e_ab, "KNOWS", &a, &b).unwrap();
684        w.create_edge(new_v7(), "KNOWS", &b, &c).unwrap();
685        w.set_properties(
686            &a,
687            None,
688            HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
689        )
690        .unwrap();
691        w.set_edge_properties(
692            &e_ab,
693            Some("KNOWS"),
694            HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
695        )
696        .unwrap();
697        w.flush().unwrap();
698        (dir, a, b, c, e_ab)
699    }
700
701    #[test]
702    fn delete_staging_orders_nodes_parquet_last() {
703        // The one hard ordering invariant: the authoritative existence record
704        // commits only after every file that refers to the deleted entities.
705        let (dir, a, _b, _c, e_ab) = chain_with_properties();
706
707        let mut staged = RewriteBatch::new();
708        stage_delete_edges(&mut staged, dir.path(), &set(&[e_ab])).unwrap();
709        stage_delete_nodes(&mut staged, dir.path(), &set(&[a])).unwrap();
710
711        let order: Vec<_> = staged.staged_paths().collect();
712        assert!(
713            order
714                .last()
715                .is_some_and(|p| p.ends_with("topology/nodes.parquet")),
716            "nodes.parquet must commit last, got {order:?}"
717        );
718        let pos = |suffix: &str| {
719            order
720                .iter()
721                .position(|p| p.to_string_lossy().contains(suffix))
722                .unwrap_or_else(|| panic!("{suffix} not staged: {order:?}"))
723        };
724        assert!(
725            pos("topology/edges/") < pos("edge_properties/"),
726            "edge files before edge properties: {order:?}"
727        );
728        // Originals untouched while staged.
729        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 3);
730        assert_eq!(
731            row_count(dir.path(), "topology/edges/_exploratory.parquet"),
732            2
733        );
734    }
735
736    #[test]
737    fn delete_nodes_and_edges_applies_all_and_reports_counts() {
738        let (dir, a, _b, _c, e_ab) = chain_with_properties();
739
740        let (nodes_removed, edges_removed) =
741            delete_nodes_and_edges(dir.path(), &set(&[a]), &set(&[e_ab])).unwrap();
742        assert_eq!((nodes_removed, edges_removed), (1, 1));
743        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 2);
744        assert_eq!(
745            row_count(dir.path(), "topology/edges/_exploratory.parquet"),
746            1
747        );
748        assert_eq!(row_count(dir.path(), "properties/_untyped.parquet"), 0);
749        assert_eq!(row_count(dir.path(), "edge_properties/KNOWS.parquet"), 0);
750
751        // No temp residue anywhere the delete touched.
752        for sub in [
753            "topology",
754            "topology/edges",
755            "properties",
756            "edge_properties",
757        ] {
758            let residue = fs::read_dir(dir.path().join(sub))
759                .unwrap()
760                .filter_map(Result::ok)
761                .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
762                .count();
763            assert_eq!(residue, 0, "temp residue under {sub}");
764        }
765    }
766
767    #[test]
768    fn dropped_staged_delete_changes_nothing() {
769        // The abort path: stage a full delete, drop without commit — the
770        // on-disk state stays byte-identical and the temps are gone.
771        let (dir, a, _b, _c, e_ab) = chain_with_properties();
772        {
773            let mut staged = RewriteBatch::new();
774            stage_delete_edges(&mut staged, dir.path(), &set(&[e_ab])).unwrap();
775            stage_delete_nodes(&mut staged, dir.path(), &set(&[a])).unwrap();
776            // Dropped here.
777        }
778        assert_eq!(row_count(dir.path(), "topology/nodes.parquet"), 3);
779        assert_eq!(
780            row_count(dir.path(), "topology/edges/_exploratory.parquet"),
781            2
782        );
783        assert_eq!(row_count(dir.path(), "properties/_untyped.parquet"), 1);
784        assert_eq!(row_count(dir.path(), "edge_properties/KNOWS.parquet"), 1);
785    }
786}