Skip to main content

graphforge_storage/
writer.rs

1//! [`GraphWriter`] — buffered Parquet write path for the UUID-first topology /
2//! properties layout (#579).
3//!
4//! Callers mint UUIDv7 identifiers (via [`graphforge_core::uuid::new_v7`]) and feed
5//! nodes, edges, and properties to the writer; it assigns integer surrogate IDs
6//! (`node_id` / `edge_id`), buffers rows in memory, and materialises them to
7//! Parquet on [`flush`](GraphWriter::flush).  The output round-trips through
8//! [`GraphCatalog`](crate::GraphCatalog).
9//!
10//! Routing depends on [`OntologyMode`]:
11//!
12//! | | edges | node properties |
13//! |---|---|---|
14//! | Strict / Advisory | `topology/edges/TYPENAME.parquet` ([`TYPED_EDGE_SCHEMA`]) | `properties/TYPENAME.parquet` |
15//! | Exploratory | `topology/edges/_exploratory.parquet` ([`EXPLORATORY_EDGE_SCHEMA`]) | `properties/_untyped.parquet` |
16//!
17//! Edge properties (#784) are written separately under
18//! `edge_properties/REL_TYPE.parquet`, keyed by `edge_uuid` and routed by
19//! relation name in **every** mode (a dedicated directory so a relation type
20//! cannot collide with a node label of the same name in `properties/`).
21//!
22//! # Behaviour and limitations (baseline write path)
23//!
24//! 1. [`flush`](GraphWriter::flush) **merges** the buffered rows with whatever is
25//!    already on disk (read-modify-write), so separate write sessions accumulate
26//!    (#733).  Each file write is atomic (temp + rename, #790) — an I/O failure
27//!    mid-write leaves the prior file intact — but the merge is **per file**: a
28//!    failure between files commits some files and not others (nodes first, so
29//!    the partial state is consistent; durability/fsync stays out of scope for
30//!    this non-production, small-graph engine).
31//!    There is no cross-session dedup: pure `CREATE` mints fresh UUIDs, so a
32//!    `node_uuid` never recurs; MATCH…CREATE upsert is deferred to #703.
33//! 2. Surrogate `node_id` / `edge_id` values start at 1 (0 is reserved as a
34//!    sentinel) and **continue from the on-disk maximum** when a writer is opened
35//!    on an existing project, so appended rows get fresh, monotonic surrogates.
36//! 3. `_untyped` property schemas are inferred from the buffered literals (union
37//!    of property names, type from the first non-null value seen).  A column
38//!    that sees conflicting scalar types uses a tagged scalar struct so values
39//!    retain their openCypher types.
40//! 4. In Advisory / Strict mode the writer trusts the caller's `rel_type` as the
41//!    typed edge file name — it performs no ontology validation here (that lives
42//!    in the execution layer, which holds the ontology handle).
43//! 5. `_untyped` property files are not auto-registered by [`GraphCatalog`] yet;
44//!    only `node_uuid` is in its read schema until the runtime catalog learns
45//!    the columns.
46//! 6. The writer only ever creates `topology/` and `properties/` (the always-on
47//!    baseline capabilities); capability-gated directories for other features
48//!    are deferred to when those capabilities exist.
49
50use std::collections::{HashMap, HashSet};
51use std::fmt;
52use std::fs;
53use std::path::{Path, PathBuf};
54use std::sync::Arc;
55use std::time::{SystemTime, UNIX_EPOCH};
56
57use arrow::array::{
58    ArrayRef, BooleanBuilder, FixedSizeBinaryArray, Float64Array, Float64Builder, Int64Builder,
59    RecordBatch, StringArray, StringBuilder, TimestampMicrosecondArray,
60    TimestampMicrosecondBuilder, UInt32Array, UInt64Array,
61};
62use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
63use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
64
65use graphforge_core::uuid::{Uuid, to_bytes};
66use graphforge_core::{GfError, OntologyMode, TypeId};
67use graphforge_ir::IrLiteral;
68
69/// Identity, labels, and properties of a buffered node matched by MERGE.
70pub type PendingNodeMatch = ([u8; 16], u64, u32, Vec<u32>, HashMap<String, IrLiteral>);
71
72use crate::schemas::{
73    EXPLORATORY_EDGE_SCHEMA, TOPOLOGY_NODES_SCHEMA, TYPED_EDGE_SCHEMA, uuid_field,
74};
75
76/// File stem for the exploratory catch-all edge file.
77const EXPLORATORY_STEM: &str = "_exploratory";
78/// File stem for the untyped catch-all property file.
79const UNTYPED_STEM: &str = "_untyped";
80/// Sentinel UUID for "no provenance" (all-zero bytes).
81/// Join-key column name for node-property files.
82const NODE_PROPERTY_UUID_FIELD: &str = "node_uuid";
83/// Join-key column name for edge-property files.
84const EDGE_PROPERTY_UUID_FIELD: &str = "edge_uuid";
85
86// ---------------------------------------------------------------------------
87// Error helpers
88// ---------------------------------------------------------------------------
89
90fn io_err(e: &std::io::Error) -> GfError {
91    GfError::Storage(e.to_string())
92}
93
94fn pq_err(e: impl fmt::Display) -> GfError {
95    GfError::Storage(e.to_string())
96}
97
98/// Largest value of a `UInt64` column named `col` across `batches`, or `0` if
99/// the column is absent/empty. Used to continue surrogate id assignment from
100/// the on-disk maximum.
101fn max_u64_column(batches: &[RecordBatch], col: &str) -> u64 {
102    use arrow::array::Array;
103    let mut max = 0u64;
104    for batch in batches {
105        if let Some(c) = batch.column_by_name(col)
106            && let Some(ids) = c.as_any().downcast_ref::<UInt64Array>()
107        {
108            for i in 0..ids.len() {
109                if !ids.is_null(i) {
110                    max = max.max(ids.value(i));
111                }
112            }
113        }
114    }
115    max
116}
117
118// ---------------------------------------------------------------------------
119// Buffered rows
120// ---------------------------------------------------------------------------
121
122struct NodeRow {
123    node_uuid: [u8; 16],
124    node_id: u64,
125    type_id: u32,
126    type_ids: Vec<u32>,
127}
128
129struct EdgeRow {
130    edge_uuid: [u8; 16],
131    src_uuid: [u8; 16],
132    dst_uuid: [u8; 16],
133    edge_id: u64,
134    src_id: u64,
135    dst_id: u64,
136    /// `Some` for the exploratory file (carries the relation name as a column);
137    /// `None` for typed edge files.
138    rel_type_name: Option<String>,
139}
140
141struct PropRow {
142    node_uuid: [u8; 16],
143    props: HashMap<String, IrLiteral>,
144}
145
146struct EdgePropRow {
147    edge_uuid: [u8; 16],
148    props: HashMap<String, IrLiteral>,
149}
150
151/// Shared accessor over a buffered property row so the dynamic-schema inference
152/// (column ordering + type coercion) works identically for node properties
153/// (keyed by `node_uuid`) and edge properties (keyed by `edge_uuid`).
154///
155/// `props_mut` + `from_parts` additionally let the in-place SET/REMOVE rewrite
156/// (#791) mutate decoded rows and mint a fresh row for an entity that had no
157/// property file row yet — generically across both row kinds.
158trait PropRowLike {
159    fn uuid_bytes(&self) -> &[u8; 16];
160    fn props(&self) -> &HashMap<String, IrLiteral>;
161    fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral>;
162    fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self;
163}
164
165impl PropRowLike for PropRow {
166    fn uuid_bytes(&self) -> &[u8; 16] {
167        &self.node_uuid
168    }
169    fn props(&self) -> &HashMap<String, IrLiteral> {
170        &self.props
171    }
172    fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
173        &mut self.props
174    }
175    fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
176        Self {
177            node_uuid: uuid,
178            props,
179        }
180    }
181}
182
183impl PropRowLike for EdgePropRow {
184    fn uuid_bytes(&self) -> &[u8; 16] {
185        &self.edge_uuid
186    }
187    fn props(&self) -> &HashMap<String, IrLiteral> {
188        &self.props
189    }
190    fn props_mut(&mut self) -> &mut HashMap<String, IrLiteral> {
191        &mut self.props
192    }
193    fn from_parts(uuid: [u8; 16], props: HashMap<String, IrLiteral>) -> Self {
194        Self {
195            edge_uuid: uuid,
196            props,
197        }
198    }
199}
200
201// ---------------------------------------------------------------------------
202// GraphWriter
203// ---------------------------------------------------------------------------
204
205/// Buffered Parquet writer for graph topology and properties.
206///
207/// See the [module docs](self) for routing rules and limitations.
208pub struct GraphWriter {
209    dir: PathBuf,
210    mode: OntologyMode,
211    /// One timestamp captured at open time, reused for every row's
212    /// `created_at` / `updated_at` (microseconds since the Unix epoch, UTC).
213    now_micros: i64,
214    next_node_id: u64,
215    next_edge_id: u64,
216    /// Maps every `create_node` UUID to its surrogate so edges can resolve
217    /// `src_id` / `dst_id`.
218    uuid_to_node_id: HashMap<[u8; 16], u64>,
219    nodes: Vec<NodeRow>,
220    /// Keyed by edge file stem (`TYPENAME` or `_exploratory`).
221    edges: HashMap<String, Vec<EdgeRow>>,
222    /// Keyed by property file stem (`TYPENAME` or `_untyped`).
223    properties: HashMap<String, Vec<PropRow>>,
224    /// Edge properties, keyed by relation-type file stem (the rel name, e.g.
225    /// `KNOWS`), written under `edge_properties/<stem>.parquet`.
226    edge_properties: HashMap<String, Vec<EdgePropRow>>,
227    /// Edges created since the last commit, captured during `flush_edges` for
228    /// the adjacency delta segment (#765). Drained by `flush`/`take_pending_delta`.
229    pending_delta: Vec<crate::adjacency_delta::DeltaEdge>,
230}
231
232impl GraphWriter {
233    /// Open (creating if necessary) a project directory for writing.
234    ///
235    /// `mode` controls edge / property routing.  The current wall-clock time is
236    /// captured once and reused for all row timestamps.
237    ///
238    /// # Errors
239    /// Returns [`GfError::Storage`] if the directory cannot be created.
240    pub fn open(dir: &Path, mode: OntologyMode) -> Result<Self, GfError> {
241        let now_micros = SystemTime::now()
242            .duration_since(UNIX_EPOCH)
243            .map_or(0, |d| i64::try_from(d.as_micros()).unwrap_or(i64::MAX));
244        Self::open_at(dir, mode, now_micros)
245    }
246
247    /// Like [`open`](Self::open) but with an injected timestamp (microseconds
248    /// since the Unix epoch).  Used by tests for deterministic output.
249    ///
250    /// # Errors
251    /// Returns [`GfError::Storage`] if the directory cannot be created.
252    pub fn open_at(dir: &Path, mode: OntologyMode, now_micros: i64) -> Result<Self, GfError> {
253        fs::create_dir_all(dir).map_err(|e| io_err(&e))?;
254        // Continue surrogate assignment from the on-disk maximum so a writer
255        // opened on an existing project appends rather than colliding with /
256        // overwriting prior rows. Absent files → max 0 → start at 1.
257        let max_node_id =
258            max_u64_column(&crate::catalog::read_nodes(dir).map_err(pq_err)?, "node_id");
259        let max_edge_id = crate::catalog::max_edge_id(dir).map_err(pq_err)?;
260        Ok(Self {
261            dir: dir.to_path_buf(),
262            mode,
263            now_micros,
264            next_node_id: max_node_id + 1,
265            next_edge_id: max_edge_id + 1,
266            uuid_to_node_id: HashMap::new(),
267            nodes: Vec::new(),
268            edges: HashMap::new(),
269            properties: HashMap::new(),
270            edge_properties: HashMap::new(),
271            pending_delta: Vec::new(),
272        })
273    }
274
275    /// Buffer a new node and return its assigned `node_id` surrogate.
276    ///
277    /// # Errors
278    /// Currently infallible; returns `Result` for forward compatibility.
279    pub fn create_node(&mut self, node_uuid: Uuid, type_id: TypeId) -> Result<u64, GfError> {
280        self.create_node_with_labels(node_uuid, &[type_id])
281    }
282
283    /// Buffer a node with its complete label set.
284    ///
285    /// The first label is the immutable primary label used for legacy property
286    /// file routing. Label membership and `labels()` semantics use the complete
287    /// set. An empty slice creates an unlabelled node.
288    pub fn create_node_with_labels(
289        &mut self,
290        node_uuid: Uuid,
291        type_ids: &[TypeId],
292    ) -> Result<u64, GfError> {
293        let bytes = to_bytes(&node_uuid);
294        let node_id = self.next_node_id;
295        self.next_node_id += 1;
296        // Last-writer-wins on duplicate UUID (no dedup detection at this layer).
297        self.uuid_to_node_id.insert(bytes, node_id);
298        self.nodes.push(NodeRow {
299            node_uuid: bytes,
300            node_id,
301            type_id: type_ids.first().map_or(u32::MAX, |id| id.0),
302            type_ids: type_ids.iter().map(|id| id.0).collect(),
303        });
304        Ok(node_id)
305    }
306
307    /// Register an **already-persisted** node's identity so a subsequent
308    /// [`create_edge`](Self::create_edge) can resolve it as an endpoint — without
309    /// writing a new node row or minting a fresh surrogate.
310    ///
311    /// Used by mixed `MATCH … CREATE …` execution (#703): a node bound by the
312    /// preceding `MATCH` is referenced (its `node_uuid`/`node_id` come from the
313    /// matched row), not created. Unlike [`create_node`](Self::create_node), this
314    /// does **not** push a [`NodeRow`] or advance `next_node_id`; it only teaches
315    /// the UUID→surrogate map.
316    pub fn register_existing_node(&mut self, node_uuid: Uuid, node_id: u64) {
317        self.uuid_to_node_id.insert(to_bytes(&node_uuid), node_id);
318    }
319
320    /// Return the surrogate ID for a node already known to this write session.
321    /// This includes both nodes buffered earlier in the statement and persisted
322    /// nodes registered from a matched input row.
323    #[must_use]
324    pub fn node_id_for_uuid(&self, node_uuid: &Uuid) -> Option<u64> {
325        self.uuid_to_node_id.get(&to_bytes(node_uuid)).copied()
326    }
327
328    /// Buffer a new edge and return its assigned `edge_id` surrogate.
329    ///
330    /// Both endpoints must have been registered via
331    /// [`create_node`](Self::create_node) first so their `node_id` surrogates
332    /// can be resolved.
333    ///
334    /// # Errors
335    /// Returns [`GfError::Storage`] if either endpoint UUID is unknown.
336    pub fn create_edge(
337        &mut self,
338        edge_uuid: Uuid,
339        rel_type: &str,
340        src_uuid: &Uuid,
341        dst_uuid: &Uuid,
342    ) -> Result<u64, GfError> {
343        let src_bytes = to_bytes(src_uuid);
344        let dst_bytes = to_bytes(dst_uuid);
345        let src_id = *self.uuid_to_node_id.get(&src_bytes).ok_or_else(|| {
346            GfError::Storage(format!(
347                "create_edge: source {} has no node_id; call create_node first",
348                graphforge_core::uuid::to_string(src_uuid)
349            ))
350        })?;
351        let dst_id = *self.uuid_to_node_id.get(&dst_bytes).ok_or_else(|| {
352            GfError::Storage(format!(
353                "create_edge: destination {} has no node_id; call create_node first",
354                graphforge_core::uuid::to_string(dst_uuid)
355            ))
356        })?;
357
358        let edge_id = self.next_edge_id;
359        self.next_edge_id += 1;
360
361        let (stem, rel_type_name) = match self.mode {
362            OntologyMode::Exploratory => (EXPLORATORY_STEM.to_owned(), Some(rel_type.to_owned())),
363            OntologyMode::Advisory | OntologyMode::Strict => (rel_type.to_owned(), None),
364        };
365
366        self.edges.entry(stem).or_default().push(EdgeRow {
367            edge_uuid: to_bytes(&edge_uuid),
368            src_uuid: src_bytes,
369            dst_uuid: dst_bytes,
370            edge_id,
371            src_id,
372            dst_id,
373            rel_type_name,
374        });
375        Ok(edge_id)
376    }
377
378    /// Buffer a property row for a node.
379    ///
380    /// In Strict / Advisory mode with a known `entity_type`, properties route to
381    /// `properties/TYPENAME.parquet`; otherwise (exploratory, or no entity type)
382    /// they route to `properties/_untyped.parquet`.
383    ///
384    /// # Errors
385    /// Currently infallible; returns `Result` for forward compatibility.
386    pub fn set_properties(
387        &mut self,
388        node_uuid: &Uuid,
389        entity_type: Option<&str>,
390        props: HashMap<String, IrLiteral>,
391    ) -> Result<(), GfError> {
392        let stem = match (self.mode, entity_type) {
393            (OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
394            _ => UNTYPED_STEM.to_owned(),
395        };
396        self.properties.entry(stem).or_default().push(PropRow {
397            node_uuid: to_bytes(node_uuid),
398            props,
399        });
400        Ok(())
401    }
402
403    /// Buffer a property row for an edge, keyed by `edge_uuid`.
404    ///
405    /// Edge properties route to `edge_properties/<REL_TYPE>.parquet` by relation
406    /// name in **every** mode (unlike node properties, which fall back to
407    /// `_untyped` in exploratory mode). The read side resolves the file stem from
408    /// the relation name, so a single namespace keyed by rel type keeps write and
409    /// read in lock-step and avoids colliding with the node `properties/`
410    /// directory. A `None` `rel_type` (an edge created without a known relation
411    /// name) routes to the `_untyped` catch-all.
412    ///
413    /// # Errors
414    /// Currently infallible; returns `Result` for forward compatibility.
415    pub fn set_edge_properties(
416        &mut self,
417        edge_uuid: &Uuid,
418        rel_type: Option<&str>,
419        props: HashMap<String, IrLiteral>,
420    ) -> Result<(), GfError> {
421        let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
422        self.edge_properties
423            .entry(stem)
424            .or_default()
425            .push(EdgePropRow {
426                edge_uuid: to_bytes(edge_uuid),
427                props,
428            });
429        Ok(())
430    }
431
432    // -----------------------------------------------------------------------
433    // Pending-buffer inspection and mutation (#792)
434    //
435    // A mixed write statement (CREATE … DELETE/SET/REMOVE …) needs later
436    // clauses to see and edit the entities earlier clauses buffered: DELETE
437    // must find edges created in-statement (and cancel them in the buffer
438    // instead of rewriting files), and SET/REMOVE on a created entity must
439    // land in its buffered rows (a file rewrite keyed on an uncommitted uuid
440    // would miss).
441    // -----------------------------------------------------------------------
442
443    /// Whether a node with this uuid is buffered (created in this statement
444    /// and not yet flushed or cancelled).
445    #[must_use]
446    pub fn contains_pending_node(&self, node_uuid: &[u8; 16]) -> bool {
447        self.nodes.iter().any(|r| &r.node_uuid == node_uuid)
448    }
449
450    /// Return distinct label tokens on buffered nodes selected by UUID.
451    #[must_use]
452    pub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32> {
453        self.nodes
454            .iter()
455            .filter(|row| targets.contains(&row.node_uuid))
456            .flat_map(|row| row.type_ids.iter().copied())
457            .collect()
458    }
459
460    /// Materialize the currently buffered node topology without consuming it.
461    /// Statement-local reads use this as an in-memory overlay before commit.
462    ///
463    /// # Errors
464    /// Returns [`GfError::Parquet`] if the buffered values cannot form the
465    /// canonical topology batch.
466    pub fn pending_nodes_batch(&self) -> Result<RecordBatch, GfError> {
467        let n = self.nodes.len();
468        if n == 0 {
469            return Ok(RecordBatch::new_empty(TOPOLOGY_NODES_SCHEMA.clone()));
470        }
471        let uuids =
472            FixedSizeBinaryArray::try_from_iter(self.nodes.iter().map(|r| r.node_uuid.to_vec()))
473                .map_err(pq_err)?;
474        let node_ids = UInt64Array::from(self.nodes.iter().map(|r| r.node_id).collect::<Vec<_>>());
475        let type_ids = UInt32Array::from(self.nodes.iter().map(|r| r.type_id).collect::<Vec<_>>());
476        let nullable_label_sets =
477            arrow::array::ListArray::from_iter_primitive::<arrow::datatypes::UInt32Type, _, _>(
478                self.nodes
479                    .iter()
480                    .map(|row| Some(row.type_ids.iter().copied().map(Some))),
481            );
482        let label_sets = arrow::array::ListArray::new(
483            Arc::new(Field::new("item", DataType::UInt32, false)),
484            nullable_label_sets.offsets().clone(),
485            nullable_label_sets.values().clone(),
486            None,
487        );
488        let ts = self.timestamp_array(n);
489        RecordBatch::try_new(
490            TOPOLOGY_NODES_SCHEMA.clone(),
491            vec![
492                Arc::new(uuids),
493                Arc::new(node_ids),
494                Arc::new(type_ids),
495                Arc::new(label_sets),
496                Arc::new(ts.clone()),
497                Arc::new(ts),
498            ],
499        )
500        .map_err(pq_err)
501    }
502
503    /// Find a buffered node whose labels and properties satisfy a MERGE pattern.
504    #[must_use]
505    #[allow(clippy::type_complexity)]
506    pub fn find_pending_node(
507        &self,
508        labels: &[u32],
509        properties: &[(String, IrLiteral)],
510    ) -> Option<PendingNodeMatch> {
511        self.find_pending_nodes(labels, properties)
512            .into_iter()
513            .next()
514    }
515
516    /// Return every buffered node matching all requested labels and properties.
517    #[must_use]
518    pub fn find_pending_nodes(
519        &self,
520        labels: &[u32],
521        properties: &[(String, IrLiteral)],
522    ) -> Vec<PendingNodeMatch> {
523        self.nodes
524            .iter()
525            .filter_map(|node| {
526                if !labels.iter().all(|wanted| node.type_ids.contains(wanted)) {
527                    return None;
528                }
529                let props = self
530                    .properties
531                    .values()
532                    .flatten()
533                    .filter(|row| row.node_uuid == node.node_uuid)
534                    .flat_map(|row| {
535                        row.props
536                            .iter()
537                            .map(|(key, value)| (key.clone(), value.clone()))
538                    })
539                    .collect::<HashMap<_, _>>();
540                properties
541                    .iter()
542                    .all(|(name, value)| props.get(name) == Some(value))
543                    .then(|| {
544                        (
545                            node.node_uuid,
546                            node.node_id,
547                            node.type_id,
548                            node.type_ids.clone(),
549                            props,
550                        )
551                    })
552            })
553            .collect()
554    }
555
556    /// Whether an edge with this uuid is buffered.
557    #[must_use]
558    pub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool {
559        self.edges
560            .values()
561            .any(|rows| rows.iter().any(|r| &r.edge_uuid == edge_uuid))
562    }
563
564    /// Find a buffered edge matching type, endpoints, direction, and properties.
565    #[must_use]
566    #[allow(clippy::type_complexity)]
567    pub fn find_pending_edge(
568        &self,
569        rel_type: &str,
570        src: &[u8; 16],
571        dst: &[u8; 16],
572        undirected: bool,
573        properties: &[(String, IrLiteral)],
574    ) -> Option<([u8; 16], [u8; 16], [u8; 16], HashMap<String, IrLiteral>)> {
575        self.edges.iter().find_map(|(stem, edges)| {
576            edges.iter().find_map(|edge| {
577                let edge_type = edge.rel_type_name.as_deref().unwrap_or(stem);
578                let direct = edge.src_uuid == *src && edge.dst_uuid == *dst;
579                let reverse = edge.src_uuid == *dst && edge.dst_uuid == *src;
580                if edge_type != rel_type || !(direct || undirected && reverse) {
581                    return None;
582                }
583                let props = self
584                    .edge_properties
585                    .values()
586                    .flatten()
587                    .filter(|row| row.edge_uuid == edge.edge_uuid)
588                    .flat_map(|row| {
589                        row.props
590                            .iter()
591                            .map(|(key, value)| (key.clone(), value.clone()))
592                    })
593                    .collect::<HashMap<_, _>>();
594                properties
595                    .iter()
596                    .all(|(name, value)| props.get(name) == Some(value))
597                    .then_some((edge.edge_uuid, edge.src_uuid, edge.dst_uuid, props))
598            })
599        })
600    }
601
602    /// The uuids of buffered edges incident (as src or dst) to any of `nodes`.
603    ///
604    /// The pending complement of
605    /// [`incident_edge_uuids`](crate::incident_edge_uuids), which only sees
606    /// committed files: openCypher's "cannot delete a node that still has
607    /// relationships" must also count edges created earlier in the same
608    /// statement.
609    #[must_use]
610    pub fn pending_incident_edge_uuids<S: std::hash::BuildHasher>(
611        &self,
612        nodes: &HashSet<[u8; 16], S>,
613    ) -> Vec<[u8; 16]> {
614        self.edges
615            .values()
616            .flatten()
617            .filter(|r| nodes.contains(&r.src_uuid) || nodes.contains(&r.dst_uuid))
618            .map(|r| r.edge_uuid)
619            .collect()
620    }
621
622    /// Drop buffered nodes (and their buffered property rows) whose uuid is in
623    /// `targets`, so a created-then-deleted node never hits disk. Forgets the
624    /// uuid→surrogate mapping too: the entity no longer exists, so a later
625    /// `create_edge` referencing it must fail. Returns the node rows dropped.
626    pub fn cancel_nodes<S: std::hash::BuildHasher>(
627        &mut self,
628        targets: &HashSet<[u8; 16], S>,
629    ) -> u64 {
630        let before = self.nodes.len();
631        self.nodes.retain(|r| !targets.contains(&r.node_uuid));
632        let dropped = (before - self.nodes.len()) as u64;
633        // Drop emptied stems too — flush builds columns per buffered stem and
634        // a zero-row stem has nothing to build.
635        self.properties.retain(|_, rows| {
636            rows.retain(|r| !targets.contains(&r.node_uuid));
637            !rows.is_empty()
638        });
639        self.uuid_to_node_id
640            .retain(|uuid, _| !targets.contains(uuid));
641        dropped
642    }
643
644    /// Drop buffered edges (and their buffered property rows) whose uuid is in
645    /// `targets`. Returns the edge rows dropped.
646    pub fn cancel_edges<S: std::hash::BuildHasher>(
647        &mut self,
648        targets: &HashSet<[u8; 16], S>,
649    ) -> u64 {
650        let mut dropped = 0u64;
651        self.edges.retain(|_, rows| {
652            let before = rows.len();
653            rows.retain(|r| !targets.contains(&r.edge_uuid));
654            dropped += (before - rows.len()) as u64;
655            !rows.is_empty()
656        });
657        self.edge_properties.retain(|_, rows| {
658            rows.retain(|r| !targets.contains(&r.edge_uuid));
659            !rows.is_empty()
660        });
661        dropped
662    }
663
664    /// Merge `props` into the buffered property row of a pending node
665    /// (SET on an entity created earlier in this statement), inserting a row
666    /// if it has none yet. Same stem routing as
667    /// [`set_properties`](Self::set_properties).
668    pub fn merge_pending_node_props(
669        &mut self,
670        node_uuid: &[u8; 16],
671        entity_type: Option<&str>,
672        props: HashMap<String, IrLiteral>,
673    ) {
674        let stem = match (self.mode, entity_type) {
675            (OntologyMode::Advisory | OntologyMode::Strict, Some(t)) => t.to_owned(),
676            _ => UNTYPED_STEM.to_owned(),
677        };
678        let rows = self.properties.entry(stem).or_default();
679        if let Some(row) = rows.iter_mut().find(|r| &r.node_uuid == node_uuid) {
680            row.props.extend(props);
681        } else {
682            rows.push(PropRow {
683                node_uuid: *node_uuid,
684                props,
685            });
686        }
687    }
688
689    /// Add labels to a node buffered by this writer, preserving its primary label.
690    pub fn add_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
691        let Some(row) = self
692            .nodes
693            .iter_mut()
694            .find(|row| &row.node_uuid == node_uuid)
695        else {
696            return 0;
697        };
698        let before = row.type_ids.len();
699        row.type_ids.extend(labels.iter().copied());
700        row.type_ids.sort_unstable();
701        row.type_ids.dedup();
702        (row.type_ids.len() - before) as u64
703    }
704
705    /// Remove labels from a node buffered by this writer. The immutable scalar
706    /// `type_id` remains only as the property-file routing key; `type_ids` is
707    /// the authoritative membership set.
708    pub fn remove_pending_node_labels(&mut self, node_uuid: &[u8; 16], labels: &[u32]) -> u64 {
709        let Some(row) = self
710            .nodes
711            .iter_mut()
712            .find(|row| &row.node_uuid == node_uuid)
713        else {
714            return 0;
715        };
716        let before = row.type_ids.len();
717        row.type_ids.retain(|label| !labels.contains(label));
718        (before - row.type_ids.len()) as u64
719    }
720
721    /// Edge analogue of
722    /// [`merge_pending_node_props`](Self::merge_pending_node_props); same stem
723    /// routing as [`set_edge_properties`](Self::set_edge_properties).
724    pub fn merge_pending_edge_props(
725        &mut self,
726        edge_uuid: &[u8; 16],
727        rel_type: Option<&str>,
728        props: HashMap<String, IrLiteral>,
729    ) {
730        let stem = rel_type.unwrap_or(UNTYPED_STEM).to_owned();
731        let rows = self.edge_properties.entry(stem).or_default();
732        if let Some(row) = rows.iter_mut().find(|r| &r.edge_uuid == edge_uuid) {
733            row.props.extend(props);
734        } else {
735            rows.push(EdgePropRow {
736                edge_uuid: *edge_uuid,
737                props,
738            });
739        }
740    }
741
742    /// Remove `keys` from a pending node's buffered property rows (REMOVE on
743    /// an entity created earlier in this statement). Absent keys/rows are
744    /// no-ops (openCypher). Scans every stem — a REMOVE clause does not know
745    /// the routing the CREATE used.
746    pub fn remove_pending_node_props(&mut self, node_uuid: &[u8; 16], keys: &HashSet<String>) {
747        for rows in self.properties.values_mut() {
748            for row in rows.iter_mut().filter(|r| &r.node_uuid == node_uuid) {
749                row.props.retain(|k, _| !keys.contains(k));
750            }
751        }
752    }
753
754    /// Edge analogue of
755    /// [`remove_pending_node_props`](Self::remove_pending_node_props).
756    pub fn remove_pending_edge_props(&mut self, edge_uuid: &[u8; 16], keys: &HashSet<String>) {
757        for rows in self.edge_properties.values_mut() {
758            for row in rows.iter_mut().filter(|r| &r.edge_uuid == edge_uuid) {
759                row.props.retain(|k, _| !keys.contains(k));
760            }
761        }
762    }
763
764    /// Merge all buffered rows with any existing on-disk data and write the
765    /// result, then clear the row buffers.
766    ///
767    /// Only creates a subdirectory when there are rows to write into it.  Each
768    /// target file is read, concatenated with the new rows (property files are
769    /// decoded and re-inferred so the dynamic schema evolves), and rewritten —
770    /// so separate write sessions accumulate (#733).  All files stage and
771    /// commit as one batch (#790), nodes first: a failure while building any
772    /// file leaves the prior state fully intact, and a (rare) rename-phase
773    /// failure can commit a node without its edges, never the reverse.
774    ///
775    /// A batch that stages topology files bumps the project
776    /// `topology_generation` counter before committing (#759); property-only
777    /// flushes do not bump.
778    ///
779    /// # Errors
780    /// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
781    pub fn flush(&mut self) -> Result<(), GfError> {
782        let mut staged = RewriteBatch::new();
783        self.flush_into(&mut staged)?;
784        let pending = self.take_pending_delta();
785        if let Some(generation) = crate::generation::commit_topology_aware(staged, &self.dir)? {
786            // A pure-append flush (only CREATEs reach `GraphWriter`): record the
787            // delta segment so the adjacency index can serve the new edges
788            // without a rebuild. A node-only flush writes an empty segment so
789            // the chain stays contiguous. See `write_segment_best_effort`.
790            self.write_segment_best_effort(generation, &pending);
791        }
792        Ok(())
793    }
794
795    /// Drain the edges captured for the next adjacency delta segment. The
796    /// statement driver (#792) calls this after `flush_into` to write or
797    /// discard the segment around its own commit (#765).
798    #[must_use]
799    pub fn take_pending_delta(&mut self) -> Vec<crate::adjacency_delta::DeltaEdge> {
800        let mut edges = std::mem::take(&mut self.pending_delta);
801        // Ascending edge_id = creation order (edges buffer per stem, so the
802        // drain interleaves stems); the segment's documented order. Correctness
803        // does not depend on it — `apply_delta_segments` re-sorts by (key, edge).
804        edges.sort_unstable_by_key(|e| e.edge_id);
805        edges
806    }
807
808    /// Best-effort write of the delta segment for `generation` — only when the
809    /// adjacency capability directory exists (never grow `deltas/` for a project
810    /// that has no index). A failed write costs at most one future rebuild and
811    /// must never fail a already-committed flush, so the error is swallowed.
812    pub fn write_segment_best_effort(
813        &self,
814        generation: u64,
815        edges: &[crate::adjacency_delta::DeltaEdge],
816    ) {
817        if crate::adjacency::adjacency_dir(&self.dir).exists() {
818            let _ = crate::adjacency_delta::write_delta_segment(&self.dir, generation, edges);
819        }
820    }
821
822    /// Stage all buffered rows into `staged` (committed by the caller) and
823    /// clear the row buffers.
824    ///
825    /// Reads **through** `staged` and restages: a file this statement already
826    /// staged (e.g. a DELETE rewrite of the same property file, #792) is the
827    /// merge base and its entry is replaced in place with the net content —
828    /// files new to the batch append after it, so created edges commit after
829    /// `topology/nodes.parquet` whether or not a delete staged it earlier.
830    ///
831    /// On success the buffers are cleared even though nothing is committed
832    /// yet; the writer is not reusable if the caller's commit fails.
833    ///
834    /// # Errors
835    /// Returns [`GfError::Storage`] on any I/O, Arrow, or Parquet failure.
836    pub fn flush_into(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
837        self.flush_nodes(staged)?;
838        self.flush_edges(staged)?;
839        self.flush_properties(staged)?;
840        self.flush_edge_properties(staged)?;
841        Ok(())
842    }
843
844    fn flush_nodes(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
845        if self.nodes.is_empty() {
846            return Ok(());
847        }
848        let topology = self.dir.join("topology");
849        fs::create_dir_all(&topology).map_err(|e| io_err(&e))?;
850
851        let batch = self.pending_nodes_batch()?;
852
853        // Merge with any rows already on disk so separate write sessions
854        // accumulate (#733) rather than overwriting. The schema is fixed, so a
855        // concat of [existing, new] always succeeds.
856        let path = topology.join("nodes.parquet");
857        let read_path = staged
858            .staged_temp(&path)
859            .map_or_else(|| path.clone(), Path::to_path_buf);
860        let existing = crate::catalog::normalize_topology_nodes(
861            crate::catalog::read_parquet_or_empty(&read_path, TOPOLOGY_NODES_SCHEMA.clone())
862                .map_err(pq_err)?,
863        )
864        .map_err(pq_err)?;
865        let merged = concat_with_existing(&TOPOLOGY_NODES_SCHEMA, existing, batch)?;
866        staged.restage(&path, TOPOLOGY_NODES_SCHEMA.clone(), &merged)?;
867        self.nodes.clear();
868        Ok(())
869    }
870
871    fn flush_edges(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
872        if self.edges.is_empty() {
873            return Ok(());
874        }
875        let edges_dir = self.dir.join("topology").join("edges");
876        fs::create_dir_all(&edges_dir).map_err(|e| io_err(&e))?;
877
878        // Drain so we don't hold a borrow on self.edges while writing.
879        let buffered: Vec<(String, Vec<EdgeRow>)> = self.edges.drain().collect();
880        for (stem, rows) in buffered {
881            let exploratory = stem == EXPLORATORY_STEM;
882            // Capture created edges for the adjacency delta segment (#765): the
883            // typed stem is the relation name; exploratory rows carry their own.
884            for r in &rows {
885                self.pending_delta.push(crate::adjacency_delta::DeltaEdge {
886                    rel_type_name: if exploratory {
887                        r.rel_type_name.clone().unwrap_or_default()
888                    } else {
889                        stem.clone()
890                    },
891                    edge_id: r.edge_id,
892                    src_id: r.src_id,
893                    dst_id: r.dst_id,
894                });
895            }
896            let schema = if exploratory {
897                EXPLORATORY_EDGE_SCHEMA.clone()
898            } else {
899                TYPED_EDGE_SCHEMA.clone()
900            };
901            let batch = self.edge_batch(&rows, &schema, exploratory)?;
902            // Merge with this stem's existing file so appends accumulate (#733);
903            // stems not in this buffer are never opened, so they are untouched.
904            let path = edges_dir.join(format!("{stem}.parquet"));
905            let read_path = staged
906                .staged_temp(&path)
907                .map_or_else(|| path.clone(), Path::to_path_buf);
908            let existing = crate::catalog::read_parquet_or_empty(&read_path, schema.clone())
909                .map_err(pq_err)?;
910            let merged = concat_with_existing(&schema, existing, batch)?;
911            staged.restage(&path, schema, &merged)?;
912        }
913        Ok(())
914    }
915
916    fn edge_batch(
917        &self,
918        rows: &[EdgeRow],
919        schema: &SchemaRef,
920        exploratory: bool,
921    ) -> Result<RecordBatch, GfError> {
922        let edge_uuids =
923            FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.edge_uuid.to_vec()))
924                .map_err(pq_err)?;
925        let src_uuids =
926            FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.src_uuid.to_vec()))
927                .map_err(pq_err)?;
928        let dst_uuids =
929            FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.dst_uuid.to_vec()))
930                .map_err(pq_err)?;
931        let edge_ids = UInt64Array::from(rows.iter().map(|r| r.edge_id).collect::<Vec<_>>());
932        let src_ids = UInt64Array::from(rows.iter().map(|r| r.src_id).collect::<Vec<_>>());
933        let dst_ids = UInt64Array::from(rows.iter().map(|r| r.dst_id).collect::<Vec<_>>());
934        let ts = self.timestamp_array(rows.len());
935        let mut cols: Vec<ArrayRef> = vec![
936            Arc::new(edge_uuids),
937            Arc::new(src_uuids),
938            Arc::new(dst_uuids),
939            Arc::new(edge_ids),
940            Arc::new(src_ids),
941            Arc::new(dst_ids),
942            Arc::new(ts),
943        ];
944        if exploratory {
945            let names = StringArray::from(
946                rows.iter()
947                    .map(|r| r.rel_type_name.clone().unwrap_or_default())
948                    .collect::<Vec<_>>(),
949            );
950            cols.push(Arc::new(names));
951        }
952        RecordBatch::try_new(schema.clone(), cols).map_err(pq_err)
953    }
954
955    fn flush_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
956        if self.properties.is_empty() {
957            return Ok(());
958        }
959        let buffered: Vec<(String, Vec<PropRow>)> = self.properties.drain().collect();
960        for (stem, new_rows) in buffered {
961            // Decode any rows already on disk and prepend them, then re-run the
962            // schema inference over the combined set (#733). Re-using
963            // `build_property_columns` keeps one source of truth for the
964            // first-seen ordering and type-coercion rules across flushes.
965            let existing = read_props_through(staged, &node_props_path(&self.dir, &stem))?;
966            let mut rows = decode_property_rows(&existing)?;
967            rows.extend(new_rows);
968            let (schema, cols) = build_property_columns(&stem, &rows)?;
969            stage_property_file(staged, &self.dir, "properties", &stem, schema, cols)?;
970        }
971        Ok(())
972    }
973
974    /// Edge analogue of [`flush_properties`](Self::flush_properties): merge the
975    /// buffered edge-property rows with any on-disk rows (decode + re-infer) and
976    /// stage `edge_properties/<stem>.parquet`. The join key is `edge_uuid`.
977    fn flush_edge_properties(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError> {
978        if self.edge_properties.is_empty() {
979            return Ok(());
980        }
981        let buffered: Vec<(String, Vec<EdgePropRow>)> = self.edge_properties.drain().collect();
982        for (stem, new_rows) in buffered {
983            let existing = read_props_through(staged, &edge_props_path(&self.dir, &stem))?;
984            let mut rows = decode_edge_property_rows(&existing)?;
985            rows.extend(new_rows);
986            let (schema, cols) = build_property_columns_keyed(
987                EDGE_PROPERTY_UUID_FIELD,
988                "graphforge.rel_type",
989                &stem,
990                &rows,
991            )?;
992            stage_property_file(staged, &self.dir, "edge_properties", &stem, schema, cols)?;
993        }
994        Ok(())
995    }
996
997    fn timestamp_array(&self, n: usize) -> TimestampMicrosecondArray {
998        TimestampMicrosecondArray::from(vec![self.now_micros; n])
999            .with_timezone_opt(Some(Arc::from("UTC")))
1000    }
1001}
1002
1003// ---------------------------------------------------------------------------
1004// Property schema inference
1005// ---------------------------------------------------------------------------
1006
1007/// Arrow type a property column is built as, inferred from the literals seen.
1008// Not `Copy`: `List` boxes its inner type (a homogeneous `List<inner>` column).
1009#[derive(Clone, PartialEq, Eq)]
1010enum ColType {
1011    Int,
1012    Float,
1013    Bool,
1014    Str,
1015    HetScalar,
1016    Duration,
1017    DateTime,
1018    Date,
1019    LocalDateTime,
1020    Time,
1021    ZonedTime,
1022    ZonedDateTime,
1023    /// A homogeneous `List<inner>` column (#1006).
1024    List(Box<ColType>),
1025}
1026
1027impl ColType {
1028    fn of(lit: &IrLiteral) -> Option<Self> {
1029        match lit {
1030            IrLiteral::Null
1031            // Query-parameter maps are not a storage property type.
1032            | IrLiteral::Map(_)
1033            // Typed UUID parameters are identity predicates, not properties.
1034            | IrLiteral::Uuid(_) => None,
1035            IrLiteral::Int(_) => Some(Self::Int),
1036            IrLiteral::Float(_) => Some(Self::Float),
1037            IrLiteral::Bool(_) => Some(Self::Bool),
1038            IrLiteral::Str(_) => Some(Self::Str),
1039            IrLiteral::Duration { .. } => Some(Self::Duration),
1040            IrLiteral::DateTime(_) => Some(Self::DateTime),
1041            IrLiteral::Date(_) => Some(Self::Date),
1042            IrLiteral::LocalDateTime { .. } => Some(Self::LocalDateTime),
1043            IrLiteral::Time(_) => Some(Self::Time),
1044            IrLiteral::ZonedTime { .. } => Some(Self::ZonedTime),
1045            IrLiteral::ZonedDateTime { .. } => Some(Self::ZonedDateTime),
1046            // A homogeneous list: infer the inner type from the first non-null
1047            // element. A list whose elements are all null (or an empty list)
1048            // yields no type and the column falls back to `Str`. (#1006)
1049            IrLiteral::List(items) => items
1050                .iter()
1051                .find_map(Self::of)
1052                .map(|inner| Self::List(Box::new(inner))),
1053        }
1054    }
1055
1056    fn data_type(&self) -> DataType {
1057        match self {
1058            Self::Int => DataType::Int64,
1059            Self::Float => DataType::Float64,
1060            Self::Bool => DataType::Boolean,
1061            // `Str` is also the coercion target for mixed-type columns.
1062            Self::Str => DataType::Utf8,
1063            Self::HetScalar => DataType::Struct(heterogeneous_scalar_fields()),
1064            Self::Duration => DataType::Struct(crate::schemas::duration_struct_fields()),
1065            Self::DateTime => DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1066            Self::Date => DataType::Struct(crate::schemas::date_struct_fields()),
1067            Self::LocalDateTime => DataType::Struct(crate::schemas::localdatetime_struct_fields()),
1068            Self::Time => DataType::Time64(TimeUnit::Nanosecond),
1069            Self::ZonedTime => DataType::Struct(crate::schemas::time_struct_fields()),
1070            Self::ZonedDateTime => DataType::Struct(crate::schemas::datetime_struct_fields()),
1071            Self::List(inner) => {
1072                DataType::List(Arc::new(Field::new("item", inner.data_type(), true)))
1073            }
1074        }
1075    }
1076
1077    fn is_scalar(&self) -> bool {
1078        matches!(
1079            self,
1080            Self::Int | Self::Float | Self::Bool | Self::Str | Self::HetScalar
1081        )
1082    }
1083}
1084
1085fn heterogeneous_scalar_fields() -> arrow::datatypes::Fields {
1086    arrow::datatypes::Fields::from(vec![
1087        Field::new("__het_tag", DataType::Int8, false),
1088        Field::new("__het_int", DataType::Int64, true),
1089        Field::new("__het_float", DataType::Float64, true),
1090        Field::new("__het_str", DataType::Utf8, true),
1091        Field::new("__het_bool", DataType::Boolean, true),
1092    ])
1093}
1094
1095/// Build the dynamic schema and column arrays for a **node**-property file
1096/// (join key `node_uuid`, metadata key `graphforge.entity_type`).
1097fn build_property_columns(
1098    entity_type: &str,
1099    rows: &[PropRow],
1100) -> Result<(Schema, Vec<ArrayRef>), GfError> {
1101    build_property_columns_keyed(
1102        NODE_PROPERTY_UUID_FIELD,
1103        "graphforge.entity_type",
1104        entity_type,
1105        rows,
1106    )
1107}
1108
1109/// Build the dynamic schema and column arrays for a property file, keyed by an
1110/// arbitrary uuid join column.
1111///
1112/// Shared by node properties (`node_uuid`) and edge properties (`edge_uuid`).
1113/// Column order is the first-seen order of property names across `rows`
1114/// (deterministic).  Each column's type is inferred from the first non-null
1115/// value; conflicting scalar types use a tagged struct that preserves each value.
1116///
1117/// `uuid_field_name` is the leading join-key column; `meta_key`/`meta_value`
1118/// is the schema-level metadata identifying the file's entity or relation type.
1119fn build_property_columns_keyed<R: PropRowLike>(
1120    uuid_field_name: &str,
1121    meta_key: &str,
1122    meta_value: &str,
1123    rows: &[R],
1124) -> Result<(Schema, Vec<ArrayRef>), GfError> {
1125    // First-seen-ordered list of property names + inferred column type.
1126    // `seen` tracks column order independently of `col_types`: a column may be
1127    // seen (ordered) before any concrete value fixes its type, so order-dedup
1128    // must not key on `col_types` membership.
1129    let mut order: Vec<String> = Vec::new();
1130    let mut seen: HashSet<String> = HashSet::new();
1131    let mut col_types: HashMap<String, ColType> = HashMap::new();
1132    for row in rows {
1133        for (name, lit) in row.props() {
1134            reject_map_property_value(name, lit)?;
1135            if seen.insert(name.clone()) {
1136                order.push(name.clone());
1137            }
1138            // Only concrete literals contribute a type; `Null` contributes
1139            // none, so the first *non-null* value determines the column type.
1140            // A column that never sees a concrete value defaults to `Str` via
1141            // `unwrap_or(ColType::Str)` below.
1142            if let Some(t) = ColType::of(lit) {
1143                col_types
1144                    .entry(name.clone())
1145                    .and_modify(|existing| {
1146                        if *existing != t {
1147                            *existing = if existing.is_scalar() && t.is_scalar() {
1148                                ColType::HetScalar
1149                            } else {
1150                                ColType::Str
1151                            };
1152                        }
1153                    })
1154                    .or_insert(t);
1155            }
1156        }
1157    }
1158
1159    let mut fields: Vec<Field> = vec![uuid_field(uuid_field_name)];
1160    for name in &order {
1161        let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
1162        fields.push(Field::new(name, ct.data_type(), true));
1163    }
1164    let meta: HashMap<String, String> = [(meta_key.to_owned(), meta_value.to_owned())]
1165        .into_iter()
1166        .collect();
1167    let schema = Schema::new(fields).with_metadata(meta);
1168
1169    let mut cols: Vec<ArrayRef> = Vec::with_capacity(order.len() + 1);
1170    let uuids = FixedSizeBinaryArray::try_from_iter(rows.iter().map(|r| r.uuid_bytes().to_vec()))
1171        .map_err(pq_err)?;
1172    cols.push(Arc::new(uuids));
1173
1174    for name in &order {
1175        let ct = col_types.get(name).cloned().unwrap_or(ColType::Str);
1176        cols.push(build_property_array(name, ct, rows));
1177    }
1178    Ok((schema, cols))
1179}
1180
1181fn reject_map_property_value(name: &str, lit: &IrLiteral) -> Result<(), GfError> {
1182    if contains_uuid_literal(lit) {
1183        return Err(GfError::Validation(format!(
1184            "property `{name}` cannot store typed UUID query parameters"
1185        )));
1186    }
1187    if contains_map_literal(lit) {
1188        return Err(GfError::Storage(format!(
1189            "property `{name}` cannot store map values"
1190        )));
1191    }
1192    Ok(())
1193}
1194
1195fn contains_map_literal(lit: &IrLiteral) -> bool {
1196    match lit {
1197        IrLiteral::Map(_) => true,
1198        IrLiteral::List(items) => items.iter().any(contains_map_literal),
1199        _ => false,
1200    }
1201}
1202
1203fn contains_uuid_literal(lit: &IrLiteral) -> bool {
1204    match lit {
1205        IrLiteral::Uuid(_) => true,
1206        IrLiteral::List(items) => items.iter().any(contains_uuid_literal),
1207        IrLiteral::Map(entries) => entries
1208            .iter()
1209            .any(|(_, value)| contains_uuid_literal(value)),
1210        _ => false,
1211    }
1212}
1213
1214/// Build one nullable property column, appending nulls for rows that omit the
1215/// property (or whose value does not match the column's inferred type — those
1216/// are stringified when the column is `Str`, else null).
1217#[allow(
1218    clippy::too_many_lines,
1219    reason = "one builder arm per ColType; the per-type append loops read clearest inline"
1220)]
1221fn build_property_array<R: PropRowLike>(name: &str, ct: ColType, rows: &[R]) -> ArrayRef {
1222    match ct {
1223        ColType::Int => {
1224            let mut b = Int64Builder::new();
1225            for row in rows {
1226                match row.props().get(name) {
1227                    Some(IrLiteral::Int(v)) => b.append_value(*v),
1228                    _ => b.append_null(),
1229                }
1230            }
1231            Arc::new(b.finish())
1232        }
1233        ColType::Float => {
1234            let mut b = Float64Builder::new();
1235            for row in rows {
1236                match row.props().get(name) {
1237                    Some(IrLiteral::Float(v)) => b.append_value(*v),
1238                    _ => b.append_null(),
1239                }
1240            }
1241            Arc::new(b.finish())
1242        }
1243        ColType::Bool => {
1244            let mut b = BooleanBuilder::new();
1245            for row in rows {
1246                match row.props().get(name) {
1247                    Some(IrLiteral::Bool(v)) => b.append_value(*v),
1248                    _ => b.append_null(),
1249                }
1250            }
1251            Arc::new(b.finish())
1252        }
1253        ColType::Duration => {
1254            // A typed duration is a `Struct{months, days, seconds, nanos}` (all
1255            // Int64; Parquet cannot persist Arrow `Interval`); shared field defs
1256            // via `duration_struct_fields`.
1257            use arrow::array::{Int64Builder, StructArray};
1258            use arrow::buffer::NullBuffer;
1259            let (mut mb, mut db, mut sb, mut nb) = (
1260                Int64Builder::new(),
1261                Int64Builder::new(),
1262                Int64Builder::new(),
1263                Int64Builder::new(),
1264            );
1265            let mut valid = Vec::with_capacity(rows.len());
1266            for row in rows {
1267                if let Some(IrLiteral::Duration {
1268                    months,
1269                    days,
1270                    seconds,
1271                    nanos,
1272                }) = row.props().get(name)
1273                {
1274                    mb.append_value(*months);
1275                    db.append_value(*days);
1276                    sb.append_value(*seconds);
1277                    nb.append_value(*nanos);
1278                    valid.push(true);
1279                } else {
1280                    mb.append_null();
1281                    db.append_null();
1282                    sb.append_null();
1283                    nb.append_null();
1284                    valid.push(false);
1285                }
1286            }
1287            Arc::new(StructArray::new(
1288                crate::schemas::duration_struct_fields(),
1289                vec![
1290                    Arc::new(mb.finish()),
1291                    Arc::new(db.finish()),
1292                    Arc::new(sb.finish()),
1293                    Arc::new(nb.finish()),
1294                ],
1295                Some(NullBuffer::from(valid)),
1296            ))
1297        }
1298        ColType::DateTime => {
1299            let mut b = TimestampMicrosecondBuilder::new();
1300            for row in rows {
1301                match row.props().get(name) {
1302                    Some(IrLiteral::DateTime(v)) => b.append_value(*v),
1303                    _ => b.append_null(),
1304                }
1305            }
1306            Arc::new(b.finish().with_timezone_opt(Some(Arc::from("UTC"))))
1307        }
1308        ColType::Date => {
1309            // A typed `date` is a self-describing `Struct{epoch_day: Int64}` —
1310            // i64 days, full year range (#1011).
1311            use arrow::array::{Int64Builder, StructArray};
1312            use arrow::buffer::NullBuffer;
1313            let mut b = Int64Builder::new();
1314            let mut valid = Vec::with_capacity(rows.len());
1315            for row in rows {
1316                if let Some(IrLiteral::Date(v)) = row.props().get(name) {
1317                    b.append_value(*v);
1318                    valid.push(true);
1319                } else {
1320                    b.append_null();
1321                    valid.push(false);
1322                }
1323            }
1324            Arc::new(StructArray::new(
1325                crate::schemas::date_struct_fields(),
1326                vec![Arc::new(b.finish())],
1327                Some(NullBuffer::from(valid)),
1328            ))
1329        }
1330        ColType::LocalDateTime => {
1331            // A typed `localdatetime` is a `Struct{date: Int64, time: Time64(ns)}`
1332            // (shared field defs via `localdatetime_struct_fields`).
1333            use arrow::array::{Int64Builder, StructArray, Time64NanosecondBuilder};
1334            use arrow::buffer::NullBuffer;
1335            let (mut date_b, mut time_b) = (Int64Builder::new(), Time64NanosecondBuilder::new());
1336            let mut valid = Vec::with_capacity(rows.len());
1337            for row in rows {
1338                if let Some(IrLiteral::LocalDateTime { days, nanos }) = row.props().get(name) {
1339                    date_b.append_value(*days);
1340                    time_b.append_value(*nanos);
1341                    valid.push(true);
1342                } else {
1343                    date_b.append_null();
1344                    time_b.append_null();
1345                    valid.push(false);
1346                }
1347            }
1348            Arc::new(StructArray::new(
1349                crate::schemas::localdatetime_struct_fields(),
1350                vec![Arc::new(date_b.finish()), Arc::new(time_b.finish())],
1351                Some(NullBuffer::from(valid)),
1352            ))
1353        }
1354        ColType::Time => {
1355            use arrow::array::Time64NanosecondBuilder;
1356            let mut b = Time64NanosecondBuilder::new();
1357            for row in rows {
1358                match row.props().get(name) {
1359                    Some(IrLiteral::Time(v)) => b.append_value(*v),
1360                    _ => b.append_null(),
1361                }
1362            }
1363            Arc::new(b.finish())
1364        }
1365        ColType::ZonedTime => {
1366            // A typed `time` is a `Struct{time: Time64(ns), offset: Int32}`.
1367            use arrow::array::{Int32Builder, StructArray, Time64NanosecondBuilder};
1368            use arrow::buffer::NullBuffer;
1369            let (mut time_b, mut off_b) = (Time64NanosecondBuilder::new(), Int32Builder::new());
1370            let mut valid = Vec::with_capacity(rows.len());
1371            for row in rows {
1372                if let Some(IrLiteral::ZonedTime { nanos, offset }) = row.props().get(name) {
1373                    time_b.append_value(*nanos);
1374                    off_b.append_value(*offset);
1375                    valid.push(true);
1376                } else {
1377                    time_b.append_null();
1378                    off_b.append_null();
1379                    valid.push(false);
1380                }
1381            }
1382            Arc::new(StructArray::new(
1383                crate::schemas::time_struct_fields(),
1384                vec![Arc::new(time_b.finish()), Arc::new(off_b.finish())],
1385                Some(NullBuffer::from(valid)),
1386            ))
1387        }
1388        ColType::ZonedDateTime => {
1389            // A typed `datetime` is a
1390            // `Struct{date: Int64, time: Time64(ns), offset: Int32, zone: Utf8}`.
1391            use arrow::array::{
1392                Int32Builder, Int64Builder, StringBuilder, StructArray, Time64NanosecondBuilder,
1393            };
1394            use arrow::buffer::NullBuffer;
1395            let (mut date_b, mut time_b, mut off_b, mut zone_b) = (
1396                Int64Builder::new(),
1397                Time64NanosecondBuilder::new(),
1398                Int32Builder::new(),
1399                StringBuilder::new(),
1400            );
1401            let mut valid = Vec::with_capacity(rows.len());
1402            for row in rows {
1403                if let Some(IrLiteral::ZonedDateTime {
1404                    days,
1405                    nanos,
1406                    offset,
1407                    zone,
1408                }) = row.props().get(name)
1409                {
1410                    date_b.append_value(*days);
1411                    time_b.append_value(*nanos);
1412                    off_b.append_value(*offset);
1413                    // None (offset-only) is stored as a NULL zone, not "".
1414                    zone_b.append_option(zone.as_deref());
1415                    valid.push(true);
1416                } else {
1417                    date_b.append_null();
1418                    time_b.append_null();
1419                    off_b.append_null();
1420                    zone_b.append_null();
1421                    valid.push(false);
1422                }
1423            }
1424            Arc::new(StructArray::new(
1425                crate::schemas::datetime_struct_fields(),
1426                vec![
1427                    Arc::new(date_b.finish()),
1428                    Arc::new(time_b.finish()),
1429                    Arc::new(off_b.finish()),
1430                    Arc::new(zone_b.finish()),
1431                ],
1432                Some(NullBuffer::from(valid)),
1433            ))
1434        }
1435        ColType::Str => {
1436            let mut b = StringBuilder::new();
1437            for row in rows {
1438                match row.props().get(name) {
1439                    Some(IrLiteral::Null) | None => b.append_null(),
1440                    Some(other) => b.append_value(literal_to_string(other)),
1441                }
1442            }
1443            Arc::new(b.finish())
1444        }
1445        ColType::HetScalar => build_heterogeneous_scalar_array(name, rows),
1446        ColType::List(inner) => {
1447            // Flatten every list's elements into one-property synthetic rows, then
1448            // build the child array with the SAME per-type machinery (so temporal
1449            // element types reuse their struct builders); `offsets` delimit each
1450            // row's slice, and a row that is not a list is a null list slot. (#1006)
1451            use arrow::array::ListArray;
1452            use arrow::buffer::{NullBuffer, OffsetBuffer};
1453            let mut elem_rows: Vec<PropRow> = Vec::new();
1454            let mut offsets: Vec<i32> = vec![0];
1455            let mut valid = Vec::with_capacity(rows.len());
1456            for row in rows {
1457                if let Some(IrLiteral::List(items)) = row.props().get(name) {
1458                    for it in items {
1459                        let mut props = HashMap::with_capacity(1);
1460                        props.insert("item".to_string(), it.clone());
1461                        elem_rows.push(PropRow {
1462                            node_uuid: [0u8; 16],
1463                            props,
1464                        });
1465                    }
1466                    valid.push(true);
1467                } else {
1468                    valid.push(false);
1469                }
1470                offsets.push(i32::try_from(elem_rows.len()).unwrap_or(i32::MAX));
1471            }
1472            let child = build_property_array("item", (*inner).clone(), &elem_rows);
1473            let field = Arc::new(Field::new("item", inner.data_type(), true));
1474            Arc::new(ListArray::new(
1475                field,
1476                OffsetBuffer::new(offsets.into()),
1477                child,
1478                Some(NullBuffer::from(valid)),
1479            ))
1480        }
1481    }
1482}
1483
1484fn build_heterogeneous_scalar_array<R: PropRowLike>(name: &str, rows: &[R]) -> ArrayRef {
1485    use arrow::array::{BooleanBuilder, Float64Builder, Int8Builder, Int64Builder, StructArray};
1486    use arrow::buffer::NullBuffer;
1487
1488    let mut tags = Int8Builder::new();
1489    let mut ints = Int64Builder::new();
1490    let mut floats = Float64Builder::new();
1491    let mut strings = StringBuilder::new();
1492    let mut bools = BooleanBuilder::new();
1493    let mut valid = Vec::with_capacity(rows.len());
1494    for row in rows {
1495        let value = row.props().get(name);
1496        let tag = match value {
1497            Some(IrLiteral::Int(value)) => {
1498                ints.append_value(*value);
1499                floats.append_null();
1500                strings.append_null();
1501                bools.append_null();
1502                Some(0)
1503            }
1504            Some(IrLiteral::Float(value)) => {
1505                ints.append_null();
1506                floats.append_value(*value);
1507                strings.append_null();
1508                bools.append_null();
1509                Some(1)
1510            }
1511            Some(IrLiteral::Str(value)) => {
1512                ints.append_null();
1513                floats.append_null();
1514                strings.append_value(value);
1515                bools.append_null();
1516                Some(2)
1517            }
1518            Some(IrLiteral::Bool(value)) => {
1519                ints.append_null();
1520                floats.append_null();
1521                strings.append_null();
1522                bools.append_value(*value);
1523                Some(3)
1524            }
1525            _ => {
1526                ints.append_null();
1527                floats.append_null();
1528                strings.append_null();
1529                bools.append_null();
1530                None
1531            }
1532        };
1533        tags.append_value(tag.unwrap_or_default());
1534        valid.push(tag.is_some());
1535    }
1536    Arc::new(StructArray::new(
1537        heterogeneous_scalar_fields(),
1538        vec![
1539            Arc::new(tags.finish()),
1540            Arc::new(ints.finish()),
1541            Arc::new(floats.finish()),
1542            Arc::new(strings.finish()),
1543            Arc::new(bools.finish()),
1544        ],
1545        Some(NullBuffer::from(valid)),
1546    ))
1547}
1548
1549/// Stringify a literal for a `Utf8`-coerced (mixed-type) property column.
1550fn literal_to_string(lit: &IrLiteral) -> String {
1551    match lit {
1552        IrLiteral::Null => String::new(),
1553        IrLiteral::Bool(b) => b.to_string(),
1554        IrLiteral::Int(i) => i.to_string(),
1555        IrLiteral::Float(f) => f.to_string(),
1556        IrLiteral::Str(s) => s.clone(),
1557        IrLiteral::Uuid(bytes) => {
1558            let mut encoded = String::with_capacity(32);
1559            for byte in bytes {
1560                std::fmt::Write::write_fmt(&mut encoded, format_args!("{byte:02x}"))
1561                    .expect("writing to a String cannot fail");
1562            }
1563            encoded
1564        }
1565        // A duration in a mixed (stringified) column: a deterministic
1566        // months/days/seconds/nanos form (the canonical `P…` render lives in graphforge-rel).
1567        IrLiteral::Duration {
1568            months,
1569            days,
1570            seconds,
1571            nanos,
1572        } => format!("{months}mo{days}d{seconds}s{nanos}ns"),
1573        IrLiteral::DateTime(t) => t.to_string(),
1574        IrLiteral::Date(d) => d.to_string(),
1575        // Temporal values in a mixed (stringified) column: deterministic forms
1576        // (the canonical renders live in graphforge-rel).
1577        IrLiteral::LocalDateTime { days, nanos } => format!("{days}d{nanos}ns"),
1578        IrLiteral::Time(nanos) => format!("{nanos}ns"),
1579        IrLiteral::ZonedTime { nanos, offset } => format!("{nanos}ns{offset:+}s"),
1580        IrLiteral::ZonedDateTime {
1581            days,
1582            nanos,
1583            offset,
1584            zone,
1585        } => format!(
1586            "{days}d{nanos}ns{offset:+}s{}",
1587            zone.as_deref().unwrap_or("")
1588        ),
1589        // A list in a mixed (stringified) column: a deterministic bracketed form.
1590        IrLiteral::List(items) => {
1591            let parts: Vec<String> = items.iter().map(literal_to_string).collect();
1592            format!("[{}]", parts.join(","))
1593        }
1594        IrLiteral::Map(entries) => {
1595            let parts: Vec<String> = entries
1596                .iter()
1597                .map(|(key, value)| format!("{key}:{}", literal_to_string(value)))
1598                .collect();
1599            format!("{{{}}}", parts.join(","))
1600        }
1601    }
1602}
1603
1604/// Decode a previously-written property Parquet (the dynamic per-stem schema
1605/// `build_property_columns` produces) back into [`PropRow`]s, so a later flush
1606/// can re-run inference over `[decoded ++ new]` and merge (#733).
1607///
1608/// `node_uuid` (the join key, `FixedSizeBinary(16)`) populates `PropRow.node_uuid`;
1609/// every other column maps by its Arrow type to an [`IrLiteral`]. A **null** slot
1610/// omits the key for that row — it must never become a concrete value, or it
1611/// would wrongly pin a previously-all-null column's type on re-inference.
1612///
1613/// `IrLiteral` is a closed 7-variant set, so the writer only ever produces these
1614/// column types; an unexpected type is a defensive error.
1615fn decode_property_rows(batches: &[RecordBatch]) -> Result<Vec<PropRow>, GfError> {
1616    let mut out = Vec::new();
1617    for batch in batches {
1618        decode_property_batch(batch, NODE_PROPERTY_UUID_FIELD, |node_uuid, props| {
1619            out.push(PropRow { node_uuid, props });
1620        })?;
1621    }
1622    Ok(out)
1623}
1624
1625/// Edge analogue of [`decode_property_rows`]: decode `edge_properties/*.parquet`
1626/// back into [`EdgePropRow`]s (join key `edge_uuid`) for the read-merge-rewrite
1627/// flush cycle.
1628fn decode_edge_property_rows(batches: &[RecordBatch]) -> Result<Vec<EdgePropRow>, GfError> {
1629    let mut out = Vec::new();
1630    for batch in batches {
1631        decode_property_batch(batch, EDGE_PROPERTY_UUID_FIELD, |edge_uuid, props| {
1632            out.push(EdgePropRow { edge_uuid, props });
1633        })?;
1634    }
1635    Ok(out)
1636}
1637
1638/// Read the non-null property keys currently stored for one entity.
1639///
1640/// Used by `SET entity = map` to compute the authoritative replacement
1641/// complement even when the query plan projected only a subset of properties.
1642pub fn read_entity_property_keys(
1643    dir: &Path,
1644    stem: &str,
1645    uuid: &[u8; 16],
1646    is_edge: bool,
1647) -> Result<HashSet<String>, GfError> {
1648    let path = if is_edge {
1649        edge_props_path(dir, stem)
1650    } else {
1651        node_props_path(dir, stem)
1652    };
1653    let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1654        return Ok(HashSet::new());
1655    };
1656    let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1657    let rows = if is_edge {
1658        decode_edge_property_rows(&batches)?
1659            .into_iter()
1660            .map(|row| (row.edge_uuid, row.props))
1661            .collect::<Vec<_>>()
1662    } else {
1663        decode_property_rows(&batches)?
1664            .into_iter()
1665            .map(|row| (row.node_uuid, row.props))
1666            .collect::<Vec<_>>()
1667    };
1668    Ok(rows
1669        .into_iter()
1670        .find_map(|(row_uuid, props)| (row_uuid == *uuid).then(|| props.into_keys().collect()))
1671        .unwrap_or_default())
1672}
1673
1674/// Read the complete non-null property map for one persisted entity.
1675///
1676/// Returns an empty map when the property file or entity row is absent.
1677pub fn read_entity_properties(
1678    dir: &Path,
1679    stem: &str,
1680    uuid: &[u8; 16],
1681    is_edge: bool,
1682) -> Result<HashMap<String, IrLiteral>, GfError> {
1683    let path = if is_edge {
1684        edge_props_path(dir, stem)
1685    } else {
1686        node_props_path(dir, stem)
1687    };
1688    let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1689        return Ok(HashMap::new());
1690    };
1691    let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1692    let rows = if is_edge {
1693        decode_edge_property_rows(&batches)?
1694            .into_iter()
1695            .map(|row| (row.edge_uuid, row.props))
1696            .collect::<Vec<_>>()
1697    } else {
1698        decode_property_rows(&batches)?
1699            .into_iter()
1700            .map(|row| (row.node_uuid, row.props))
1701            .collect::<Vec<_>>()
1702    };
1703    Ok(rows
1704        .into_iter()
1705        .find_map(|(row_uuid, props)| (row_uuid == *uuid).then_some(props))
1706        .unwrap_or_default())
1707}
1708
1709/// Read every UUID-keyed node property row from one persisted stem.
1710pub fn read_node_property_rows(
1711    dir: &Path,
1712    stem: &str,
1713) -> Result<HashMap<[u8; 16], HashMap<String, IrLiteral>>, GfError> {
1714    let path = node_props_path(dir, stem);
1715    if !path.try_exists().map_err(|error| io_err(&error))? {
1716        return Ok(HashMap::new());
1717    }
1718    let file = fs::File::open(&path).map_err(|error| io_err(&error))?;
1719    let schema = ParquetRecordBatchReaderBuilder::try_new(file)
1720        .map_err(pq_err)?
1721        .schema()
1722        .clone();
1723    let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1724    Ok(decode_property_rows(&batches)?
1725        .into_iter()
1726        .map(|row| (row.node_uuid, row.props))
1727        .collect())
1728}
1729
1730/// Count non-null properties owned by the selected persisted entities across
1731/// every dynamic-schema property partition.
1732pub fn count_entity_properties<S: std::hash::BuildHasher>(
1733    dir: &Path,
1734    targets: &HashSet<[u8; 16], S>,
1735    is_edge: bool,
1736) -> Result<u64, GfError> {
1737    if targets.is_empty() {
1738        return Ok(0);
1739    }
1740    let property_dir = dir.join(if is_edge {
1741        "edge_properties"
1742    } else {
1743        "properties"
1744    });
1745    let Ok(entries) = std::fs::read_dir(property_dir) else {
1746        return Ok(0);
1747    };
1748    let mut count = 0u64;
1749    for entry in entries {
1750        let path = entry
1751            .map_err(|error| GfError::Storage(error.to_string()))?
1752            .path();
1753        if path
1754            .extension()
1755            .is_none_or(|extension| extension != "parquet")
1756        {
1757            continue;
1758        }
1759        let Some(schema) = crate::catalog::discover_parquet_schema(&path) else {
1760            continue;
1761        };
1762        let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(pq_err)?;
1763        if is_edge {
1764            for row in decode_edge_property_rows(&batches)? {
1765                if targets.contains(&row.edge_uuid) {
1766                    count += row.props.len() as u64;
1767                }
1768            }
1769        } else {
1770            for row in decode_property_rows(&batches)? {
1771                if targets.contains(&row.node_uuid) {
1772                    count += row.props.len() as u64;
1773                }
1774            }
1775        }
1776    }
1777    Ok(count)
1778}
1779
1780/// Decode one property batch row-by-row, invoking `emit(uuid, props)` per row.
1781///
1782/// `uuid_field_name` is the join-key column (`node_uuid` / `edge_uuid`); every
1783/// other column maps by its Arrow type to an [`IrLiteral`]. A **null** slot
1784/// omits the key for that row — it must never become a concrete value, or it
1785/// would wrongly pin a previously-all-null column's type on re-inference.
1786#[allow(
1787    clippy::too_many_lines,
1788    reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
1789)]
1790fn decode_property_batch(
1791    batch: &RecordBatch,
1792    uuid_field_name: &str,
1793    mut emit: impl FnMut([u8; 16], HashMap<String, IrLiteral>),
1794) -> Result<(), GfError> {
1795    use arrow::array::Array;
1796
1797    let schema = batch.schema();
1798    let uuid_col = batch
1799        .column_by_name(uuid_field_name)
1800        .and_then(|c| c.as_any().downcast_ref::<FixedSizeBinaryArray>())
1801        .ok_or_else(|| {
1802            GfError::Storage(format!("property file missing {uuid_field_name} column"))
1803        })?;
1804    for r in 0..batch.num_rows() {
1805        let mut uuid = [0u8; 16];
1806        uuid.copy_from_slice(uuid_col.value(r));
1807        let mut props: HashMap<String, IrLiteral> = HashMap::new();
1808        for (c, field) in schema.fields().iter().enumerate() {
1809            if field.name() == uuid_field_name {
1810                continue;
1811            }
1812            let col = batch.column(c);
1813            if col.is_null(r) {
1814                continue; // null slot: omit the key (never fabricate a value)
1815            }
1816            let lit = decode_value(col, field, r)?;
1817            props.insert(field.name().clone(), lit);
1818        }
1819        emit(uuid, props);
1820    }
1821    Ok(())
1822}
1823
1824/// Decode one property value at row `r` from its column, dispatched by Arrow
1825/// type (structs by field NAMES). A `List` recurses element-wise on the inner
1826/// field's type, so list-of-temporals reuse the struct dispatch (#1006). Shared
1827/// by [`decode_property_batch`] and its own list arm.
1828#[allow(
1829    clippy::too_many_lines,
1830    reason = "one decode arm per Arrow type, plus per-shape struct dispatch; clearest inline"
1831)]
1832fn decode_value(
1833    col: &arrow::array::ArrayRef,
1834    field: &arrow::datatypes::Field,
1835    r: usize,
1836) -> Result<IrLiteral, GfError> {
1837    use arrow::array::{
1838        Array, BooleanArray, Int32Array, Int64Array, ListArray, StructArray, Time64NanosecondArray,
1839    };
1840    Ok(match field.data_type() {
1841        DataType::Int64 => IrLiteral::Int(downcast::<Int64Array>(col, field)?.value(r)),
1842        DataType::Float64 => IrLiteral::Float(downcast::<Float64Array>(col, field)?.value(r)),
1843        DataType::Boolean => IrLiteral::Bool(downcast::<BooleanArray>(col, field)?.value(r)),
1844        DataType::Utf8 => IrLiteral::Str(downcast::<StringArray>(col, field)?.value(r).to_owned()),
1845        // A typed temporal struct (#920). Dispatch by the struct's field
1846        // NAMES — every persisted Struct used to be assumed a duration,
1847        // but `localdatetime` (and later `time`/`datetime`) are also
1848        // structs, so the shape must select the decode.
1849        DataType::Struct(fields) => {
1850            let s = downcast::<StructArray>(col, field)?;
1851            let names: Vec<&str> = fields.iter().map(|f| f.name().as_str()).collect();
1852            match names.as_slice() {
1853                [
1854                    "__het_tag",
1855                    "__het_int",
1856                    "__het_float",
1857                    "__het_str",
1858                    "__het_bool",
1859                ] => {
1860                    let tag = s
1861                        .column(0)
1862                        .as_any()
1863                        .downcast_ref::<arrow::array::Int8Array>()
1864                        .ok_or_else(|| GfError::Storage("heterogeneous tag not Int8".into()))?
1865                        .value(r);
1866                    match tag {
1867                        0 => IrLiteral::Int(
1868                            s.column(1)
1869                                .as_any()
1870                                .downcast_ref::<Int64Array>()
1871                                .ok_or_else(|| {
1872                                    GfError::Storage("heterogeneous int not Int64".into())
1873                                })?
1874                                .value(r),
1875                        ),
1876                        1 => IrLiteral::Float(
1877                            s.column(2)
1878                                .as_any()
1879                                .downcast_ref::<Float64Array>()
1880                                .ok_or_else(|| {
1881                                    GfError::Storage("heterogeneous float not Float64".into())
1882                                })?
1883                                .value(r),
1884                        ),
1885                        2 => IrLiteral::Str(
1886                            s.column(3)
1887                                .as_any()
1888                                .downcast_ref::<StringArray>()
1889                                .ok_or_else(|| {
1890                                    GfError::Storage("heterogeneous string not Utf8".into())
1891                                })?
1892                                .value(r)
1893                                .to_owned(),
1894                        ),
1895                        3 => IrLiteral::Bool(
1896                            s.column(4)
1897                                .as_any()
1898                                .downcast_ref::<BooleanArray>()
1899                                .ok_or_else(|| {
1900                                    GfError::Storage("heterogeneous bool not Boolean".into())
1901                                })?
1902                                .value(r),
1903                        ),
1904                        _ => {
1905                            return Err(GfError::Storage(format!(
1906                                "unsupported heterogeneous property tag {tag}"
1907                            )));
1908                        }
1909                    }
1910                }
1911                ["months", "days", "seconds", "nanos"] => {
1912                    let i64_at = |idx: usize| -> Result<i64, GfError> {
1913                        Ok(s.column(idx)
1914                            .as_any()
1915                            .downcast_ref::<Int64Array>()
1916                            .ok_or_else(|| {
1917                                GfError::Storage("duration struct child not Int64".into())
1918                            })?
1919                            .value(r))
1920                    };
1921                    IrLiteral::Duration {
1922                        months: i64_at(0)?,
1923                        days: i64_at(1)?,
1924                        seconds: i64_at(2)?,
1925                        nanos: i64_at(3)?,
1926                    }
1927                }
1928                ["epoch_day"] => IrLiteral::Date(
1929                    s.column(0)
1930                        .as_any()
1931                        .downcast_ref::<Int64Array>()
1932                        .ok_or_else(|| GfError::Storage("date epoch_day not Int64".into()))?
1933                        .value(r),
1934                ),
1935                ["date", "time"] => {
1936                    let days = s
1937                        .column(0)
1938                        .as_any()
1939                        .downcast_ref::<Int64Array>()
1940                        .ok_or_else(|| GfError::Storage("localdatetime date not Int64".into()))?
1941                        .value(r);
1942                    let nanos = s
1943                        .column(1)
1944                        .as_any()
1945                        .downcast_ref::<Time64NanosecondArray>()
1946                        .ok_or_else(|| {
1947                            GfError::Storage("localdatetime time not Time64(ns)".into())
1948                        })?
1949                        .value(r);
1950                    IrLiteral::LocalDateTime { days, nanos }
1951                }
1952                ["time", "offset"] => {
1953                    let nanos = s
1954                        .column(0)
1955                        .as_any()
1956                        .downcast_ref::<Time64NanosecondArray>()
1957                        .ok_or_else(|| GfError::Storage("time not Time64(ns)".into()))?
1958                        .value(r);
1959                    let offset = s
1960                        .column(1)
1961                        .as_any()
1962                        .downcast_ref::<Int32Array>()
1963                        .ok_or_else(|| GfError::Storage("time offset not Int32".into()))?
1964                        .value(r);
1965                    IrLiteral::ZonedTime { nanos, offset }
1966                }
1967                ["date", "time", "offset", "zone"] => {
1968                    let days = s
1969                        .column(0)
1970                        .as_any()
1971                        .downcast_ref::<Int64Array>()
1972                        .ok_or_else(|| GfError::Storage("datetime date not Int64".into()))?
1973                        .value(r);
1974                    let nanos = s
1975                        .column(1)
1976                        .as_any()
1977                        .downcast_ref::<Time64NanosecondArray>()
1978                        .ok_or_else(|| GfError::Storage("datetime time not Time64(ns)".into()))?
1979                        .value(r);
1980                    let offset = s
1981                        .column(2)
1982                        .as_any()
1983                        .downcast_ref::<Int32Array>()
1984                        .ok_or_else(|| GfError::Storage("datetime offset not Int32".into()))?
1985                        .value(r);
1986                    let zone_col = s
1987                        .column(3)
1988                        .as_any()
1989                        .downcast_ref::<StringArray>()
1990                        .ok_or_else(|| GfError::Storage("datetime zone not Utf8".into()))?;
1991                    // A NULL zone child means offset-only (no named zone).
1992                    let zone = (!zone_col.is_null(r)).then(|| zone_col.value(r).to_owned());
1993                    IrLiteral::ZonedDateTime {
1994                        days,
1995                        nanos,
1996                        offset,
1997                        zone,
1998                    }
1999                }
2000                _ => {
2001                    return Err(GfError::Storage(format!(
2002                        "property column {} has unsupported struct shape {names:?}",
2003                        field.name()
2004                    )));
2005                }
2006            }
2007        }
2008        DataType::Time64(TimeUnit::Nanosecond) => {
2009            IrLiteral::Time(downcast::<Time64NanosecondArray>(col, field)?.value(r))
2010        }
2011        DataType::Timestamp(TimeUnit::Microsecond, _) => {
2012            IrLiteral::DateTime(downcast::<TimestampMicrosecondArray>(col, field)?.value(r))
2013        }
2014        // A homogeneous list (#1006): decode each element by recursing on
2015        // the inner field's type (so list-of-temporals reuse the struct
2016        // dispatch). Nested lists recurse naturally.
2017        DataType::List(inner) => {
2018            let larr = downcast::<ListArray>(col, field)?;
2019            let elems = larr.value(r);
2020            let mut items = Vec::with_capacity(elems.len());
2021            for j in 0..elems.len() {
2022                if elems.is_null(j) {
2023                    items.push(IrLiteral::Null);
2024                } else {
2025                    items.push(decode_value(&elems, inner, j)?);
2026                }
2027            }
2028            IrLiteral::List(items)
2029        }
2030        other => {
2031            return Err(GfError::Storage(format!(
2032                "property column {} has unsupported type {other:?}",
2033                field.name()
2034            )));
2035        }
2036    })
2037}
2038
2039/// Downcast a column to a concrete Arrow array type, erroring with the column
2040/// name if the dynamic type does not match its declared field type.
2041fn downcast<'a, A: 'static>(
2042    col: &'a arrow::array::ArrayRef,
2043    field: &arrow::datatypes::Field,
2044) -> Result<&'a A, GfError> {
2045    col.as_any().downcast_ref::<A>().ok_or_else(|| {
2046        GfError::Storage(format!(
2047            "property column {} could not be read as its declared type",
2048            field.name()
2049        ))
2050    })
2051}
2052
2053// ---------------------------------------------------------------------------
2054// SET / REMOVE property rewrite primitives (#791)
2055// ---------------------------------------------------------------------------
2056//
2057// These rewrite committed `properties/<stem>.parquet` / `edge_properties/
2058// <stem>.parquet` files, mirroring the writer's decode → mutate → re-infer →
2059// write cycle (see [`GraphWriter::flush_properties`]). The `stage_*` forms
2060// stage into a caller-owned [`RewriteBatch`] so one statement's rewrites
2061// across stems commit all-or-nothing (#790); the original four functions are
2062// stage-and-commit wrappers for single-stem callers.
2063//
2064// The execution layer accumulates per-uuid updates/removals from the matched
2065// rows, then calls these once per file stem. SET **merges** into a uuid's
2066// existing property map (overwriting same-named keys) and **inserts** a fresh
2067// row for a uuid that had no property row yet; REMOVE drops keys (a column that
2068// becomes all-absent disappears on re-inference). Both return the number of
2069// distinct entities whose file row was written.
2070
2071/// Apply per-uuid property `updates` (SET) to the rows decoded from `existing`,
2072/// merging into each uuid's map and inserting a row for any uuid not present.
2073/// Returns the rebuilt row set and the number of distinct uuids touched.
2074fn apply_property_updates<R: PropRowLike>(
2075    mut rows: Vec<R>,
2076    updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2077) -> (Vec<R>, u64) {
2078    // uuid → index into `rows` (first occurrence wins; the decode never emits
2079    // duplicate uuids, but a defensive first-wins keeps this total).
2080    let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
2081    for (i, row) in rows.iter().enumerate() {
2082        index.entry(*row.uuid_bytes()).or_insert(i);
2083    }
2084    let mut touched = 0u64;
2085    for (uuid, new_props) in updates {
2086        if new_props.is_empty() {
2087            continue;
2088        }
2089        touched += 1;
2090        if let Some(&i) = index.get(uuid) {
2091            rows[i].props_mut().extend(new_props.clone());
2092        } else {
2093            index.insert(*uuid, rows.len());
2094            rows.push(R::from_parts(*uuid, new_props.clone()));
2095        }
2096    }
2097    (rows, touched)
2098}
2099
2100/// Apply per-uuid property `removals` (REMOVE) to the rows decoded from
2101/// `existing`. Removing an absent key or an absent uuid is a no-op (openCypher).
2102/// Returns the rebuilt row set and the number of distinct uuids touched (a uuid
2103/// is counted even if every named key was already absent — the entity was
2104/// targeted).
2105fn apply_property_removals<R: PropRowLike>(
2106    mut rows: Vec<R>,
2107    removals: &HashMap<[u8; 16], HashSet<String>>,
2108) -> (Vec<R>, u64) {
2109    let mut index: HashMap<[u8; 16], usize> = HashMap::with_capacity(rows.len());
2110    for (i, row) in rows.iter().enumerate() {
2111        index.entry(*row.uuid_bytes()).or_insert(i);
2112    }
2113    let mut touched = 0u64;
2114    for (uuid, keys) in removals {
2115        if keys.is_empty() {
2116            continue;
2117        }
2118        touched += 1;
2119        if let Some(&i) = index.get(uuid) {
2120            let props = rows[i].props_mut();
2121            for k in keys {
2122                props.remove(k);
2123            }
2124        }
2125    }
2126    (rows, touched)
2127}
2128
2129/// Stage a rewrite of `properties/<stem>.parquet` applying per-`node_uuid` SET
2130/// `updates` into `staged` (committed by the caller, #790).
2131///
2132/// Reads the current file (decode → re-infer over the merged set), inserts a
2133/// row for any node that had no property row yet, and stages the rebuilt file.
2134/// A write is skipped only when the merged set is empty. Returns the number of
2135/// distinct nodes whose row was set.
2136///
2137/// # Errors
2138/// Propagates Parquet / Arrow / IO errors from reading or staging the file.
2139// The execution layer always accumulates updates with the default hasher;
2140// generalizing the nested maps over `BuildHasher` would add two type params per
2141// fn for no caller benefit.
2142#[allow(clippy::implicit_hasher)]
2143pub fn stage_set_node_properties(
2144    staged: &mut RewriteBatch,
2145    dir: &Path,
2146    stem: &str,
2147    updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2148) -> Result<u64, GfError> {
2149    let existing = read_props_through(staged, &node_props_path(dir, stem))?;
2150    let rows = decode_property_rows(&existing)?;
2151    let (rows, touched) = apply_property_updates(rows, updates);
2152    stage_node_property_file(staged, dir, stem, &rows)?;
2153    Ok(touched)
2154}
2155
2156/// Stage a rewrite of `properties/<stem>.parquet` applying per-`node_uuid`
2157/// REMOVE `removals`. Returns the number of distinct nodes targeted.
2158///
2159/// # Errors
2160/// Propagates Parquet / Arrow / IO errors from reading or staging the file.
2161#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2162pub fn stage_remove_node_properties(
2163    staged: &mut RewriteBatch,
2164    dir: &Path,
2165    stem: &str,
2166    removals: &HashMap<[u8; 16], HashSet<String>>,
2167) -> Result<u64, GfError> {
2168    let existing = read_props_through(staged, &node_props_path(dir, stem))?;
2169    let rows = decode_property_rows(&existing)?;
2170    let (rows, touched) = apply_property_removals(rows, removals);
2171    stage_node_property_file(staged, dir, stem, &rows)?;
2172    Ok(touched)
2173}
2174
2175/// Stage a rewrite of `edge_properties/<rel_stem>.parquet` applying
2176/// per-`edge_uuid` SET `updates`. Edge analogue of
2177/// [`stage_set_node_properties`]; the join key is `edge_uuid` and the file is
2178/// routed by relation name.
2179///
2180/// # Errors
2181/// Propagates Parquet / Arrow / IO errors from reading or staging the file.
2182#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2183pub fn stage_set_edge_properties(
2184    staged: &mut RewriteBatch,
2185    dir: &Path,
2186    rel_stem: &str,
2187    updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2188) -> Result<u64, GfError> {
2189    let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
2190    let rows = decode_edge_property_rows(&existing)?;
2191    let (rows, touched) = apply_property_updates(rows, updates);
2192    stage_edge_property_file(staged, dir, rel_stem, &rows)?;
2193    Ok(touched)
2194}
2195
2196/// Stage a rewrite of `edge_properties/<rel_stem>.parquet` applying
2197/// per-`edge_uuid` REMOVE `removals`. Edge analogue of
2198/// [`stage_remove_node_properties`].
2199///
2200/// # Errors
2201/// Propagates Parquet / Arrow / IO errors from reading or staging the file.
2202#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2203pub fn stage_remove_edge_properties(
2204    staged: &mut RewriteBatch,
2205    dir: &Path,
2206    rel_stem: &str,
2207    removals: &HashMap<[u8; 16], HashSet<String>>,
2208) -> Result<u64, GfError> {
2209    let existing = read_props_through(staged, &edge_props_path(dir, rel_stem))?;
2210    let rows = decode_edge_property_rows(&existing)?;
2211    let (rows, touched) = apply_property_removals(rows, removals);
2212    stage_edge_property_file(staged, dir, rel_stem, &rows)?;
2213    Ok(touched)
2214}
2215
2216/// Rewrite `properties/<stem>.parquet` applying per-`node_uuid` SET `updates`,
2217/// staged and committed as one batch (#790).
2218///
2219/// # Errors
2220/// Propagates Parquet / Arrow / IO errors from reading or rewriting the file.
2221#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2222pub fn set_node_properties(
2223    dir: &Path,
2224    stem: &str,
2225    updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2226) -> Result<u64, GfError> {
2227    let mut staged = RewriteBatch::new();
2228    let touched = stage_set_node_properties(&mut staged, dir, stem, updates)?;
2229    crate::generation::commit_topology_aware(staged, dir)?;
2230    Ok(touched)
2231}
2232
2233/// Rewrite `properties/<stem>.parquet` applying per-`node_uuid` REMOVE
2234/// `removals`, staged and committed as one batch (#790).
2235///
2236/// # Errors
2237/// Propagates Parquet / Arrow / IO errors from reading or rewriting the file.
2238#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2239pub fn remove_node_properties(
2240    dir: &Path,
2241    stem: &str,
2242    removals: &HashMap<[u8; 16], HashSet<String>>,
2243) -> Result<u64, GfError> {
2244    let mut staged = RewriteBatch::new();
2245    let touched = stage_remove_node_properties(&mut staged, dir, stem, removals)?;
2246    crate::generation::commit_topology_aware(staged, dir)?;
2247    Ok(touched)
2248}
2249
2250/// Rewrite `edge_properties/<rel_stem>.parquet` applying per-`edge_uuid` SET
2251/// `updates`, staged and committed as one batch (#790).
2252///
2253/// # Errors
2254/// Propagates Parquet / Arrow / IO errors from reading or rewriting the file.
2255#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2256pub fn set_edge_properties_rewrite(
2257    dir: &Path,
2258    rel_stem: &str,
2259    updates: &HashMap<[u8; 16], HashMap<String, IrLiteral>>,
2260) -> Result<u64, GfError> {
2261    let mut staged = RewriteBatch::new();
2262    let touched = stage_set_edge_properties(&mut staged, dir, rel_stem, updates)?;
2263    staged.commit()?;
2264    Ok(touched)
2265}
2266
2267/// Rewrite `edge_properties/<rel_stem>.parquet` applying per-`edge_uuid`
2268/// REMOVE `removals`, staged and committed as one batch (#790).
2269///
2270/// # Errors
2271/// Propagates Parquet / Arrow / IO errors from reading or rewriting the file.
2272#[allow(clippy::implicit_hasher)] // see `stage_set_node_properties`
2273pub fn remove_edge_properties(
2274    dir: &Path,
2275    rel_stem: &str,
2276    removals: &HashMap<[u8; 16], HashSet<String>>,
2277) -> Result<u64, GfError> {
2278    let mut staged = RewriteBatch::new();
2279    let touched = stage_remove_edge_properties(&mut staged, dir, rel_stem, removals)?;
2280    staged.commit()?;
2281    Ok(touched)
2282}
2283
2284/// Rebuild `properties/<stem>.parquet` from `rows` and stage it. A write is
2285/// skipped only when `rows` is empty (an absent file with no inserts) — the
2286/// dynamic-schema builder cannot emit a zero-row key column, and there is
2287/// nothing to persist. A REMOVE that empties a row's last property keeps the
2288/// row (now with no property columns); its property map decodes back as empty.
2289fn stage_node_property_file(
2290    staged: &mut RewriteBatch,
2291    dir: &Path,
2292    stem: &str,
2293    rows: &[PropRow],
2294) -> Result<(), GfError> {
2295    if rows.is_empty() {
2296        return Ok(());
2297    }
2298    let (schema, cols) = build_property_columns(stem, rows)?;
2299    stage_property_file(staged, dir, "properties", stem, schema, cols)
2300}
2301
2302/// Edge analogue of [`stage_node_property_file`] (key `edge_uuid`, file routed
2303/// by relation name under `edge_properties/`).
2304fn stage_edge_property_file(
2305    staged: &mut RewriteBatch,
2306    dir: &Path,
2307    stem: &str,
2308    rows: &[EdgePropRow],
2309) -> Result<(), GfError> {
2310    if rows.is_empty() {
2311        return Ok(());
2312    }
2313    let (schema, cols) =
2314        build_property_columns_keyed(EDGE_PROPERTY_UUID_FIELD, "graphforge.rel_type", stem, rows)?;
2315    stage_property_file(staged, dir, "edge_properties", stem, schema, cols)
2316}
2317
2318/// Stage a rebuilt property file under `<dir>/<subdir>/<stem>.parquet` (the
2319/// staging core creates the subdirectory), replacing any content this
2320/// statement already staged for it. Callers guard the empty-row case before
2321/// reaching here.
2322fn stage_property_file(
2323    staged: &mut RewriteBatch,
2324    dir: &Path,
2325    subdir: &str,
2326    stem: &str,
2327    schema: Schema,
2328    cols: Vec<ArrayRef>,
2329) -> Result<(), GfError> {
2330    let batch = RecordBatch::try_new(Arc::new(schema), cols).map_err(pq_err)?;
2331    staged.restage(
2332        &dir.join(subdir).join(format!("{stem}.parquet")),
2333        batch.schema(),
2334        &batch,
2335    )
2336}
2337
2338/// `properties/<stem>.parquet` under `dir`.
2339fn node_props_path(dir: &Path, stem: &str) -> PathBuf {
2340    dir.join("properties").join(format!("{stem}.parquet"))
2341}
2342
2343/// `edge_properties/<stem>.parquet` under `dir`.
2344fn edge_props_path(dir: &Path, stem: &str) -> PathBuf {
2345    dir.join("edge_properties").join(format!("{stem}.parquet"))
2346}
2347
2348/// Read a dynamic-schema property file's batches for staging, **through**
2349/// `staged` (#792): content already staged for `path` in this statement is
2350/// the base, else the committed file; an absent file reads as empty.
2351fn read_props_through(staged: &RewriteBatch, path: &Path) -> Result<Vec<RecordBatch>, GfError> {
2352    let read_path = staged
2353        .staged_temp(path)
2354        .map_or_else(|| path.to_path_buf(), Path::to_path_buf);
2355    match crate::catalog::discover_parquet_schema(&read_path) {
2356        Some(schema) => crate::catalog::read_parquet_or_empty(&read_path, schema).map_err(pq_err),
2357        None => Ok(Vec::new()),
2358    }
2359}
2360
2361// ---------------------------------------------------------------------------
2362// Parquet write helper
2363// ---------------------------------------------------------------------------
2364
2365use crate::staging::RewriteBatch;
2366
2367/// Concatenate already-on-disk `existing` batches with the freshly-built `new`
2368/// batch under a fixed `schema` (used for the node + typed/exploratory edge
2369/// files, whose schema never varies). Zero-row placeholder batches from an
2370/// absent file concat harmlessly.
2371fn concat_with_existing(
2372    schema: &SchemaRef,
2373    existing: Vec<RecordBatch>,
2374    new: RecordBatch,
2375) -> Result<RecordBatch, GfError> {
2376    let mut all = existing;
2377    all.push(new);
2378    arrow::compute::concat_batches(schema, &all).map_err(pq_err)
2379}
2380
2381// ---------------------------------------------------------------------------
2382// Tests
2383// ---------------------------------------------------------------------------
2384
2385#[cfg(test)]
2386mod tests {
2387    use std::fs::File;
2388
2389    use super::*;
2390    use graphforge_core::uuid::new_v7;
2391    use tempfile::TempDir;
2392
2393    const TS: i64 = 1_700_000_000_000_000;
2394
2395    #[test]
2396    fn create_node_persists_complete_label_set_and_primary_label() {
2397        let dir = TempDir::new().unwrap();
2398        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2399        w.create_node_with_labels(new_v7(), &[TypeId(4), TypeId(9)])
2400            .unwrap();
2401        w.flush().unwrap();
2402
2403        let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2404        let batch = &nodes[0];
2405        let primary = batch
2406            .column_by_name("type_id")
2407            .unwrap()
2408            .as_any()
2409            .downcast_ref::<UInt32Array>()
2410            .unwrap();
2411        assert_eq!(primary.value(0), 4);
2412        let sets = batch
2413            .column_by_name("type_ids")
2414            .unwrap()
2415            .as_any()
2416            .downcast_ref::<arrow::array::ListArray>()
2417            .unwrap();
2418        let labels = sets.value(0);
2419        let labels = labels.as_any().downcast_ref::<UInt32Array>().unwrap();
2420        assert_eq!(labels.values(), &[4, 9]);
2421    }
2422
2423    #[test]
2424    fn surrogate_ids_are_monotonic_from_one() {
2425        let dir = TempDir::new().unwrap();
2426        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2427        assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 1);
2428        assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 2);
2429        assert_eq!(w.create_node(new_v7(), TypeId(0)).unwrap(), 3);
2430    }
2431
2432    #[test]
2433    fn create_edge_with_unknown_endpoint_errors() {
2434        let dir = TempDir::new().unwrap();
2435        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2436        let a = new_v7();
2437        w.create_node(a, TypeId(0)).unwrap();
2438        // `b` was never created.
2439        let b = new_v7();
2440        let e = w.create_edge(new_v7(), "KNOWS", &a, &b);
2441        assert!(matches!(e, Err(GfError::Storage(_))), "got {e:?}");
2442
2443        let unknown_source = new_v7();
2444        let source_error = w.create_edge(new_v7(), "KNOWS", &unknown_source, &a);
2445        assert!(matches!(&source_error, Err(GfError::Storage(_))));
2446        assert!(source_error.unwrap_err().to_string().contains("source"));
2447    }
2448
2449    #[test]
2450    fn register_existing_node_resolves_edge_endpoint() {
2451        // #703: a MATCH-bound node is referenced (not created). Registering its
2452        // identity lets a subsequent edge resolve it without writing a node row
2453        // or advancing the surrogate counter.
2454        let dir = TempDir::new().unwrap();
2455        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2456        // `a` already exists on disk with surrogate 42 (e.g. from a prior write).
2457        let a = new_v7();
2458        w.register_existing_node(a, 42);
2459        // A freshly-minted `b` (next surrogate is 1 — register did not advance it).
2460        let b = new_v7();
2461        assert_eq!(w.create_node(b, TypeId(0)).unwrap(), 1);
2462        // The edge resolves both endpoints; its src_id is the registered 42.
2463        let edge_id = w
2464            .create_edge(new_v7(), "KNOWS", &a, &b)
2465            .expect("edge with a registered endpoint resolves");
2466        assert_eq!(edge_id, 1);
2467    }
2468
2469    #[test]
2470    fn register_existing_node_does_not_write_a_node_row() {
2471        // Registering an existing node must not buffer a NodeRow — only the one
2472        // genuinely-created node is flushed.
2473        let dir = TempDir::new().unwrap();
2474        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2475        w.register_existing_node(new_v7(), 7);
2476        let created = new_v7();
2477        w.create_node(created, TypeId(0)).unwrap();
2478        w.flush().unwrap();
2479        // Exactly one node on disk (the created one), not two.
2480        let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2481        let rows: usize = nodes.iter().map(arrow::array::RecordBatch::num_rows).sum();
2482        assert_eq!(
2483            rows, 1,
2484            "register_existing_node must not persist a node row"
2485        );
2486    }
2487
2488    #[test]
2489    fn empty_flush_creates_no_directories() {
2490        let dir = TempDir::new().unwrap();
2491        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2492        w.flush().unwrap();
2493        assert!(!dir.path().join("topology").exists());
2494        assert!(!dir.path().join("properties").exists());
2495    }
2496
2497    #[test]
2498    fn null_first_property_column_infers_later_concrete_type() {
2499        use arrow::array::Array;
2500        use arrow::datatypes::DataType;
2501        use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
2502
2503        let dir = TempDir::new().unwrap();
2504        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2505        let a = new_v7();
2506        let b = new_v7();
2507        let c = new_v7();
2508        w.create_node(a, TypeId(0)).unwrap();
2509        w.create_node(b, TypeId(0)).unwrap();
2510        w.create_node(c, TypeId(0)).unwrap();
2511        // First row: `score` is Null. Later rows: consistently Int.
2512        w.set_properties(
2513            &a,
2514            None,
2515            HashMap::from([("score".to_owned(), IrLiteral::Null)]),
2516        )
2517        .unwrap();
2518        w.set_properties(
2519            &b,
2520            None,
2521            HashMap::from([("score".to_owned(), IrLiteral::Int(10))]),
2522        )
2523        .unwrap();
2524        w.set_properties(
2525            &c,
2526            None,
2527            HashMap::from([("score".to_owned(), IrLiteral::Int(20))]),
2528        )
2529        .unwrap();
2530        w.flush().unwrap();
2531
2532        let path = dir.path().join("properties").join("_untyped.parquet");
2533        let file = File::open(&path).unwrap();
2534        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2535        let schema = builder.schema().clone();
2536        // Exactly one `score` column (no duplicate from the null row), typed Int64.
2537        let score_fields: Vec<_> = schema
2538            .fields()
2539            .iter()
2540            .filter(|f| f.name() == "score")
2541            .collect();
2542        assert_eq!(score_fields.len(), 1, "expected a single score column");
2543        assert_eq!(
2544            score_fields[0].data_type(),
2545            &DataType::Int64,
2546            "null-first then Int should infer Int64, not Utf8"
2547        );
2548
2549        // Round-trip: 3 rows, first null then 10, 20.
2550        let mut reader = builder.build().unwrap();
2551        let batch = reader.next().unwrap().unwrap();
2552        assert_eq!(batch.num_rows(), 3);
2553        let scores = batch
2554            .column(schema.index_of("score").unwrap())
2555            .as_any()
2556            .downcast_ref::<arrow::array::Int64Array>()
2557            .unwrap();
2558        assert!(scores.is_null(0));
2559        assert_eq!(scores.value(1), 10);
2560        assert_eq!(scores.value(2), 20);
2561    }
2562
2563    #[test]
2564    fn mixed_type_property_column_uses_tagged_scalars() {
2565        use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
2566
2567        let dir = TempDir::new().unwrap();
2568        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2569        let a = new_v7();
2570        let b = new_v7();
2571        w.create_node(a, TypeId(0)).unwrap();
2572        w.create_node(b, TypeId(0)).unwrap();
2573        w.set_properties(
2574            &a,
2575            None,
2576            HashMap::from([("x".to_owned(), IrLiteral::Int(1))]),
2577        )
2578        .unwrap();
2579        w.set_properties(
2580            &b,
2581            None,
2582            HashMap::from([("x".to_owned(), IrLiteral::Str("two".to_owned()))]),
2583        )
2584        .unwrap();
2585        w.flush().unwrap();
2586
2587        let path = dir.path().join("properties").join("_untyped.parquet");
2588        let file = File::open(&path).unwrap();
2589        let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2590        let schema = builder.schema().clone();
2591        let x = schema.field_with_name("x").unwrap();
2592        assert_eq!(
2593            x.data_type(),
2594            &DataType::Struct(heterogeneous_scalar_fields())
2595        );
2596    }
2597
2598    #[test]
2599    fn every_heterogeneous_scalar_tag_round_trips_exactly() {
2600        let dir = TempDir::new().unwrap();
2601        let cases = [
2602            IrLiteral::Int(-1),
2603            IrLiteral::Float(2.25),
2604            IrLiteral::Str("three".into()),
2605            IrLiteral::Bool(true),
2606        ];
2607        let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2608        let mut expected = HashMap::new();
2609        for value in cases {
2610            let node = new_v7();
2611            writer.create_node(node, TypeId(0)).unwrap();
2612            writer
2613                .set_properties(
2614                    &node,
2615                    None,
2616                    HashMap::from([("mixed".into(), value.clone())]),
2617                )
2618                .unwrap();
2619            expected.insert(to_bytes(&node), value);
2620        }
2621        writer.flush().unwrap();
2622
2623        let reopened = read_node_props(dir.path(), "_untyped");
2624        assert_eq!(reopened.len(), expected.len());
2625        for (node, value) in expected {
2626            assert_eq!(reopened[&node].get("mixed"), Some(&value));
2627        }
2628    }
2629
2630    // -----------------------------------------------------------------------
2631    // SET / REMOVE rewrite primitives (#791)
2632    // -----------------------------------------------------------------------
2633
2634    /// Read a node-property file back into a `uuid → props` map (mirrors the
2635    /// decode the rewrite primitives use), for assertions.
2636    fn read_node_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
2637        read_node_property_rows(dir, stem).unwrap()
2638    }
2639
2640    fn read_edge_props(dir: &Path, stem: &str) -> HashMap<[u8; 16], HashMap<String, IrLiteral>> {
2641        let batches = crate::catalog::read_edge_properties(dir, stem).unwrap();
2642        let mut out = HashMap::new();
2643        for row in decode_edge_property_rows(&batches).unwrap() {
2644            out.insert(row.edge_uuid, row.props);
2645        }
2646        out
2647    }
2648
2649    #[test]
2650    fn set_node_properties_sets_new_and_overwrites_existing() {
2651        let dir = TempDir::new().unwrap();
2652        assert!(
2653            read_node_property_rows(dir.path(), "_untyped")
2654                .unwrap()
2655                .is_empty()
2656        );
2657        fs::create_dir_all(dir.path().join("properties")).unwrap();
2658        fs::write(dir.path().join("properties/_untyped.parquet"), b"invalid").unwrap();
2659        assert!(read_node_property_rows(dir.path(), "_untyped").is_err());
2660        fs::remove_file(dir.path().join("properties/_untyped.parquet")).unwrap();
2661        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2662        let a = new_v7();
2663        w.create_node(a, TypeId(0)).unwrap();
2664        w.set_properties(
2665            &a,
2666            None,
2667            HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2668        )
2669        .unwrap();
2670        w.flush().unwrap();
2671
2672        let ab = to_bytes(&a);
2673        let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2674        // Overwrite `age` and add a new `name`.
2675        let updates = HashMap::from([(
2676            ab,
2677            HashMap::from([
2678                ("age".to_owned(), IrLiteral::Int(31)),
2679                ("name".to_owned(), IrLiteral::Str("Al".to_owned())),
2680            ]),
2681        )]);
2682        let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
2683        assert_eq!(touched, 1);
2684        assert_eq!(
2685            crate::generation::read_search_generation(dir.path()).unwrap(),
2686            search_generation + 1
2687        );
2688
2689        let props = read_node_props(dir.path(), "_untyped");
2690        assert_eq!(props[&ab]["age"], IrLiteral::Int(31));
2691        assert_eq!(props[&ab]["name"], IrLiteral::Str("Al".to_owned()));
2692    }
2693
2694    #[test]
2695    fn set_node_properties_inserts_row_for_propertyless_node() {
2696        // A node with no property row yet must get a fresh row on SET.
2697        let dir = TempDir::new().unwrap();
2698        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2699        let a = new_v7();
2700        w.create_node(a, TypeId(0)).unwrap();
2701        w.flush().unwrap(); // no properties written → no _untyped file
2702
2703        let ab = to_bytes(&a);
2704        let updates =
2705            HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(42))]))]);
2706        let touched = set_node_properties(dir.path(), "_untyped", &updates).unwrap();
2707        assert_eq!(touched, 1);
2708
2709        let props = read_node_props(dir.path(), "_untyped");
2710        assert_eq!(props[&ab]["age"], IrLiteral::Int(42));
2711    }
2712
2713    #[test]
2714    fn set_node_properties_routes_by_stem_in_strict_mode() {
2715        // Strict/Advisory route to properties/<Entity>.parquet, not _untyped.
2716        let dir = TempDir::new().unwrap();
2717        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
2718        let a = new_v7();
2719        w.create_node(a, TypeId(1)).unwrap();
2720        w.set_properties(
2721            &a,
2722            Some("Person"),
2723            HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2724        )
2725        .unwrap();
2726        w.flush().unwrap();
2727
2728        let ab = to_bytes(&a);
2729        let updates =
2730            HashMap::from([(ab, HashMap::from([("age".to_owned(), IrLiteral::Int(99))]))]);
2731        set_node_properties(dir.path(), "Person", &updates).unwrap();
2732
2733        assert!(
2734            dir.path()
2735                .join("properties")
2736                .join("Person.parquet")
2737                .exists()
2738        );
2739        let props = read_node_props(dir.path(), "Person");
2740        assert_eq!(props[&ab]["age"], IrLiteral::Int(99));
2741    }
2742
2743    #[test]
2744    fn remove_node_properties_drops_key_and_column_when_last() {
2745        let dir = TempDir::new().unwrap();
2746        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2747        let a = new_v7();
2748        w.create_node(a, TypeId(0)).unwrap();
2749        w.set_properties(
2750            &a,
2751            None,
2752            HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2753        )
2754        .unwrap();
2755        w.flush().unwrap();
2756
2757        let ab = to_bytes(&a);
2758        let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2759        let removals = HashMap::from([(ab, HashSet::from(["age".to_owned()]))]);
2760        let touched = remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
2761        assert_eq!(touched, 1);
2762        assert_eq!(
2763            crate::generation::read_search_generation(dir.path()).unwrap(),
2764            search_generation + 1
2765        );
2766
2767        // The only property was removed → the row's map is empty and the `age`
2768        // column is gone from the re-inferred schema.
2769        let props = read_node_props(dir.path(), "_untyped");
2770        assert!(props.get(&ab).map_or(true, HashMap::is_empty));
2771    }
2772
2773    #[test]
2774    fn remove_node_properties_missing_key_and_uuid_are_noops() {
2775        let dir = TempDir::new().unwrap();
2776        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2777        let a = new_v7();
2778        w.create_node(a, TypeId(0)).unwrap();
2779        w.set_properties(
2780            &a,
2781            None,
2782            HashMap::from([("age".to_owned(), IrLiteral::Int(30))]),
2783        )
2784        .unwrap();
2785        w.flush().unwrap();
2786
2787        let ab = to_bytes(&a);
2788        // Remove a key that isn't there, plus a uuid that doesn't exist.
2789        let removals = HashMap::from([
2790            (ab, HashSet::from(["nope".to_owned()])),
2791            (to_bytes(&new_v7()), HashSet::from(["age".to_owned()])),
2792        ]);
2793        remove_node_properties(dir.path(), "_untyped", &removals).unwrap();
2794
2795        // `age` survives untouched.
2796        let props = read_node_props(dir.path(), "_untyped");
2797        assert_eq!(props[&ab]["age"], IrLiteral::Int(30));
2798    }
2799
2800    #[test]
2801    fn set_and_remove_edge_properties_round_trip() {
2802        let dir = TempDir::new().unwrap();
2803        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2804        let a = new_v7();
2805        let b = new_v7();
2806        w.create_node(a, TypeId(0)).unwrap();
2807        w.create_node(b, TypeId(0)).unwrap();
2808        let e = new_v7();
2809        w.create_edge(e, "KNOWS", &a, &b).unwrap();
2810        w.set_edge_properties(
2811            &e,
2812            Some("KNOWS"),
2813            HashMap::from([("since".to_owned(), IrLiteral::Int(2019))]),
2814        )
2815        .unwrap();
2816        w.flush().unwrap();
2817
2818        let eb = to_bytes(&e);
2819        let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2820        // SET overwrites since.
2821        let updates = HashMap::from([(
2822            eb,
2823            HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
2824        )]);
2825        assert_eq!(
2826            set_edge_properties_rewrite(dir.path(), "KNOWS", &updates).unwrap(),
2827            1
2828        );
2829        assert_eq!(
2830            read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2831            IrLiteral::Int(2020)
2832        );
2833
2834        // REMOVE since.
2835        let removals = HashMap::from([(eb, HashSet::from(["since".to_owned()]))]);
2836        assert_eq!(
2837            remove_edge_properties(dir.path(), "KNOWS", &removals).unwrap(),
2838            1
2839        );
2840        let props = read_edge_props(dir.path(), "KNOWS");
2841        assert!(props.get(&eb).map_or(true, HashMap::is_empty));
2842        assert_eq!(
2843            crate::generation::read_search_generation(dir.path()).unwrap(),
2844            search_generation
2845        );
2846    }
2847
2848    #[test]
2849    fn set_node_properties_empty_map_writes_nothing() {
2850        let dir = TempDir::new().unwrap();
2851        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2852        w.create_node(new_v7(), TypeId(0)).unwrap();
2853        w.flush().unwrap();
2854
2855        let search_generation = crate::generation::read_search_generation(dir.path()).unwrap();
2856        let touched = set_node_properties(dir.path(), "_untyped", &HashMap::new()).unwrap();
2857        assert_eq!(touched, 0);
2858        assert_eq!(
2859            crate::generation::read_search_generation(dir.path()).unwrap(),
2860            search_generation
2861        );
2862        // No property file was created from an empty update set.
2863        assert!(
2864            !dir.path()
2865                .join("properties")
2866                .join("_untyped.parquet")
2867                .exists()
2868        );
2869    }
2870
2871    #[test]
2872    fn staged_set_is_invisible_until_commit_across_stems() {
2873        // One RewriteBatch spanning a node-property stem AND an edge-property
2874        // stem (#790): nothing changes until commit, then both apply at once.
2875        let dir = TempDir::new().unwrap();
2876        let (a, e) = (new_v7(), new_v7());
2877        let b = new_v7();
2878        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2879        w.create_node(a, TypeId(0)).unwrap();
2880        w.create_node(b, TypeId(0)).unwrap();
2881        w.create_edge(e, "KNOWS", &a, &b).unwrap();
2882        w.set_properties(
2883            &a,
2884            None,
2885            HashMap::from([("name".to_owned(), IrLiteral::Str("old".into()))]),
2886        )
2887        .unwrap();
2888        w.set_edge_properties(
2889            &e,
2890            Some("KNOWS"),
2891            HashMap::from([("since".to_owned(), IrLiteral::Int(2000))]),
2892        )
2893        .unwrap();
2894        w.flush().unwrap();
2895
2896        let (ab, eb) = (to_bytes(&a), to_bytes(&e));
2897        let node_updates = HashMap::from([(
2898            ab,
2899            HashMap::from([("name".to_owned(), IrLiteral::Str("new".into()))]),
2900        )]);
2901        let edge_updates = HashMap::from([(
2902            eb,
2903            HashMap::from([("since".to_owned(), IrLiteral::Int(2024))]),
2904        )]);
2905
2906        let mut staged = RewriteBatch::new();
2907        let touched = stage_set_node_properties(&mut staged, dir.path(), "_untyped", &node_updates)
2908            .unwrap()
2909            + stage_set_edge_properties(&mut staged, dir.path(), "KNOWS", &edge_updates).unwrap();
2910        assert_eq!(touched, 2, "one node + one edge written");
2911
2912        // Invisible while staged.
2913        assert_eq!(
2914            read_node_props(dir.path(), "_untyped")[&ab]["name"],
2915            IrLiteral::Str("old".into())
2916        );
2917        assert_eq!(
2918            read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2919            IrLiteral::Int(2000)
2920        );
2921
2922        staged.commit().unwrap();
2923        assert_eq!(
2924            read_node_props(dir.path(), "_untyped")[&ab]["name"],
2925            IrLiteral::Str("new".into())
2926        );
2927        assert_eq!(
2928            read_edge_props(dir.path(), "KNOWS")[&eb]["since"],
2929            IrLiteral::Int(2024)
2930        );
2931    }
2932
2933    // -----------------------------------------------------------------------
2934    // Pending-buffer API (#792)
2935    // -----------------------------------------------------------------------
2936
2937    #[test]
2938    fn pending_nodes_batch_is_canonical_and_non_consuming() {
2939        let dir = TempDir::new().unwrap();
2940        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2941        let empty = w.pending_nodes_batch().unwrap();
2942        assert_eq!(empty.schema(), TOPOLOGY_NODES_SCHEMA.clone());
2943        assert_eq!(empty.num_rows(), 0);
2944
2945        let node = new_v7();
2946        w.create_node_with_labels(node, &[TypeId(3), TypeId(7)])
2947            .unwrap();
2948        for _ in 0..2 {
2949            let batch = w.pending_nodes_batch().unwrap();
2950            assert_eq!(batch.schema(), TOPOLOGY_NODES_SCHEMA.clone());
2951            assert_eq!(batch.num_rows(), 1);
2952            assert_eq!(
2953                batch
2954                    .column_by_name("node_id")
2955                    .unwrap()
2956                    .as_any()
2957                    .downcast_ref::<UInt64Array>()
2958                    .unwrap()
2959                    .value(0),
2960                1
2961            );
2962        }
2963        assert!(w.contains_pending_node(&to_bytes(&node)));
2964    }
2965
2966    #[test]
2967    fn cancel_nodes_drops_rows_props_and_mapping() {
2968        let dir = TempDir::new().unwrap();
2969        let (a, b) = (new_v7(), new_v7());
2970        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
2971        w.create_node(a, TypeId(0)).unwrap();
2972        w.create_node(b, TypeId(0)).unwrap();
2973        w.set_properties(
2974            &a,
2975            None,
2976            HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
2977        )
2978        .unwrap();
2979
2980        assert!(w.contains_pending_node(&to_bytes(&a)));
2981        let dropped = w.cancel_nodes(&HashSet::from([to_bytes(&a)]));
2982        assert_eq!(dropped, 1);
2983        assert!(!w.contains_pending_node(&to_bytes(&a)));
2984
2985        // The mapping is forgotten: an edge to the cancelled node must fail.
2986        let err = w.create_edge(new_v7(), "KNOWS", &b, &a);
2987        assert!(err.is_err(), "edge to a cancelled node must fail");
2988
2989        w.flush().unwrap();
2990        // Only b hit disk; a's property row never did.
2991        let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
2992        let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
2993        assert_eq!(total, 1, "only b persisted");
2994        assert!(
2995            !read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
2996            "cancelled node's props never hit disk"
2997        );
2998    }
2999
3000    #[test]
3001    fn cancel_edges_drops_rows_and_edge_props() {
3002        let dir = TempDir::new().unwrap();
3003        let (a, b) = (new_v7(), new_v7());
3004        let (e1, e2) = (new_v7(), new_v7());
3005        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3006        w.create_node(a, TypeId(0)).unwrap();
3007        w.create_node(b, TypeId(0)).unwrap();
3008        w.create_edge(e1, "KNOWS", &a, &b).unwrap();
3009        w.create_edge(e2, "KNOWS", &b, &a).unwrap();
3010        w.set_edge_properties(
3011            &e1,
3012            Some("KNOWS"),
3013            HashMap::from([("since".to_owned(), IrLiteral::Int(2020))]),
3014        )
3015        .unwrap();
3016
3017        assert!(w.contains_pending_edge(&to_bytes(&e1)));
3018        assert_eq!(w.cancel_edges(&HashSet::from([to_bytes(&e1)])), 1);
3019        assert!(!w.contains_pending_edge(&to_bytes(&e1)));
3020
3021        w.flush().unwrap();
3022        assert!(
3023            !read_edge_props(dir.path(), "KNOWS").contains_key(&to_bytes(&e1)),
3024            "cancelled edge's props never hit disk"
3025        );
3026        let edges = crate::catalog::read_parquet_or_empty(
3027            &dir.path().join("topology/edges/_exploratory.parquet"),
3028            EXPLORATORY_EDGE_SCHEMA.clone(),
3029        )
3030        .unwrap();
3031        let total: usize = edges.iter().map(RecordBatch::num_rows).sum();
3032        assert_eq!(total, 1, "only e2 persisted");
3033    }
3034
3035    #[test]
3036    fn pending_incident_edge_uuids_sees_buffered_edges() {
3037        let dir = TempDir::new().unwrap();
3038        let (a, b, c) = (new_v7(), new_v7(), new_v7());
3039        let e_ab = new_v7();
3040        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3041        w.create_node(a, TypeId(0)).unwrap();
3042        w.create_node(b, TypeId(0)).unwrap();
3043        w.create_node(c, TypeId(0)).unwrap();
3044        w.create_edge(e_ab, "KNOWS", &a, &b).unwrap();
3045
3046        // Incident from either endpoint; c has none.
3047        let hits = w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&b)]));
3048        assert_eq!(hits, vec![to_bytes(&e_ab)]);
3049        assert!(
3050            w.pending_incident_edge_uuids(&HashSet::from([to_bytes(&c)]))
3051                .is_empty()
3052        );
3053    }
3054
3055    #[test]
3056    fn pending_query_and_label_edits_are_exact_before_flush_and_reopen() {
3057        let dir = TempDir::new().unwrap();
3058        let (alice, bob, edge) = (new_v7(), new_v7(), new_v7());
3059        let (alice_bytes, bob_bytes, edge_bytes) =
3060            (to_bytes(&alice), to_bytes(&bob), to_bytes(&edge));
3061        let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS).unwrap();
3062        writer
3063            .create_node_with_labels(alice, &[TypeId(3), TypeId(7)])
3064            .unwrap();
3065        writer.create_node(bob, TypeId(3)).unwrap();
3066        writer
3067            .set_properties(
3068                &alice,
3069                Some("Person"),
3070                HashMap::from([
3071                    ("name".into(), IrLiteral::Str("Alice".into())),
3072                    ("age".into(), IrLiteral::Int(42)),
3073                ]),
3074            )
3075            .unwrap();
3076
3077        assert_eq!(
3078            writer.pending_node_labels(&HashSet::from([alice_bytes, bob_bytes])),
3079            HashSet::from([3, 7])
3080        );
3081        let matched = writer
3082            .find_pending_node(&[3, 7], &[("name".into(), IrLiteral::Str("Alice".into()))])
3083            .unwrap();
3084        assert_eq!(matched.0, alice_bytes);
3085        assert_eq!(matched.2, 3);
3086        assert_eq!(matched.3, vec![3, 7]);
3087        assert_eq!(matched.4["age"], IrLiteral::Int(42));
3088        assert!(writer.find_pending_node(&[9], &[]).is_none());
3089        assert!(
3090            writer
3091                .find_pending_node(&[3], &[("name".into(), IrLiteral::Str("Bob".into()))])
3092                .is_none()
3093        );
3094
3095        assert_eq!(writer.add_pending_node_labels(&alice_bytes, &[7, 9]), 1);
3096        assert_eq!(writer.add_pending_node_labels(&[0xff; 16], &[1]), 0);
3097        assert_eq!(writer.remove_pending_node_labels(&alice_bytes, &[7, 99]), 1);
3098        assert_eq!(writer.remove_pending_node_labels(&[0xff; 16], &[1]), 0);
3099        assert_eq!(
3100            writer.pending_node_labels(&HashSet::from([alice_bytes])),
3101            HashSet::from([3, 9])
3102        );
3103
3104        writer.create_edge(edge, "KNOWS", &alice, &bob).unwrap();
3105        writer
3106            .set_edge_properties(
3107                &edge,
3108                Some("KNOWS"),
3109                HashMap::from([("since".into(), IrLiteral::Int(2024))]),
3110            )
3111            .unwrap();
3112        let direct = writer
3113            .find_pending_edge(
3114                "KNOWS",
3115                &alice_bytes,
3116                &bob_bytes,
3117                false,
3118                &[("since".into(), IrLiteral::Int(2024))],
3119            )
3120            .unwrap();
3121        assert_eq!(direct.0, edge_bytes);
3122        assert_eq!(direct.1, alice_bytes);
3123        assert_eq!(direct.2, bob_bytes);
3124        assert_eq!(direct.3["since"], IrLiteral::Int(2024));
3125        assert!(
3126            writer
3127                .find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, false, &[])
3128                .is_none()
3129        );
3130        assert!(
3131            writer
3132                .find_pending_edge("KNOWS", &bob_bytes, &alice_bytes, true, &[])
3133                .is_some()
3134        );
3135        assert!(
3136            writer
3137                .find_pending_edge("IGNORES", &alice_bytes, &bob_bytes, false, &[])
3138                .is_none()
3139        );
3140
3141        writer.flush().unwrap();
3142        let mut reopened = GraphWriter::open_at(dir.path(), OntologyMode::Strict, TS + 1).unwrap();
3143        assert_eq!(reopened.create_node(new_v7(), TypeId(3)).unwrap(), 3);
3144        assert_eq!(
3145            read_node_props(dir.path(), "Person")[&alice_bytes]["name"],
3146            IrLiteral::Str("Alice".into())
3147        );
3148        assert_eq!(
3149            read_edge_props(dir.path(), "KNOWS")[&edge_bytes]["since"],
3150            IrLiteral::Int(2024)
3151        );
3152    }
3153
3154    #[test]
3155    fn merge_and_remove_pending_props_edit_buffered_rows() {
3156        let dir = TempDir::new().unwrap();
3157        let a = new_v7();
3158        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3159        w.create_node(a, TypeId(0)).unwrap();
3160        w.set_properties(
3161            &a,
3162            None,
3163            HashMap::from([
3164                ("name".to_owned(), IrLiteral::Str("old".into())),
3165                ("age".to_owned(), IrLiteral::Int(30)),
3166            ]),
3167        )
3168        .unwrap();
3169
3170        // SET on the pending node: overwrite one key, add another.
3171        w.merge_pending_node_props(
3172            &to_bytes(&a),
3173            None,
3174            HashMap::from([
3175                ("name".to_owned(), IrLiteral::Str("new".into())),
3176                ("city".to_owned(), IrLiteral::Str("Oslo".into())),
3177            ]),
3178        );
3179        // REMOVE on the pending node: drop a key; absent keys are no-ops.
3180        w.remove_pending_node_props(
3181            &to_bytes(&a),
3182            &HashSet::from(["age".to_owned(), "absent".to_owned()]),
3183        );
3184
3185        w.flush().unwrap();
3186        let props = &read_node_props(dir.path(), "_untyped")[&to_bytes(&a)];
3187        assert_eq!(props["name"], IrLiteral::Str("new".into()));
3188        assert_eq!(props["city"], IrLiteral::Str("Oslo".into()));
3189        assert!(!props.contains_key("age"), "removed before flush");
3190    }
3191
3192    #[test]
3193    fn flush_into_composes_with_staged_delete_in_one_batch() {
3194        // The #792 statement shape: DELETE a committed node and CREATE a new
3195        // one in the same statement — one RewriteBatch, one commit, with the
3196        // flush reading through the delete's staged nodes.parquet content.
3197        let dir = TempDir::new().unwrap();
3198        let (a, b) = (new_v7(), new_v7());
3199        let mut seed = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3200        seed.create_node(a, TypeId(0)).unwrap();
3201        seed.create_node(b, TypeId(0)).unwrap();
3202        seed.set_properties(
3203            &a,
3204            None,
3205            HashMap::from([("name".to_owned(), IrLiteral::Str("A".into()))]),
3206        )
3207        .unwrap();
3208        seed.flush().unwrap();
3209
3210        let mut staged = RewriteBatch::new();
3211        let removed = crate::mutator::stage_delete_nodes(
3212            &mut staged,
3213            dir.path(),
3214            &HashSet::from([to_bytes(&a)]),
3215        )
3216        .unwrap();
3217        assert_eq!(removed, 1);
3218
3219        let d = new_v7();
3220        let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3221        w.create_node(d, TypeId(0)).unwrap();
3222        w.flush_into(&mut staged).unwrap();
3223
3224        // nodes.parquet staged exactly once (net content), still last-ish in
3225        // commit order relative to the delete's property rewrite.
3226        let staged_nodes = staged
3227            .staged_paths()
3228            .filter(|p| p.ends_with("topology/nodes.parquet"))
3229            .count();
3230        assert_eq!(staged_nodes, 1, "net content, no double-stage");
3231
3232        // Nothing visible before commit.
3233        let pre: usize = crate::catalog::read_nodes(dir.path())
3234            .unwrap()
3235            .iter()
3236            .map(RecordBatch::num_rows)
3237            .sum();
3238        assert_eq!(pre, 2);
3239
3240        staged.commit().unwrap();
3241        let nodes = crate::catalog::read_nodes(dir.path()).unwrap();
3242        let total: usize = nodes.iter().map(RecordBatch::num_rows).sum();
3243        assert_eq!(total, 2, "b survives, a deleted, d created");
3244        assert!(
3245            !read_node_props(dir.path(), "_untyped").contains_key(&to_bytes(&a)),
3246            "deleted node's props gone"
3247        );
3248    }
3249
3250    #[test]
3251    fn every_persisted_property_family_round_trips_through_parquet_reopen() {
3252        let dir = TempDir::new().unwrap();
3253        let node = new_v7();
3254        let propertyless = new_v7();
3255        let values = HashMap::from([
3256            ("int".into(), IrLiteral::Int(-7)),
3257            ("float".into(), IrLiteral::Float(2.5)),
3258            ("bool".into(), IrLiteral::Bool(true)),
3259            ("str".into(), IrLiteral::Str("value".into())),
3260            (
3261                "duration".into(),
3262                IrLiteral::Duration {
3263                    months: 1,
3264                    days: -2,
3265                    seconds: 3,
3266                    nanos: 4,
3267                },
3268            ),
3269            ("datetime".into(), IrLiteral::DateTime(TS)),
3270            ("date".into(), IrLiteral::Date(19_000)),
3271            (
3272                "local_datetime".into(),
3273                IrLiteral::LocalDateTime {
3274                    days: 19_001,
3275                    nanos: 123,
3276                },
3277            ),
3278            ("time".into(), IrLiteral::Time(456)),
3279            (
3280                "zoned_time".into(),
3281                IrLiteral::ZonedTime {
3282                    nanos: 789,
3283                    offset: -21_600,
3284                },
3285            ),
3286            (
3287                "zoned_datetime".into(),
3288                IrLiteral::ZonedDateTime {
3289                    days: 19_002,
3290                    nanos: 987,
3291                    offset: 3_600,
3292                    zone: Some("Europe/Paris".into()),
3293                },
3294            ),
3295            (
3296                "offset_datetime".into(),
3297                IrLiteral::ZonedDateTime {
3298                    days: 19_003,
3299                    nanos: 654,
3300                    offset: 0,
3301                    zone: None,
3302                },
3303            ),
3304            (
3305                "ints".into(),
3306                IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Null, IrLiteral::Int(3)]),
3307            ),
3308            (
3309                "dates".into(),
3310                IrLiteral::List(vec![IrLiteral::Date(19_004), IrLiteral::Date(19_005)]),
3311            ),
3312            ("empty".into(), IrLiteral::List(Vec::new())),
3313            ("null".into(), IrLiteral::Null),
3314        ]);
3315        let mut writer = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, TS).unwrap();
3316        writer.create_node(node, TypeId(0)).unwrap();
3317        writer.create_node(propertyless, TypeId(0)).unwrap();
3318        writer.set_properties(&node, None, values.clone()).unwrap();
3319        writer
3320            .set_properties(
3321                &propertyless,
3322                None,
3323                values
3324                    .keys()
3325                    .map(|name| (name.clone(), IrLiteral::Null))
3326                    .collect(),
3327            )
3328            .unwrap();
3329        writer.flush().unwrap();
3330
3331        let reopened = read_node_props(dir.path(), "_untyped");
3332        let actual = reopened.get(&to_bytes(&node)).unwrap();
3333        for (name, expected) in &values {
3334            if matches!(expected, IrLiteral::Null) {
3335                assert!(!actual.contains_key(name));
3336            } else if name == "empty" {
3337                assert_eq!(actual.get(name), Some(&IrLiteral::Str("[]".into())));
3338            } else {
3339                assert_eq!(actual.get(name), Some(expected), "property {name}");
3340            }
3341        }
3342        assert!(
3343            reopened
3344                .get(&to_bytes(&propertyless))
3345                .is_none_or(HashMap::is_empty)
3346        );
3347    }
3348
3349    #[test]
3350    fn property_literal_rendering_and_nested_invalid_values_are_deterministic() {
3351        let uuid = [0xabu8; 16];
3352        let cases = [
3353            (IrLiteral::Null, "".into()),
3354            (IrLiteral::Bool(true), "true".into()),
3355            (IrLiteral::Int(-2), "-2".into()),
3356            (IrLiteral::Float(1.25), "1.25".into()),
3357            (IrLiteral::Str("s".into()), "s".into()),
3358            (IrLiteral::Uuid(uuid), "ab".repeat(16)),
3359            (
3360                IrLiteral::Duration {
3361                    months: 1,
3362                    days: 2,
3363                    seconds: 3,
3364                    nanos: 4,
3365                },
3366                "1mo2d3s4ns".into(),
3367            ),
3368            (IrLiteral::DateTime(5), "5".into()),
3369            (IrLiteral::Date(6), "6".into()),
3370            (
3371                IrLiteral::LocalDateTime { days: 7, nanos: 8 },
3372                "7d8ns".into(),
3373            ),
3374            (IrLiteral::Time(9), "9ns".into()),
3375            (
3376                IrLiteral::ZonedTime {
3377                    nanos: 10,
3378                    offset: -1,
3379                },
3380                "10ns-1s".into(),
3381            ),
3382            (
3383                IrLiteral::ZonedDateTime {
3384                    days: 11,
3385                    nanos: 12,
3386                    offset: 13,
3387                    zone: Some("UTC".into()),
3388                },
3389                "11d12ns+13sUTC".into(),
3390            ),
3391            (
3392                IrLiteral::List(vec![IrLiteral::Int(1), IrLiteral::Str("x".into())]),
3393                "[1,x]".into(),
3394            ),
3395            (
3396                IrLiteral::Map(vec![("a".into(), IrLiteral::Bool(false))]),
3397                "{a:false}".into(),
3398            ),
3399        ];
3400        for (literal, expected) in cases {
3401            assert_eq!(literal_to_string(&literal), expected);
3402        }
3403
3404        for invalid in [
3405            IrLiteral::Uuid(uuid),
3406            IrLiteral::List(vec![IrLiteral::Uuid(uuid)]),
3407            IrLiteral::Map(vec![("nested".into(), IrLiteral::Uuid(uuid))]),
3408        ] {
3409            assert_eq!(
3410                reject_map_property_value("p", &invalid).unwrap_err().code(),
3411                "GF_VALIDATION"
3412            );
3413        }
3414        for invalid in [
3415            IrLiteral::Map(vec![]),
3416            IrLiteral::List(vec![IrLiteral::Map(vec![])]),
3417        ] {
3418            assert_eq!(
3419                reject_map_property_value("p", &invalid).unwrap_err().code(),
3420                "GF_IO"
3421            );
3422        }
3423    }
3424
3425    #[test]
3426    fn persisted_property_decoder_rejects_unsupported_shape_type_and_dynamic_array() {
3427        use arrow::array::{Int32Array, Int64Array, StructArray, UInt8Array};
3428        use arrow::datatypes::{DataType, Field, Fields};
3429
3430        let unsupported: arrow::array::ArrayRef = Arc::new(UInt8Array::from(vec![1]));
3431        let unsupported_field = Field::new("unsupported", DataType::UInt8, false);
3432        assert!(decode_value(&unsupported, &unsupported_field, 0).is_err());
3433
3434        let fields: Fields = vec![Field::new("other", DataType::Int32, false)].into();
3435        let structure: arrow::array::ArrayRef = Arc::new(StructArray::new(
3436            fields.clone(),
3437            vec![Arc::new(Int32Array::from(vec![1]))],
3438            None,
3439        ));
3440        let structure_field = Field::new("structure", DataType::Struct(fields), false);
3441        assert!(decode_value(&structure, &structure_field, 0).is_err());
3442
3443        let wrong_dynamic: arrow::array::ArrayRef = Arc::new(Int64Array::from(vec![1]));
3444        let declared = Field::new("declared", DataType::UInt64, false);
3445        assert!(decode_value(&wrong_dynamic, &declared, 0).is_err());
3446    }
3447}