Skip to main content

graphforge_plan/
lib.rs

1//! GraphForge DataFusion integration, optimizer rules, and custom plan nodes.
2//!
3//! # Custom logical plan nodes (#577)
4//!
5//! Six graph-native operators cannot be expressed in relational algebra.  This
6//! module defines their **logical plan stubs** — enough for the lowering layer
7//! to produce a valid [`LogicalPlan`] that DataFusion accepts without panicking.
8//! Physical implementations come in milestone 13 (Execution Baseline).
9//!
10//! | Node | Triggered by |
11//! |---|---|
12//! | [`VarLenExpandNode`] | `Expand` with `max_hops != Some(1)` |
13//! | [`OptionalMatchNode`] | `Optional { child }` |
14//! | [`PathUniqueNode`] | `Expand` with path-uniqueness flag |
15//! | [`OntologyInferNode`] | `Expand` on transitive/symmetric relation |
16//! | [`GraphMergeNode`] | `Merge { pattern }` |
17//! | [`UnwindNode`] | `Unwind { list_expr, alias }` |
18#![forbid(unsafe_code)]
19// The `name()` methods return string literals but the trait signature requires
20// `&str` tied to `&self`; the lint fires because the bound is unnecessary for
21// the literal values.
22#![allow(clippy::unnecessary_literal_bound)]
23
24use std::cmp::Ordering;
25use std::collections::HashMap;
26use std::fmt;
27use std::hash::{Hash, Hasher};
28use std::path::PathBuf;
29use std::sync::Arc;
30
31use datafusion::arrow::datatypes::{DataType, Field, Fields, Schema};
32use datafusion::common::{DFSchema, DFSchemaRef, Result as DfResult, TableReference};
33use datafusion::logical_expr::{Expr, ExprSchemable, LogicalPlan, UserDefinedLogicalNodeCore};
34
35use graphforge_core::OntologyMode;
36use graphforge_ir::{Direction, IrLiteral};
37
38/// Implement `PartialOrd` for custom plan nodes that contain `DFSchemaRef`.
39///
40/// `DFSchemaRef` (`Arc<DFSchema>`) does not implement `PartialOrd`, so we
41/// cannot derive it.  Returning `None` (incomparable) is the safe default
42/// used by DataFusion's own test nodes.
43macro_rules! impl_partial_ord {
44    ($t:ty) => {
45        impl PartialOrd for $t {
46            fn partial_cmp(&self, _other: &Self) -> Option<Ordering> {
47                None
48            }
49        }
50    };
51}
52
53// ---------------------------------------------------------------------------
54// Shared helpers
55// ---------------------------------------------------------------------------
56
57/// A placeholder output schema that mirrors the input schema.  Used by stubs
58/// whose real output schema is determined at physical planning time.
59fn passthrough_schema(inputs: &[&LogicalPlan]) -> DFSchemaRef {
60    inputs
61        .first()
62        .map_or_else(|| Arc::new(DFSchema::empty()), |p| p.schema().clone())
63}
64
65/// Arrow field name of the variable-length edge-list column (qualified
66/// `var_<edge_var>`).  openCypher binds the edge variable `r` in
67/// `(a)-[r:KNOWS*1..3]->(b)` to the **list of relationships** along each path.
68///
69/// Lowercase ASCII so the dotted reference `col("var_<edge>.rels")` resolves —
70/// DataFusion lowercases unquoted identifiers.
71pub const VAR_LEN_EDGE_LIST_FIELD: &str = "rels";
72
73/// The Arrow [`Field`] for the variable-length edge-list column: a nullable
74/// `List<Struct<{edge_uuid, src_uuid, dst_uuid, rel_type, <props…>}>>`.
75///
76/// This is the **single source of truth** for the column's type, shared by the
77/// lowerer (which puts it in [`VarLenExpandNode`]'s schema) and the physical
78/// `VarLenExpandExec` (which produces it).  They must be byte-identical or the
79/// positional `RecordBatch::try_new` at execution time fails the schema check.
80///
81/// The struct carries only **public** UUID identity + the relation type — never
82/// the surrogate `edge_id`/`src_id`/`dst_id` — so the UUID-only output contract
83/// holds (the graphforge-api Shaper passes the whole column through untouched).
84///
85/// `prop_fields` are the relation's persisted edge-property columns (#755),
86/// appended after the four topology fields in order. They are discovered at
87/// lowering time from `edge_properties/<REL>.parquet` (see
88/// `lower_var_len_expand`); pass an empty slice for a topology-only struct
89/// (wildcard `*`, or a relation with no persisted properties), which reproduces
90/// the original four-field layout byte-for-byte.
91#[must_use]
92pub fn var_len_edge_list_field(prop_fields: &[Field]) -> Arc<Field> {
93    let mut fields = vec![
94        Field::new("edge_uuid", DataType::FixedSizeBinary(16), false),
95        Field::new("src_uuid", DataType::FixedSizeBinary(16), false),
96        Field::new("dst_uuid", DataType::FixedSizeBinary(16), false),
97        Field::new("rel_type", DataType::Utf8, true),
98    ];
99    fields.extend(prop_fields.iter().cloned());
100    let struct_fields = Fields::from(fields);
101    // The item is NULLABLE even though the exec never emits a null element:
102    // DataFusion's list kernels (`array_pop_front` for `tail(r)`, #1023) return
103    // item-nullable lists, and the runtime asserts the kernel's output type
104    // matches the type promised at planning time — a non-null item fails that
105    // assertion the moment any list function touches the column.
106    let item = Arc::new(Field::new("item", DataType::Struct(struct_fields), true));
107    Arc::new(Field::new(
108        VAR_LEN_EDGE_LIST_FIELD,
109        DataType::List(item),
110        true,
111    ))
112}
113
114// ---------------------------------------------------------------------------
115// VarLenExpandNode
116// ---------------------------------------------------------------------------
117
118/// Physical node for variable-length path expansion.
119///
120/// Triggered when `Expand` has `max_hops != Some(1)` — patterns like
121/// `(a)-[:KNOWS*1..3]->(b)`.  These cannot be expressed as a finite join
122/// sequence, so the physical layer ([`VarLenExpandExec`](../graphforge_exec) in M13)
123/// performs an iterative BFS over the Parquet edge table.
124///
125/// # Baked execution context
126///
127/// Like [`GraphCreateNode`], the project `dir` and ontology `mode` are baked in
128/// at lowering time because the physical-planning `ExtensionPlanner` only sees
129/// the DataFusion session state, not the GraphForge project path.  The physical
130/// node reads edges directly from `dir` at execution time.
131///
132/// # Output schema
133///
134/// The output extends the input schema with the **destination** node's columns
135/// (`var_<dst>`-qualified [`TOPOLOGY_NODES_SCHEMA`](../graphforge_storage)), then a
136/// trailing edge-list column (`var_<edge>`-qualified [`var_len_edge_list_field`])
137/// binding the edge variable `r` to the openCypher *list of relationships* along
138/// each path (#709).  The edge column is **last**; the physical node produces
139/// columns in this exact order.
140#[derive(Debug, Clone, PartialEq, Eq, Hash)]
141pub struct VarLenExpandNode {
142    /// Input plan (the source node scan / prior pipeline).
143    pub input: Arc<LogicalPlan>,
144    /// Name of the relation type to expand along (or `"*"` for wildcard).
145    pub rel_type_name: String,
146    /// Minimum number of hops.
147    pub min_hops: u16,
148    /// Maximum number of hops (`None` = unbounded).
149    pub max_hops: Option<u16>,
150    /// Source node pattern-variable id (the frontier seed in the input).
151    pub src_var: u32,
152    /// Destination node pattern-variable id (bound to the reached node).
153    pub dst_var: u32,
154    /// Edge pattern-variable id (bound to the per-path relationship list).
155    pub edge_var: u32,
156    /// Edge traversal direction.
157    pub direction: Direction,
158    /// Resolved relation type id (`TypeId.0`), or `None` for a wildcard.
159    pub rel_ty: Option<u32>,
160    /// Project directory the physical node reads edges from.
161    pub dir: PathBuf,
162    /// Ontology mode (drives typed vs exploratory edge-file routing).
163    pub mode: OntologyMode,
164    schema: DFSchemaRef,
165}
166
167impl VarLenExpandNode {
168    /// Create a variable-length expand node.
169    ///
170    /// `dst_fields` is the destination node's column list (the storage layer's
171    /// `TOPOLOGY_NODES_SCHEMA` fields), passed in by the lowerer so this crate
172    /// need not depend on `graphforge-storage`.  They are qualified `var_<dst_var>` and
173    /// appended to the input schema.  `edge_field` is the trailing edge-list
174    /// column ([`var_len_edge_list_field`]), qualified `var_<edge_var>`.
175    #[allow(clippy::too_many_arguments)]
176    #[must_use]
177    pub fn new(
178        input: Arc<LogicalPlan>,
179        rel_type_name: impl Into<String>,
180        min_hops: u16,
181        max_hops: Option<u16>,
182        src_var: u32,
183        dst_var: u32,
184        edge_var: u32,
185        direction: Direction,
186        rel_ty: Option<u32>,
187        dir: PathBuf,
188        mode: OntologyMode,
189        dst_fields: Vec<Arc<Field>>,
190        edge_field: Arc<Field>,
191    ) -> Self {
192        let schema = Self::build_schema(&input, dst_var, edge_var, dst_fields, edge_field);
193        Self {
194            input,
195            rel_type_name: rel_type_name.into(),
196            min_hops,
197            max_hops,
198            src_var,
199            dst_var,
200            edge_var,
201            direction,
202            rel_ty,
203            dir,
204            mode,
205            schema,
206        }
207    }
208
209    /// Build the output [`DFSchema`]: the input's qualified fields, then the
210    /// destination node's fields qualified `var_<dst_var>`, then the edge-list
211    /// column qualified `var_<edge_var>` (last).
212    fn build_schema(
213        input: &Arc<LogicalPlan>,
214        dst_var: u32,
215        edge_var: u32,
216        dst_fields: Vec<Arc<Field>>,
217        edge_field: Arc<Field>,
218    ) -> DFSchemaRef {
219        let dst_qualifier = TableReference::bare(format!("var_{dst_var}"));
220        let edge_qualifier = TableReference::bare(format!("var_{edge_var}"));
221        let mut qualified: Vec<(Option<TableReference>, Arc<Field>)> = input
222            .schema()
223            .iter()
224            .map(|(q, f)| (q.cloned(), Arc::clone(f)))
225            .collect();
226        qualified.extend(
227            dst_fields
228                .into_iter()
229                .map(|f| (Some(dst_qualifier.clone()), f)),
230        );
231        // The edge-list column is appended LAST; the physical node produces
232        // columns in this order.
233        qualified.push((Some(edge_qualifier), edge_field));
234        // Fail fast: a construction error here means the node would advertise a
235        // schema missing its `var_<dst>`/`var_<edge>` columns, breaking
236        // downstream resolution.  This only fails on a malformed field set (a
237        // programmer error), so panic with context rather than degrading
238        // silently.
239        Arc::new(
240            DFSchema::new_with_metadata(qualified, std::collections::HashMap::new())
241                .expect("VarLenExpandNode schema must include qualified destination + edge fields"),
242        )
243    }
244}
245
246impl_partial_ord!(VarLenExpandNode);
247
248impl UserDefinedLogicalNodeCore for VarLenExpandNode {
249    fn name(&self) -> &str {
250        "VarLenExpand"
251    }
252
253    fn inputs(&self) -> Vec<&LogicalPlan> {
254        vec![&self.input]
255    }
256
257    fn schema(&self) -> &DFSchemaRef {
258        &self.schema
259    }
260
261    fn expressions(&self) -> Vec<Expr> {
262        vec![]
263    }
264
265    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
266        let max = self.max_hops.map_or("*".to_owned(), |h| h.to_string());
267        let arrow = match self.direction {
268            Direction::Out => "->",
269            Direction::In => "<-",
270            Direction::Undirected => "--",
271        };
272        write!(
273            f,
274            "VarLenExpand: rel={}, hops={}..{}, dir={arrow}, edge=var_{}",
275            self.rel_type_name, self.min_hops, max, self.edge_var
276        )
277    }
278
279    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
280        let input = Arc::new(
281            inputs
282                .into_iter()
283                .next()
284                .unwrap_or_else(|| (*self.input).clone()),
285        );
286        // Reuse the already-computed destination fields (strip the qualifier);
287        // the input schema is the only piece that can change here.
288        let dst_qual = format!("var_{}", self.dst_var);
289        let dst_fields: Vec<Arc<Field>> = self
290            .schema
291            .iter()
292            .filter(|(q, _)| q.map(TableReference::table) == Some(dst_qual.as_str()))
293            .map(|(_, f)| Arc::clone(f))
294            .collect();
295        // The trailing edge-list column (exactly one field qualified `var_<edge>`).
296        let edge_qual = format!("var_{}", self.edge_var);
297        let edge_field = self
298            .schema
299            .iter()
300            .find(|(q, _)| q.map(TableReference::table) == Some(edge_qual.as_str()))
301            .map_or_else(|| var_len_edge_list_field(&[]), |(_, f)| Arc::clone(f));
302        Ok(Self::new(
303            input,
304            self.rel_type_name.clone(),
305            self.min_hops,
306            self.max_hops,
307            self.src_var,
308            self.dst_var,
309            self.edge_var,
310            self.direction,
311            self.rel_ty,
312            self.dir.clone(),
313            self.mode,
314            dst_fields,
315            edge_field,
316        ))
317    }
318}
319
320// ---------------------------------------------------------------------------
321// ExpandNode (#763, #1248)
322// ---------------------------------------------------------------------------
323
324/// Physical node for adjacency-backed single-hop expansion (#763).
325///
326/// Emitted by the lowerer **instead of** the two-join chain whenever a project
327/// read target is available. The physical `ExpandExec` (graphforge-exec) probes the
328/// session adjacency provider per frontier row; that provider owns the
329/// hit/miss/building fallback policy, so plan shape is stable across index
330/// state and downstream `LIMIT` can cancel traversal work (#1248).
331///
332/// # Output schema (join-path parity)
333///
334/// Exactly what the join chain would have produced, in the same order: the
335/// input's qualified fields, then the edge topology fields (typed or
336/// exploratory/wildcard) qualified `var_<edge_var>`, then persisted
337/// edge-property fields (forced **nullable** — the join path LEFT-joins them)
338/// also under `var_<edge_var>`, then the destination node's topology fields
339/// qualified `var_<dst_var>`. Destination type filtering and property joining
340/// stay downstream in the binder's trailing `NodeScan{dst, ty}` (#789),
341/// exactly as on the join path.
342///
343/// For `Undirected`, `ExpandExec` collapses the provider's duplicate self-loop
344/// entries per input row, matching the relational union shape.
345#[derive(Debug, Clone, PartialEq, Eq, Hash)]
346pub struct ExpandNode {
347    /// Input plan (the source node scan / prior pipeline).
348    pub input: Arc<LogicalPlan>,
349    /// Name of the relation type to expand along (`"*"` for wildcard).
350    pub rel_type_name: String,
351    /// Source node pattern-variable id (the frontier seed in the input).
352    pub src_var: u32,
353    /// Destination node pattern-variable id (bound to the reached node).
354    pub dst_var: u32,
355    /// Edge pattern-variable id (bound per-column, like the join path).
356    pub edge_var: u32,
357    /// Edge traversal direction.
358    pub direction: Direction,
359    /// Resolved relation type id (`TypeId.0`), absent for wildcard expansion.
360    pub rel_ty: Option<u32>,
361    /// Project directory the physical node reads from.
362    pub dir: PathBuf,
363    /// Ontology mode controlling the persisted edge layout.
364    pub mode: OntologyMode,
365    /// How many of the `var_<edge_var>` fields are edge-property columns
366    /// (the trailing ones); the rest are edge topology columns.
367    pub edge_prop_count: usize,
368    schema: DFSchemaRef,
369}
370
371impl ExpandNode {
372    /// Create an adjacency-backed single-hop expand node.
373    ///
374    /// `edge_fields` are the typed edge table's topology columns and
375    /// `edge_prop_fields` the relation's persisted property columns (the
376    /// lowerer discovers them from `edge_properties/<REL>.parquet` and forces
377    /// them nullable); `dst_fields` are the destination node's topology
378    /// columns. Passed in so this crate need not depend on graphforge-storage.
379    #[allow(clippy::too_many_arguments)]
380    #[must_use]
381    pub fn new(
382        input: Arc<LogicalPlan>,
383        rel_type_name: impl Into<String>,
384        src_var: u32,
385        dst_var: u32,
386        edge_var: u32,
387        direction: Direction,
388        rel_ty: Option<u32>,
389        dir: PathBuf,
390        mode: OntologyMode,
391        edge_fields: Vec<Arc<Field>>,
392        edge_prop_fields: Vec<Arc<Field>>,
393        dst_fields: Vec<Arc<Field>>,
394    ) -> Self {
395        let edge_prop_count = edge_prop_fields.len();
396        let schema = Self::build_schema(
397            &input,
398            dst_var,
399            edge_var,
400            edge_fields,
401            edge_prop_fields,
402            dst_fields,
403        );
404        Self {
405            input,
406            rel_type_name: rel_type_name.into(),
407            src_var,
408            dst_var,
409            edge_var,
410            direction,
411            rel_ty,
412            dir,
413            mode,
414            edge_prop_count,
415            schema,
416        }
417    }
418
419    /// Build the output [`DFSchema`]: input qualified fields ++ edge topology
420    /// fields (`var_<edge>`) ++ edge property fields (`var_<edge>`, forced
421    /// nullable) ++ destination node fields (`var_<dst>`).
422    fn build_schema(
423        input: &Arc<LogicalPlan>,
424        dst_var: u32,
425        edge_var: u32,
426        edge_fields: Vec<Arc<Field>>,
427        edge_prop_fields: Vec<Arc<Field>>,
428        dst_fields: Vec<Arc<Field>>,
429    ) -> DFSchemaRef {
430        let edge_qualifier = TableReference::bare(format!("var_{edge_var}"));
431        let dst_qualifier = TableReference::bare(format!("var_{dst_var}"));
432        let mut qualified: Vec<(Option<TableReference>, Arc<Field>)> = input
433            .schema()
434            .iter()
435            .map(|(q, f)| (q.cloned(), Arc::clone(f)))
436            .collect();
437        qualified.extend(
438            edge_fields
439                .into_iter()
440                .map(|f| (Some(edge_qualifier.clone()), f)),
441        );
442        // Property columns are nullable on the join path (LEFT join); force it
443        // here so the schemas agree regardless of the file's declared
444        // nullability.
445        qualified.extend(edge_prop_fields.into_iter().map(|f| {
446            let f = Arc::new(f.as_ref().clone().with_nullable(true));
447            (Some(edge_qualifier.clone()), f)
448        }));
449        qualified.extend(
450            dst_fields
451                .into_iter()
452                .map(|f| (Some(dst_qualifier.clone()), f)),
453        );
454        Arc::new(
455            DFSchema::new_with_metadata(qualified, std::collections::HashMap::new())
456                .expect("ExpandNode schema must include qualified edge + destination fields"),
457        )
458    }
459}
460
461impl_partial_ord!(ExpandNode);
462
463impl UserDefinedLogicalNodeCore for ExpandNode {
464    fn name(&self) -> &str {
465        "Expand"
466    }
467
468    fn inputs(&self) -> Vec<&LogicalPlan> {
469        vec![&self.input]
470    }
471
472    fn schema(&self) -> &DFSchemaRef {
473        &self.schema
474    }
475
476    fn expressions(&self) -> Vec<Expr> {
477        vec![]
478    }
479
480    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
481        let arrow = match self.direction {
482            Direction::Out => "->",
483            Direction::In => "<-",
484            Direction::Undirected => "--",
485        };
486        write!(
487            f,
488            "Expand: rel={}, dir={arrow}, edge=var_{}",
489            self.rel_type_name, self.edge_var
490        )
491    }
492
493    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
494        let input = Arc::new(
495            inputs
496                .into_iter()
497                .next()
498                .unwrap_or_else(|| (*self.input).clone()),
499        );
500        // Recover the edge/dst fields from this node's own schema by
501        // qualifier; only the input can change under optimizer rewrites.
502        let edge_qual = format!("var_{}", self.edge_var);
503        let all_edge: Vec<Arc<Field>> = self
504            .schema
505            .iter()
506            .filter(|(q, _)| q.map(TableReference::table) == Some(edge_qual.as_str()))
507            .map(|(_, f)| Arc::clone(f))
508            .collect();
509        let topo_count = all_edge.len().saturating_sub(self.edge_prop_count);
510        let edge_fields = all_edge[..topo_count].to_vec();
511        let edge_prop_fields = all_edge[topo_count..].to_vec();
512        let dst_qual = format!("var_{}", self.dst_var);
513        let dst_fields: Vec<Arc<Field>> = self
514            .schema
515            .iter()
516            .filter(|(q, _)| q.map(TableReference::table) == Some(dst_qual.as_str()))
517            .map(|(_, f)| Arc::clone(f))
518            .collect();
519        Ok(Self::new(
520            input,
521            self.rel_type_name.clone(),
522            self.src_var,
523            self.dst_var,
524            self.edge_var,
525            self.direction,
526            self.rel_ty,
527            self.dir.clone(),
528            self.mode,
529            edge_fields,
530            edge_prop_fields,
531            dst_fields,
532        ))
533    }
534}
535
536// ---------------------------------------------------------------------------
537// OptionalMatchNode
538// ---------------------------------------------------------------------------
539
540/// Physical node for `OPTIONAL MATCH` (LEFT OUTER semantics with openCypher
541/// null-shaping over a sub-plan).
542///
543/// The physical layer ([`OptionalMatchExec`](../graphforge_exec) in M13) left-joins the
544/// `outer` (mandatory) input against the `optional` sub-plan on the shared
545/// pattern variables ([`join_keys`](Self::join_keys)), preserving every outer
546/// row and setting the optional-side columns to **null** when there is no match
547/// — distinct from a SQL `LEFT JOIN` only in that the shared join-key columns
548/// are not duplicated on the output (they belong to the outer side).
549///
550/// # Output schema
551///
552/// `outer` fields, followed by the `optional` fields that are **not** join keys,
553/// each made nullable (an unmatched outer row nulls them all).
554#[derive(Debug, Clone, PartialEq, Eq, Hash)]
555pub struct OptionalMatchNode {
556    /// The outer (mandatory) input.
557    pub outer: Arc<LogicalPlan>,
558    /// The optional sub-plan.
559    pub optional: Arc<LogicalPlan>,
560    /// Shared-variable join keys as `(outer_col_idx, inner_col_idx)` pairs,
561    /// resolved against the qualified input schemas at lowering time.
562    pub join_keys: Vec<(usize, usize)>,
563    /// Inner (optional-side) column indices to append to the output, in order.
564    /// Excludes **every** shared-variable column — not merely the join-key
565    /// columns — because a shared variable's whole node (e.g. `a` in
566    /// `MATCH (a) OPTIONAL MATCH (a)-[:R]->(b)`) is carried by the outer side;
567    /// appending its columns again would duplicate the `var_<shared>` fields.
568    pub inner_keep_idx: Vec<usize>,
569    schema: DFSchemaRef,
570}
571
572impl OptionalMatchNode {
573    /// Create an optional-match node.
574    ///
575    /// `inner_keep_idx` lists the optional plan's column indices to append to
576    /// the output (every column of a shared variable already excluded — those
577    /// live on the outer side). Passed by the lowerer, which resolves them
578    /// against the qualified inner schema, so this crate needs no
579    /// `graphforge-rel`/`graphforge-storage` dependency. The node appends the corresponding
580    /// fields to the outer schema as **nullable** columns (an unmatched outer
581    /// row nulls them all).
582    #[must_use]
583    pub fn new(
584        outer: Arc<LogicalPlan>,
585        optional: Arc<LogicalPlan>,
586        join_keys: Vec<(usize, usize)>,
587        inner_keep_idx: Vec<usize>,
588    ) -> Self {
589        let schema = Self::build_schema(&outer, &optional, &inner_keep_idx);
590        Self {
591            outer,
592            optional,
593            join_keys,
594            inner_keep_idx,
595            schema,
596        }
597    }
598
599    /// Build the output [`DFSchema`]: outer fields, then the kept optional
600    /// fields (by [`inner_keep_idx`](Self::inner_keep_idx)), each made nullable.
601    fn build_schema(
602        outer: &Arc<LogicalPlan>,
603        optional: &Arc<LogicalPlan>,
604        inner_keep_idx: &[usize],
605    ) -> DFSchemaRef {
606        let mut qualified: Vec<(Option<TableReference>, Arc<Field>)> = outer
607            .schema()
608            .iter()
609            .map(|(q, f)| (q.cloned(), Arc::clone(f)))
610            .collect();
611        let inner: Vec<(Option<TableReference>, Arc<Field>)> = optional
612            .schema()
613            .iter()
614            .map(|(q, f)| (q.cloned(), Arc::clone(f)))
615            .collect();
616        qualified.extend(inner_keep_idx.iter().map(|&i| {
617            let (q, f) = &inner[i];
618            // Unmatched outer rows null these columns, so they must be nullable
619            // regardless of the inner plan's nullability (idempotent).
620            let nullable = Arc::new(f.as_ref().clone().with_nullable(true));
621            (q.clone(), nullable)
622        }));
623        Arc::new(
624            DFSchema::new_with_metadata(qualified, std::collections::HashMap::new())
625                .expect("OptionalMatchNode schema must be constructible from outer + inner fields"),
626        )
627    }
628}
629
630impl_partial_ord!(OptionalMatchNode);
631
632impl UserDefinedLogicalNodeCore for OptionalMatchNode {
633    fn name(&self) -> &str {
634        "OptionalMatch"
635    }
636
637    fn inputs(&self) -> Vec<&LogicalPlan> {
638        vec![&self.outer, &self.optional]
639    }
640
641    fn schema(&self) -> &DFSchemaRef {
642        &self.schema
643    }
644
645    fn expressions(&self) -> Vec<Expr> {
646        vec![]
647    }
648
649    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
650        write!(f, "OptionalMatch: keys={}", self.join_keys.len())
651    }
652
653    fn with_exprs_and_inputs(
654        &self,
655        _exprs: Vec<Expr>,
656        mut inputs: Vec<LogicalPlan>,
657    ) -> DfResult<Self> {
658        // Only the inputs can change here; the join keys and the kept-inner
659        // column indices are preserved (the schema is rebuilt from them).
660        let optional = Arc::new(inputs.pop().unwrap_or_else(|| (*self.optional).clone()));
661        let outer = Arc::new(inputs.pop().unwrap_or_else(|| (*self.outer).clone()));
662        Ok(Self::new(
663            outer,
664            optional,
665            self.join_keys.clone(),
666            self.inner_keep_idx.clone(),
667        ))
668    }
669}
670
671// ---------------------------------------------------------------------------
672// PathUniqueNode
673// ---------------------------------------------------------------------------
674
675/// Logical stub for path-uniqueness filtering (eliminates paths that visit
676/// the same node or edge more than once).
677#[derive(Debug, Clone, PartialEq, Eq, Hash)]
678pub struct PathUniqueNode {
679    /// Input plan.
680    pub input: Arc<LogicalPlan>,
681    schema: DFSchemaRef,
682}
683
684impl PathUniqueNode {
685    /// Create a new path-uniqueness stub.
686    #[must_use]
687    pub fn new(input: Arc<LogicalPlan>) -> Self {
688        let schema = passthrough_schema(&[&input]);
689        Self { input, schema }
690    }
691}
692
693impl_partial_ord!(PathUniqueNode);
694
695impl UserDefinedLogicalNodeCore for PathUniqueNode {
696    fn name(&self) -> &str {
697        "PathUnique"
698    }
699
700    fn inputs(&self) -> Vec<&LogicalPlan> {
701        vec![&self.input]
702    }
703
704    fn schema(&self) -> &DFSchemaRef {
705        &self.schema
706    }
707
708    fn expressions(&self) -> Vec<Expr> {
709        vec![]
710    }
711
712    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
713        write!(f, "PathUnique")
714    }
715
716    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
717        let input = Arc::new(
718            inputs
719                .into_iter()
720                .next()
721                .unwrap_or_else(|| (*self.input).clone()),
722        );
723        Ok(Self::new(input))
724    }
725}
726
727// ---------------------------------------------------------------------------
728// OntologyInferNode
729// ---------------------------------------------------------------------------
730
731/// Logical stub for ontology-driven semantic inference (transitive /
732/// symmetric closure expansion).
733#[derive(Debug, Clone, PartialEq, Eq, Hash)]
734pub struct OntologyInferNode {
735    /// Input plan.
736    pub input: Arc<LogicalPlan>,
737    /// The relation type whose semantic flags trigger inference.
738    pub rel_type_name: String,
739    /// The ontology rule that triggered inference, e.g. `transitive:KNOWS` (#605).
740    pub rule_id: String,
741    /// The confidence policy for the derived facts, e.g. `conservative_min` (#605).
742    pub confidence_model: String,
743    schema: DFSchemaRef,
744}
745
746impl OntologyInferNode {
747    /// Create a new ontology-inference node carrying its rule (#605).
748    #[must_use]
749    pub fn new(
750        input: Arc<LogicalPlan>,
751        rel_type_name: impl Into<String>,
752        rule_id: impl Into<String>,
753        confidence_model: impl Into<String>,
754    ) -> Self {
755        let schema = passthrough_schema(&[&input]);
756        Self {
757            input,
758            rel_type_name: rel_type_name.into(),
759            rule_id: rule_id.into(),
760            confidence_model: confidence_model.into(),
761            schema,
762        }
763    }
764}
765
766impl_partial_ord!(OntologyInferNode);
767
768impl UserDefinedLogicalNodeCore for OntologyInferNode {
769    fn name(&self) -> &str {
770        "OntologyInfer"
771    }
772
773    fn inputs(&self) -> Vec<&LogicalPlan> {
774        vec![&self.input]
775    }
776
777    fn schema(&self) -> &DFSchemaRef {
778        &self.schema
779    }
780
781    fn expressions(&self) -> Vec<Expr> {
782        vec![]
783    }
784
785    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
786        write!(
787            f,
788            "OntologyInfer: rel={} rule_id={}",
789            self.rel_type_name, self.rule_id
790        )
791    }
792
793    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
794        let input = Arc::new(
795            inputs
796                .into_iter()
797                .next()
798                .unwrap_or_else(|| (*self.input).clone()),
799        );
800        Ok(Self::new(
801            input,
802            self.rel_type_name.clone(),
803            self.rule_id.clone(),
804            self.confidence_model.clone(),
805        ))
806    }
807}
808
809// ---------------------------------------------------------------------------
810// GraphMergeNode
811// ---------------------------------------------------------------------------
812
813/// Logical stub for `MERGE` (match-or-create write semantics).
814#[derive(Debug, Clone, PartialEq, Eq, Hash)]
815pub struct GraphMergeNode {
816    schema: DFSchemaRef,
817}
818
819impl GraphMergeNode {
820    /// Create a new graph-merge stub.
821    #[must_use]
822    pub fn new() -> Self {
823        Self {
824            schema: Arc::new(DFSchema::empty()),
825        }
826    }
827}
828
829impl_partial_ord!(GraphMergeNode);
830
831impl Default for GraphMergeNode {
832    fn default() -> Self {
833        Self::new()
834    }
835}
836
837impl UserDefinedLogicalNodeCore for GraphMergeNode {
838    fn name(&self) -> &str {
839        "GraphMerge"
840    }
841
842    fn inputs(&self) -> Vec<&LogicalPlan> {
843        vec![]
844    }
845
846    fn schema(&self) -> &DFSchemaRef {
847        &self.schema
848    }
849
850    fn expressions(&self) -> Vec<Expr> {
851        vec![]
852    }
853
854    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
855        write!(f, "GraphMerge")
856    }
857
858    fn with_exprs_and_inputs(
859        &self,
860        _exprs: Vec<Expr>,
861        _inputs: Vec<LogicalPlan>,
862    ) -> DfResult<Self> {
863        Ok(Self::new())
864    }
865}
866
867// ---------------------------------------------------------------------------
868// GraphCreateNode
869// ---------------------------------------------------------------------------
870
871/// A node to create, fully resolved (no expression-arena coupling).
872///
873/// Built by the relational lowering layer from a `graphforge_ir::CreateNodeSpec`:
874/// label IDs are paired with their resolved names and property maps are
875/// evaluated to literal key/value pairs, so the execution layer needs no
876/// access to the IR arena or ontology.
877#[derive(Debug, Clone, PartialEq)]
878pub struct ResolvedNodeSpec {
879    /// Pattern variable id (`VarId.0`).
880    pub var: u32,
881    /// Complete resolved label type-id set, in pattern order.
882    pub label_ids: Vec<u32>,
883    /// Complete resolved label-name set, in pattern order.
884    pub label_names: Vec<String>,
885    /// Literal property key/value pairs (constant or constant-folded values).
886    pub properties: Vec<(String, IrLiteral)>,
887    /// Row-dependent property values that could not be folded to a literal
888    /// (e.g. `{n: x}` from a driving `UNWIND`/`MATCH`): each is a DataFusion
889    /// `Expr` over the input columns, evaluated per minted row by the execution
890    /// layer (#814). Surfaced through the node's `expressions()` so the
891    /// optimizer rewrites the columns they reference.
892    pub computed_properties: Vec<(String, Expr)>,
893    /// `true` when this var was bound by a preceding `MATCH`/`WITH` — the
894    /// executor **references** the matched node (its identity comes from the
895    /// input row) rather than minting a new one (#703).
896    pub is_reference: bool,
897}
898
899/// An edge to create, fully resolved.  See [`ResolvedNodeSpec`].
900#[derive(Debug, Clone, PartialEq)]
901pub struct ResolvedEdgeSpec {
902    /// Pattern variable id of the edge.
903    pub var: u32,
904    /// Source node variable id.
905    pub src: u32,
906    /// Destination node variable id.
907    pub dst: u32,
908    /// Resolved relation type id, if any.
909    pub rel_type_id: Option<u32>,
910    /// Resolved relation type name (for file routing), if known.
911    pub rel_type_name: Option<String>,
912    /// Edge direction.
913    pub direction: Direction,
914    /// Literal property key/value pairs (constant or constant-folded values).
915    ///
916    /// Persisted as of #784: the execution layer writes these via
917    /// `GraphWriter::set_edge_properties` to `edge_properties/<REL_TYPE>.parquet`
918    /// (keyed by `edge_uuid`), and the read side joins them back. An edge with
919    /// properties but no relation type is rejected (the props would be
920    /// unreadable), so a non-empty map implies `rel_type_name` is `Some`.
921    pub properties: Vec<(String, IrLiteral)>,
922    /// Row-dependent property values evaluated per minted row (#814); see
923    /// [`ResolvedNodeSpec::computed_properties`].
924    pub computed_properties: Vec<(String, Expr)>,
925}
926
927/// Logical node for `CREATE`: a write specification driven by an input plan.
928///
929/// Carries the resolved node/edge specs plus the project directory and ontology
930/// mode needed to drive a writer.  The directory is baked in at lowering time
931/// because the physical-planning layer (an `ExtensionPlanner`) only sees the
932/// DataFusion session state, not the GraphForge project path.
933///
934/// # Input
935///
936/// The CREATE runs **once per input row** (#703): a standalone `CREATE (:X)`
937/// lowers over the implicit single unit row (so it creates exactly once), while
938/// a mixed `MATCH … CREATE …` runs the write once per matched row, referencing
939/// MATCH-bound vars' identities (`ResolvedNodeSpec::is_reference`) and minting
940/// the rest. The input columns are **consumed** (for cardinality + referenced
941/// identities), not projected — the node still emits only the write summary.
942///
943/// Output schema is a one-row write summary: `nodes_created` / `edges_created`
944/// (`UInt64`).
945#[derive(Debug, Clone)]
946pub struct GraphCreateNode {
947    /// Input plan whose rows drive the writes (the implicit unit row for a
948    /// standalone CREATE; the MATCH results for a mixed pipeline).
949    pub input: Arc<LogicalPlan>,
950    /// Nodes to create, in pattern order.
951    pub nodes: Vec<ResolvedNodeSpec>,
952    /// Edges to create, in pattern order.
953    pub edges: Vec<ResolvedEdgeSpec>,
954    /// Target project directory.
955    pub dir: PathBuf,
956    /// Ontology mode (drives writer edge / property routing).
957    pub mode: OntologyMode,
958    /// Output schema: the one-row write summary in summary mode, or the
959    /// created-entity row schema in emit-rows mode (#814).
960    schema: DFSchemaRef,
961    /// `true` when the node emits created-entity rows (write-result RETURN);
962    /// `false` for the terminal one-row summary.
963    emit_rows: bool,
964}
965
966impl GraphCreateNode {
967    /// Build the write-summary Arrow schema
968    /// (`nodes_created`, `edges_created`, `properties_set`, `labels_added`).
969    #[must_use]
970    pub fn summary_schema() -> Arc<Schema> {
971        Arc::new(Schema::new(vec![
972            Field::new("nodes_created", DataType::UInt64, false),
973            Field::new("edges_created", DataType::UInt64, false),
974            Field::new("properties_set", DataType::UInt64, false),
975            Field::new("labels_added", DataType::UInt64, false),
976        ]))
977    }
978
979    /// Create a new graph-create node over `input` in **summary** mode: the node
980    /// emits the one-row write summary (a terminal `CREATE`).
981    #[must_use]
982    pub fn new(
983        input: Arc<LogicalPlan>,
984        nodes: Vec<ResolvedNodeSpec>,
985        edges: Vec<ResolvedEdgeSpec>,
986        dir: PathBuf,
987        mode: OntologyMode,
988    ) -> Self {
989        let schema = Arc::new(
990            DFSchema::try_from(Self::summary_schema())
991                .expect("write-summary schema is always valid"),
992        );
993        Self {
994            input,
995            nodes,
996            edges,
997            dir,
998            mode,
999            schema,
1000            emit_rows: false,
1001        }
1002    }
1003
1004    /// Create a node in **emit-rows** mode (write-result RETURN, #814): instead
1005    /// of the summary, it emits one row per input row carrying the input columns
1006    /// plus the created variables' `var_<n>`-qualified identity + property columns
1007    /// (`output_schema`), so a following `RETURN`/`WITH` can read them. Side
1008    /// effects are reported out of band (the exec's tally), since a relation
1009    /// cannot also carry the summary.
1010    #[must_use]
1011    pub fn new_emitting(
1012        input: Arc<LogicalPlan>,
1013        nodes: Vec<ResolvedNodeSpec>,
1014        edges: Vec<ResolvedEdgeSpec>,
1015        dir: PathBuf,
1016        mode: OntologyMode,
1017        output_schema: DFSchemaRef,
1018    ) -> Self {
1019        Self {
1020            input,
1021            nodes,
1022            edges,
1023            dir,
1024            mode,
1025            schema: output_schema,
1026            emit_rows: true,
1027        }
1028    }
1029
1030    /// Whether this node emits created-entity rows (`true`) rather than the
1031    /// one-row write summary (`false`).
1032    #[must_use]
1033    pub fn emits_rows(&self) -> bool {
1034        self.emit_rows
1035    }
1036}
1037
1038impl_partial_ord!(GraphCreateNode);
1039
1040// `IrLiteral` is `PartialEq` (so we derive `PartialEq`) but not `Eq`/`Hash`
1041// (it carries `f64`).  `UserDefinedLogicalNodeCore` requires `Eq + Hash`, so we
1042// provide them by hand: `Eq` is the empty marker (literal property values are
1043// never NaN-reflexivity-sensitive in practice — matching DataFusion's own
1044// float-bearing nodes), and `Hash` normalises each literal (floats via
1045// `to_bits`).  `schema` is excluded from both (it is derived from the rest).
1046impl PartialEq for GraphCreateNode {
1047    fn eq(&self, other: &Self) -> bool {
1048        self.input == other.input
1049            && self.nodes == other.nodes
1050            && self.edges == other.edges
1051            && self.dir == other.dir
1052            && self.mode == other.mode
1053            && self.emit_rows == other.emit_rows
1054    }
1055}
1056
1057impl Eq for GraphCreateNode {}
1058
1059fn hash_literal<H: Hasher>(lit: &IrLiteral, state: &mut H) {
1060    match lit {
1061        IrLiteral::Null => 0u8.hash(state),
1062        IrLiteral::Bool(b) => {
1063            1u8.hash(state);
1064            b.hash(state);
1065        }
1066        IrLiteral::Int(i) => {
1067            2u8.hash(state);
1068            i.hash(state);
1069        }
1070        IrLiteral::Float(f) => {
1071            3u8.hash(state);
1072            f.to_bits().hash(state);
1073        }
1074        IrLiteral::Str(s) => {
1075            4u8.hash(state);
1076            s.hash(state);
1077        }
1078        IrLiteral::Uuid(uuid) => {
1079            14u8.hash(state);
1080            uuid.hash(state);
1081        }
1082        IrLiteral::Duration {
1083            months,
1084            days,
1085            seconds,
1086            nanos,
1087        } => {
1088            5u8.hash(state);
1089            months.hash(state);
1090            days.hash(state);
1091            seconds.hash(state);
1092            nanos.hash(state);
1093        }
1094        IrLiteral::DateTime(t) => {
1095            6u8.hash(state);
1096            t.hash(state);
1097        }
1098        IrLiteral::Date(d) => {
1099            7u8.hash(state);
1100            d.hash(state);
1101        }
1102        IrLiteral::LocalDateTime { days, nanos } => {
1103            8u8.hash(state);
1104            days.hash(state);
1105            nanos.hash(state);
1106        }
1107        IrLiteral::Time(n) => {
1108            9u8.hash(state);
1109            n.hash(state);
1110        }
1111        IrLiteral::ZonedTime { nanos, offset } => {
1112            10u8.hash(state);
1113            nanos.hash(state);
1114            offset.hash(state);
1115        }
1116        IrLiteral::ZonedDateTime {
1117            days,
1118            nanos,
1119            offset,
1120            zone,
1121        } => {
1122            11u8.hash(state);
1123            days.hash(state);
1124            nanos.hash(state);
1125            offset.hash(state);
1126            zone.hash(state);
1127        }
1128        IrLiteral::List(items) => {
1129            12u8.hash(state);
1130            items.len().hash(state);
1131            for it in items {
1132                hash_literal(it, state);
1133            }
1134        }
1135        IrLiteral::Map(entries) => {
1136            13u8.hash(state);
1137            entries.len().hash(state);
1138            for (key, value) in entries {
1139                key.hash(state);
1140                hash_literal(value, state);
1141            }
1142        }
1143    }
1144}
1145
1146fn hash_props<H: Hasher>(props: &[(String, IrLiteral)], state: &mut H) {
1147    props.len().hash(state);
1148    for (k, v) in props {
1149        k.hash(state);
1150        hash_literal(v, state);
1151    }
1152}
1153
1154// Computed-property values are `Expr` (not `Hash`); hash the count and keys and
1155// let `PartialEq` disambiguate the exprs (matching `GraphSetNode::hash`).
1156fn hash_computed<H: Hasher>(props: &[(String, Expr)], state: &mut H) {
1157    props.len().hash(state);
1158    for (k, _) in props {
1159        k.hash(state);
1160    }
1161}
1162
1163impl Hash for GraphCreateNode {
1164    fn hash<H: Hasher>(&self, state: &mut H) {
1165        self.input.hash(state);
1166        self.nodes.len().hash(state);
1167        for n in &self.nodes {
1168            n.var.hash(state);
1169            n.label_ids.hash(state);
1170            n.label_names.hash(state);
1171            n.is_reference.hash(state);
1172            hash_props(&n.properties, state);
1173            hash_computed(&n.computed_properties, state);
1174        }
1175        self.edges.len().hash(state);
1176        for e in &self.edges {
1177            e.var.hash(state);
1178            e.src.hash(state);
1179            e.dst.hash(state);
1180            e.rel_type_id.hash(state);
1181            e.rel_type_name.hash(state);
1182            e.direction.hash(state);
1183            hash_props(&e.properties, state);
1184            hash_computed(&e.computed_properties, state);
1185        }
1186        self.dir.hash(state);
1187        self.mode.hash(state);
1188        self.emit_rows.hash(state);
1189    }
1190}
1191
1192impl UserDefinedLogicalNodeCore for GraphCreateNode {
1193    fn name(&self) -> &str {
1194        "GraphCreate"
1195    }
1196
1197    fn inputs(&self) -> Vec<&LogicalPlan> {
1198        vec![&self.input]
1199    }
1200
1201    fn schema(&self) -> &DFSchemaRef {
1202        &self.schema
1203    }
1204
1205    fn expressions(&self) -> Vec<Expr> {
1206        // Surface every row-dependent property value expr (nodes then edges, in
1207        // spec order) so the optimizer sees the columns they reference and
1208        // round-trips them through `with_exprs_and_inputs`.
1209        self.nodes
1210            .iter()
1211            .flat_map(|n| n.computed_properties.iter())
1212            .chain(self.edges.iter().flat_map(|e| e.computed_properties.iter()))
1213            .map(|(_, e)| e.clone())
1214            .collect()
1215    }
1216
1217    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
1218        write!(
1219            f,
1220            "GraphCreate: nodes={}, edges={}",
1221            self.nodes.len(),
1222            self.edges.len()
1223        )
1224    }
1225
1226    fn with_exprs_and_inputs(&self, exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
1227        let input = Arc::new(
1228            inputs
1229                .into_iter()
1230                .next()
1231                .unwrap_or_else(|| (*self.input).clone()),
1232        );
1233        // Reinstall the computed-property exprs positionally in the same
1234        // nodes-then-edges order `expressions()` emits; a short/empty `exprs`
1235        // (e.g. a probe call) keeps the originals.
1236        let mut nodes = self.nodes.clone();
1237        let mut edges = self.edges.clone();
1238        let total: usize = nodes
1239            .iter()
1240            .map(|n| n.computed_properties.len())
1241            .sum::<usize>()
1242            + edges
1243                .iter()
1244                .map(|e| e.computed_properties.len())
1245                .sum::<usize>();
1246        if exprs.len() == total {
1247            let mut it = exprs.into_iter();
1248            for n in &mut nodes {
1249                for (_, e) in &mut n.computed_properties {
1250                    *e = it.next().expect("count checked above");
1251                }
1252            }
1253            for e in &mut edges {
1254                for (_, x) in &mut e.computed_properties {
1255                    *x = it.next().expect("count checked above");
1256                }
1257            }
1258        }
1259        // Preserve emit-rows mode + its output schema (a plain `new` would reset
1260        // to summary mode and drop the created-rows schema).
1261        if self.emit_rows {
1262            Ok(Self::new_emitting(
1263                input,
1264                nodes,
1265                edges,
1266                self.dir.clone(),
1267                self.mode,
1268                self.schema.clone(),
1269            ))
1270        } else {
1271            Ok(Self::new(input, nodes, edges, self.dir.clone(), self.mode))
1272        }
1273    }
1274}
1275
1276// ---------------------------------------------------------------------------
1277// GraphDeleteNode
1278// ---------------------------------------------------------------------------
1279
1280/// One resolved `DELETE` target: a bound variable and whether it is an edge.
1281///
1282/// The binder records only the `VarId`; the lowering layer resolves whether the
1283/// var is a node or an edge by inspecting the input schema (which identity
1284/// column — `node_uuid` vs `edge_uuid` — it carries), so the execution layer
1285/// reads the right identity column without re-deriving it.
1286#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1287pub struct DeleteTarget {
1288    /// Pattern variable id (`VarId.0`).
1289    pub var: u32,
1290    /// `true` if the variable is an edge (delete the edge row); `false` for a
1291    /// node (delete the node row, and incident edges when `detach`).
1292    pub is_edge: bool,
1293}
1294
1295/// Logical node for `DELETE` / `DETACH DELETE` (#740): a delete specification
1296/// driven by an input plan (the preceding `MATCH`).
1297///
1298/// Like [`GraphCreateNode`], the project directory and ontology mode are baked
1299/// in at lowering time (the `ExtensionPlanner` sees only the DataFusion session
1300/// state). The input rows supply the identities of the entities to delete; the
1301/// node emits a one-row write summary `nodes_deleted` / `edges_deleted`.
1302#[derive(Debug, Clone)]
1303pub struct GraphDeleteNode {
1304    /// Input plan whose rows carry the matched entities' identities.
1305    pub input: Arc<LogicalPlan>,
1306    /// The variables to delete, with their resolved node/edge kind.
1307    pub targets: Vec<DeleteTarget>,
1308    /// `true` for `DETACH DELETE` (also remove a node's incident edges).
1309    pub detach: bool,
1310    /// Target project directory.
1311    pub dir: PathBuf,
1312    /// Ontology mode (drives writer edge / property routing).
1313    pub mode: OntologyMode,
1314    schema: DFSchemaRef,
1315}
1316
1317impl GraphDeleteNode {
1318    /// Build the write-summary Arrow schema (`nodes_deleted`, `edges_deleted`).
1319    #[must_use]
1320    pub fn summary_schema() -> Arc<Schema> {
1321        Arc::new(Schema::new(vec![
1322            Field::new("nodes_deleted", DataType::UInt64, false),
1323            Field::new("edges_deleted", DataType::UInt64, false),
1324        ]))
1325    }
1326
1327    /// Create a new graph-delete node over `input`.
1328    #[must_use]
1329    pub fn new(
1330        input: Arc<LogicalPlan>,
1331        targets: Vec<DeleteTarget>,
1332        detach: bool,
1333        dir: PathBuf,
1334        mode: OntologyMode,
1335    ) -> Self {
1336        let schema = Arc::new(
1337            DFSchema::try_from(Self::summary_schema())
1338                .expect("write-summary schema is always valid"),
1339        );
1340        Self {
1341            input,
1342            targets,
1343            detach,
1344            dir,
1345            mode,
1346            schema,
1347        }
1348    }
1349}
1350
1351impl_partial_ord!(GraphDeleteNode);
1352
1353// `schema` is derived from the rest, so it is excluded from `PartialEq`/`Hash`
1354// (it is not part of the node's logical identity).
1355impl PartialEq for GraphDeleteNode {
1356    fn eq(&self, other: &Self) -> bool {
1357        self.input == other.input
1358            && self.targets == other.targets
1359            && self.detach == other.detach
1360            && self.dir == other.dir
1361            && self.mode == other.mode
1362    }
1363}
1364
1365impl Eq for GraphDeleteNode {}
1366
1367impl Hash for GraphDeleteNode {
1368    fn hash<H: Hasher>(&self, state: &mut H) {
1369        self.input.hash(state);
1370        self.targets.hash(state);
1371        self.detach.hash(state);
1372        self.dir.hash(state);
1373        self.mode.hash(state);
1374    }
1375}
1376
1377impl UserDefinedLogicalNodeCore for GraphDeleteNode {
1378    fn name(&self) -> &str {
1379        "GraphDelete"
1380    }
1381
1382    fn inputs(&self) -> Vec<&LogicalPlan> {
1383        vec![&self.input]
1384    }
1385
1386    fn schema(&self) -> &DFSchemaRef {
1387        &self.schema
1388    }
1389
1390    fn expressions(&self) -> Vec<Expr> {
1391        vec![]
1392    }
1393
1394    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
1395        write!(
1396            f,
1397            "GraphDelete: targets={}, detach={}",
1398            self.targets.len(),
1399            self.detach
1400        )
1401    }
1402
1403    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
1404        let input = Arc::new(
1405            inputs
1406                .into_iter()
1407                .next()
1408                .unwrap_or_else(|| (*self.input).clone()),
1409        );
1410        Ok(Self::new(
1411            input,
1412            self.targets.clone(),
1413            self.detach,
1414            self.dir.clone(),
1415            self.mode,
1416        ))
1417    }
1418}
1419
1420// ---------------------------------------------------------------------------
1421// GraphSetNode / GraphRemoveNode
1422// ---------------------------------------------------------------------------
1423
1424/// One resolved `SET n.prop = <expr>` target (#791).
1425///
1426/// Like [`DeleteTarget`], the node/edge kind (`is_edge`) is resolved by the
1427/// lowering layer from the input schema. `value` is the lowered DataFusion
1428/// expression for the assigned value — evaluated per matched row at execution
1429/// time (mirrors how [`UnwindNode`] carries its `list_expr`).
1430///
1431/// The property-file **stem** is resolved per row in the exec layer, not baked
1432/// here: a node uses its `type_id` (→ `_untyped` in Exploratory mode, else the
1433/// entity name from [`GraphSetNode::type_id_to_entity_name`]); an edge uses its
1434/// `rel_type_name` column (edge property files are keyed by relation name in
1435/// every mode).
1436#[derive(Debug, Clone, PartialEq)]
1437pub struct SetTarget {
1438    /// Pattern variable id (`VarId.0`).
1439    pub var: u32,
1440    /// `true` if the variable is an edge (write to `edge_properties/`); `false`
1441    /// for a node (write to `properties/`).
1442    pub is_edge: bool,
1443    /// The property name being assigned.
1444    pub prop_name: String,
1445    /// The value expression, evaluated per matched row in the exec layer.
1446    pub value: Expr,
1447}
1448
1449/// One resolved `REMOVE n.prop` target (#791) — the value-less dual of
1450/// [`SetTarget`].
1451#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1452pub struct RemoveTarget {
1453    /// Pattern variable id (`VarId.0`).
1454    pub var: u32,
1455    /// `true` if the variable is an edge.
1456    pub is_edge: bool,
1457    /// The property name being removed.
1458    pub prop_name: String,
1459}
1460
1461/// Logical node for `SET <prop> = <expr>` (#791): a property-write driven by an
1462/// input plan (the preceding `MATCH`).
1463///
1464/// Mirrors [`GraphDeleteNode`] — project directory and ontology mode are baked
1465/// in at lowering time. Each target's value expression is evaluated per matched
1466/// row by the physical layer, the resulting per-row literal written to the
1467/// entity's property file. Node-target file stems are resolved per row from the
1468/// row's `type_id` via [`type_id_to_entity_name`](Self::type_id_to_entity_name)
1469/// (handles an untyped `MATCH (n)` whose rows span several entity types). The
1470/// node emits a one-row summary `properties_set`.
1471#[derive(Debug, Clone)]
1472pub struct GraphSetNode {
1473    /// Input plan whose rows carry the matched entities' identities + columns.
1474    pub input: Arc<LogicalPlan>,
1475    /// The property assignments, with resolved node/edge kind + value expr.
1476    pub targets: Vec<SetTarget>,
1477    /// Maps a node `type_id` to its property-file entity stem, for per-row node
1478    /// stem resolution (an empty map / missing id falls back to `_untyped`).
1479    pub type_id_to_entity_name: HashMap<u32, String>,
1480    /// Target project directory.
1481    pub dir: PathBuf,
1482    /// Ontology mode (drives node property-file routing).
1483    pub mode: OntologyMode,
1484    schema: DFSchemaRef,
1485}
1486
1487impl GraphSetNode {
1488    /// Build the write-summary Arrow schema (`properties_set`).
1489    #[must_use]
1490    pub fn summary_schema() -> Arc<Schema> {
1491        Arc::new(Schema::new(vec![Field::new(
1492            "properties_set",
1493            DataType::UInt64,
1494            false,
1495        )]))
1496    }
1497
1498    /// Create a new graph-set node over `input`.
1499    #[must_use]
1500    pub fn new(
1501        input: Arc<LogicalPlan>,
1502        targets: Vec<SetTarget>,
1503        type_id_to_entity_name: HashMap<u32, String>,
1504        dir: PathBuf,
1505        mode: OntologyMode,
1506    ) -> Self {
1507        let schema = Arc::new(
1508            DFSchema::try_from(Self::summary_schema())
1509                .expect("write-summary schema is always valid"),
1510        );
1511        Self {
1512            input,
1513            targets,
1514            type_id_to_entity_name,
1515            dir,
1516            mode,
1517            schema,
1518        }
1519    }
1520}
1521
1522impl_partial_ord!(GraphSetNode);
1523
1524// `schema` is derived; excluded from `PartialEq`/`Hash` (not part of logical
1525// identity). The value exprs ARE part of identity, so they are included.
1526impl PartialEq for GraphSetNode {
1527    fn eq(&self, other: &Self) -> bool {
1528        self.input == other.input
1529            && self.targets == other.targets
1530            && self.dir == other.dir
1531            && self.mode == other.mode
1532            && self.type_id_to_entity_name == other.type_id_to_entity_name
1533    }
1534}
1535
1536impl Eq for GraphSetNode {}
1537
1538impl Hash for GraphSetNode {
1539    fn hash<H: Hasher>(&self, state: &mut H) {
1540        self.input.hash(state);
1541        // `SetTarget` is not `Hash` (it holds an `Expr`, which is not `Hash`);
1542        // hash the stable scalar fields and let `PartialEq` disambiguate the
1543        // value exprs. Plan-node hashing only needs to be consistent with `Eq`
1544        // for buckets, not collision-free.
1545        for t in &self.targets {
1546            t.var.hash(state);
1547            t.is_edge.hash(state);
1548            t.prop_name.hash(state);
1549        }
1550        self.dir.hash(state);
1551        self.mode.hash(state);
1552    }
1553}
1554
1555impl UserDefinedLogicalNodeCore for GraphSetNode {
1556    fn name(&self) -> &str {
1557        "GraphSet"
1558    }
1559
1560    fn inputs(&self) -> Vec<&LogicalPlan> {
1561        vec![&self.input]
1562    }
1563
1564    fn schema(&self) -> &DFSchemaRef {
1565        &self.schema
1566    }
1567
1568    fn expressions(&self) -> Vec<Expr> {
1569        // Surface each target's value expr so the optimizer sees the columns it
1570        // references (and round-trips them through `with_exprs_and_inputs`).
1571        self.targets.iter().map(|t| t.value.clone()).collect()
1572    }
1573
1574    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
1575        write!(f, "GraphSet: targets={}", self.targets.len())
1576    }
1577
1578    fn with_exprs_and_inputs(&self, exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
1579        let input = Arc::new(
1580            inputs
1581                .into_iter()
1582                .next()
1583                .unwrap_or_else(|| (*self.input).clone()),
1584        );
1585        // Reinstall value exprs positionally; an empty/short `exprs` (e.g. a
1586        // probe call) keeps the originals.
1587        let mut targets = self.targets.clone();
1588        if exprs.len() == targets.len() {
1589            for (t, e) in targets.iter_mut().zip(exprs) {
1590                t.value = e;
1591            }
1592        }
1593        Ok(Self::new(
1594            input,
1595            targets,
1596            self.type_id_to_entity_name.clone(),
1597            self.dir.clone(),
1598            self.mode,
1599        ))
1600    }
1601}
1602
1603/// Logical node for `REMOVE <prop>` (#791): the value-less dual of
1604/// [`GraphSetNode`]. Emits a one-row summary `properties_removed`.
1605#[derive(Debug, Clone)]
1606pub struct GraphRemoveNode {
1607    /// Input plan whose rows carry the matched entities' identities + columns.
1608    pub input: Arc<LogicalPlan>,
1609    /// The properties to remove, with resolved node/edge kind.
1610    pub targets: Vec<RemoveTarget>,
1611    /// Maps a node `type_id` to its property-file entity stem (see
1612    /// [`GraphSetNode::type_id_to_entity_name`]).
1613    pub type_id_to_entity_name: HashMap<u32, String>,
1614    /// Target project directory.
1615    pub dir: PathBuf,
1616    /// Ontology mode (drives node property-file routing).
1617    pub mode: OntologyMode,
1618    schema: DFSchemaRef,
1619}
1620
1621impl GraphRemoveNode {
1622    /// Build the write-summary Arrow schema (`properties_removed`).
1623    #[must_use]
1624    pub fn summary_schema() -> Arc<Schema> {
1625        Arc::new(Schema::new(vec![Field::new(
1626            "properties_removed",
1627            DataType::UInt64,
1628            false,
1629        )]))
1630    }
1631
1632    /// Create a new graph-remove node over `input`.
1633    #[must_use]
1634    pub fn new(
1635        input: Arc<LogicalPlan>,
1636        targets: Vec<RemoveTarget>,
1637        type_id_to_entity_name: HashMap<u32, String>,
1638        dir: PathBuf,
1639        mode: OntologyMode,
1640    ) -> Self {
1641        let schema = Arc::new(
1642            DFSchema::try_from(Self::summary_schema())
1643                .expect("write-summary schema is always valid"),
1644        );
1645        Self {
1646            input,
1647            targets,
1648            type_id_to_entity_name,
1649            dir,
1650            mode,
1651            schema,
1652        }
1653    }
1654}
1655
1656impl_partial_ord!(GraphRemoveNode);
1657
1658impl PartialEq for GraphRemoveNode {
1659    fn eq(&self, other: &Self) -> bool {
1660        self.input == other.input
1661            && self.targets == other.targets
1662            && self.dir == other.dir
1663            && self.mode == other.mode
1664            && self.type_id_to_entity_name == other.type_id_to_entity_name
1665    }
1666}
1667
1668impl Eq for GraphRemoveNode {}
1669
1670impl Hash for GraphRemoveNode {
1671    fn hash<H: Hasher>(&self, state: &mut H) {
1672        self.input.hash(state);
1673        self.targets.hash(state);
1674        self.dir.hash(state);
1675        self.mode.hash(state);
1676    }
1677}
1678
1679impl UserDefinedLogicalNodeCore for GraphRemoveNode {
1680    fn name(&self) -> &str {
1681        "GraphRemove"
1682    }
1683
1684    fn inputs(&self) -> Vec<&LogicalPlan> {
1685        vec![&self.input]
1686    }
1687
1688    fn schema(&self) -> &DFSchemaRef {
1689        &self.schema
1690    }
1691
1692    fn expressions(&self) -> Vec<Expr> {
1693        vec![]
1694    }
1695
1696    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
1697        write!(f, "GraphRemove: targets={}", self.targets.len())
1698    }
1699
1700    fn with_exprs_and_inputs(&self, _exprs: Vec<Expr>, inputs: Vec<LogicalPlan>) -> DfResult<Self> {
1701        let input = Arc::new(
1702            inputs
1703                .into_iter()
1704                .next()
1705                .unwrap_or_else(|| (*self.input).clone()),
1706        );
1707        Ok(Self::new(
1708            input,
1709            self.targets.clone(),
1710            self.type_id_to_entity_name.clone(),
1711            self.dir.clone(),
1712            self.mode,
1713        ))
1714    }
1715}
1716
1717// ---------------------------------------------------------------------------
1718// UnwindNode
1719// ---------------------------------------------------------------------------
1720
1721/// Physical node for `UNWIND` — explodes a list expression into one row per
1722/// element, binding each element to an alias variable.
1723///
1724/// The physical layer ([`UnwindExec`](../graphforge_exec) in M13) evaluates `list_expr`
1725/// per input row and emits one output row per list element (null/empty list →
1726/// zero rows). The output extends the input schema with a single `alias`-named,
1727/// `alias`-qualified column of the element type (nullable — elements may be
1728/// null).
1729#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1730pub struct UnwindNode {
1731    /// Input plan.
1732    pub input: Arc<LogicalPlan>,
1733    /// The list expression to iterate.
1734    pub list_expr: Expr,
1735    /// The alias variable qualifier bound to each element (e.g. `var_0`).
1736    pub alias: String,
1737    element_field: Arc<Field>,
1738    schema: DFSchemaRef,
1739}
1740
1741impl UnwindNode {
1742    /// Create an unwind node.
1743    ///
1744    /// `element_field` is the unwound element's column (its type comes from the
1745    /// list element type, resolved by the lowerer so this crate needs no
1746    /// `graphforge-rel` dependency). It is appended to the input schema qualified by
1747    /// `alias` and forced nullable (UNWIND elements may be null).
1748    #[must_use]
1749    pub fn new(
1750        input: Arc<LogicalPlan>,
1751        list_expr: Expr,
1752        alias: impl Into<String>,
1753        element_field: &Field,
1754    ) -> Self {
1755        let alias = alias.into();
1756        // The element column is named after the alias and forced nullable
1757        // (UNWIND elements may be null); build_schema appends it unqualified so a
1758        // bare `RETURN x` (lowered to `col(alias)`) resolves it.
1759        let element_field = Arc::new(element_field.clone().with_name(&alias).with_nullable(true));
1760        let schema = Self::build_schema(&input, &alias, Arc::clone(&element_field));
1761        Self {
1762            input,
1763            list_expr,
1764            alias,
1765            element_field,
1766            schema,
1767        }
1768    }
1769
1770    /// Build the output [`DFSchema`]: the input's qualified fields followed by
1771    /// the element column (already named after the alias, nullable), appended
1772    /// **unqualified** so a bare `UNWIND … AS x RETURN x` — which lowers `x` to
1773    /// `col(alias)`, an unqualified `Column` — resolves it.
1774    fn build_schema(
1775        input: &Arc<LogicalPlan>,
1776        alias: &str,
1777        element_field: Arc<Field>,
1778    ) -> DFSchemaRef {
1779        let mut qualified: Vec<(Option<TableReference>, Arc<Field>)> = input
1780            .schema()
1781            .iter()
1782            .map(|(q, f)| (q.cloned(), Arc::clone(f)))
1783            .collect();
1784        if let DataType::Struct(fields) = element_field.data_type()
1785            && fields.iter().any(|field| {
1786                matches!(
1787                    field.name().as_str(),
1788                    "node_uuid" | "edge_uuid" | "src_uuid" | "dst_uuid" | "nodes" | "relationships"
1789                )
1790            })
1791        {
1792            let alias_ref = TableReference::bare(alias.to_owned());
1793            qualified.extend(fields.iter().map(|field| {
1794                (
1795                    Some(alias_ref.clone()),
1796                    Arc::new(field.as_ref().clone().with_nullable(true)),
1797                )
1798            }));
1799        } else {
1800            qualified.push((None, element_field));
1801        }
1802        Arc::new(
1803            DFSchema::new_with_metadata(qualified, std::collections::HashMap::new())
1804                .expect("UnwindNode schema must include the element column"),
1805        )
1806    }
1807
1808    /// Recover the appended element field from the output schema (last column),
1809    /// for reconstructing the node in [`with_exprs_and_inputs`].
1810    fn element_field(&self) -> Arc<Field> {
1811        Arc::clone(&self.element_field)
1812    }
1813
1814    fn bound_element_field(&self, list_expr: &Expr, input: &LogicalPlan) -> Arc<Field> {
1815        match list_expr.get_type(input.schema().as_ref()) {
1816            Ok(
1817                DataType::List(field)
1818                | DataType::LargeList(field)
1819                | DataType::FixedSizeList(field, _),
1820            ) => Arc::new(
1821                field
1822                    .as_ref()
1823                    .clone()
1824                    .with_name(&self.alias)
1825                    .with_nullable(true),
1826            ),
1827            _ => self.element_field(),
1828        }
1829    }
1830}
1831
1832impl_partial_ord!(UnwindNode);
1833
1834impl UserDefinedLogicalNodeCore for UnwindNode {
1835    fn name(&self) -> &str {
1836        "Unwind"
1837    }
1838
1839    fn inputs(&self) -> Vec<&LogicalPlan> {
1840        vec![&self.input]
1841    }
1842
1843    fn schema(&self) -> &DFSchemaRef {
1844        &self.schema
1845    }
1846
1847    fn expressions(&self) -> Vec<Expr> {
1848        vec![self.list_expr.clone()]
1849    }
1850
1851    fn fmt_for_explain(&self, f: &mut fmt::Formatter) -> fmt::Result {
1852        write!(f, "Unwind: alias={}", self.alias)
1853    }
1854
1855    fn with_exprs_and_inputs(
1856        &self,
1857        mut exprs: Vec<Expr>,
1858        inputs: Vec<LogicalPlan>,
1859    ) -> DfResult<Self> {
1860        let list_expr = exprs.pop().unwrap_or_else(|| self.list_expr.clone());
1861        let input = Arc::new(
1862            inputs
1863                .into_iter()
1864                .next()
1865                .unwrap_or_else(|| (*self.input).clone()),
1866        );
1867        let element_field = self.bound_element_field(&list_expr, &input);
1868        Ok(Self::new(
1869            input,
1870            list_expr,
1871            self.alias.clone(),
1872            &element_field,
1873        ))
1874    }
1875}
1876
1877// ---------------------------------------------------------------------------
1878// Prevent-predicate-push-down override for all stubs
1879// ---------------------------------------------------------------------------
1880
1881// All stubs inherit the default (block all push-down) which is safe.
1882
1883// ---------------------------------------------------------------------------
1884// Tests
1885// ---------------------------------------------------------------------------
1886
1887#[cfg(test)]
1888mod tests {
1889    use super::*;
1890    use datafusion::logical_expr::{Extension, LogicalPlan};
1891
1892    fn empty_plan() -> Arc<LogicalPlan> {
1893        use datafusion::logical_expr::LogicalPlanBuilder;
1894        Arc::new(LogicalPlanBuilder::empty(false).build().unwrap())
1895    }
1896
1897    fn wrap(node: impl UserDefinedLogicalNodeCore) -> LogicalPlan {
1898        LogicalPlan::Extension(Extension {
1899            node: Arc::new(node),
1900        })
1901    }
1902
1903    /// Two representative destination-node fields (mirrors the storage layer's
1904    /// `TOPOLOGY_NODES_SCHEMA`, kept minimal for tests).
1905    fn dst_node_fields() -> Vec<Arc<Field>> {
1906        vec![
1907            Arc::new(Field::new("node_id", DataType::UInt64, false)),
1908            Arc::new(Field::new("type_id", DataType::UInt32, false)),
1909        ]
1910    }
1911
1912    /// Build a `VarLenExpandNode` with test defaults (src=0, dst=1, edge=2, Out).
1913    fn var_len(rel: &str, min_hops: u16, max_hops: Option<u16>) -> VarLenExpandNode {
1914        VarLenExpandNode::new(
1915            empty_plan(),
1916            rel,
1917            min_hops,
1918            max_hops,
1919            0,
1920            1,
1921            2,
1922            Direction::Out,
1923            Some(7),
1924            PathBuf::from("/tmp/gf"),
1925            OntologyMode::Strict,
1926            dst_node_fields(),
1927            var_len_edge_list_field(&[]),
1928        )
1929    }
1930
1931    /// A `LogicalPlan` scan over `fields` qualified `var_<var>`, for tests that
1932    /// need an optional sub-plan with real columns.
1933    fn scan_plan(var: u32, fields: Vec<Field>) -> Arc<LogicalPlan> {
1934        use datafusion::logical_expr::LogicalPlanBuilder;
1935        use datafusion::logical_expr::logical_plan::LogicalTableSource;
1936        let schema = Arc::new(Schema::new(fields));
1937        Arc::new(
1938            LogicalPlanBuilder::scan(
1939                format!("var_{var}"),
1940                Arc::new(LogicalTableSource::new(schema)),
1941                None,
1942            )
1943            .unwrap()
1944            .build()
1945            .unwrap(),
1946        )
1947    }
1948
1949    /// Build an `OptionalMatchNode` over an empty outer and the given optional
1950    /// sub-plan, keeping the inner columns named by `inner_keep_idx`, for tests.
1951    fn opt_match(
1952        optional: Arc<LogicalPlan>,
1953        join_keys: Vec<(usize, usize)>,
1954        inner_keep_idx: Vec<usize>,
1955    ) -> OptionalMatchNode {
1956        OptionalMatchNode::new(empty_plan(), optional, join_keys, inner_keep_idx)
1957    }
1958
1959    /// Build an `UnwindNode` with the given list expr + alias, defaulting the
1960    /// element column to a nullable Int64 named `elem`.
1961    fn unwind(list_expr: datafusion::logical_expr::Expr, alias: &str) -> UnwindNode {
1962        let element_field = Field::new("elem", DataType::Int64, true);
1963        UnwindNode::new(empty_plan(), list_expr, alias, &element_field)
1964    }
1965
1966    #[test]
1967    fn var_len_expand_name() {
1968        let n = var_len("KNOWS", 1, Some(3));
1969        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "VarLenExpand");
1970        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
1971    }
1972
1973    #[test]
1974    fn var_len_expand_schema_extends_input_with_dst_and_edge_fields() {
1975        use datafusion::common::TableReference;
1976        // empty input has 0 fields; output = 0 + 2 dst fields + 1 edge-list col.
1977        let n = var_len("KNOWS", 1, Some(2));
1978        let schema = UserDefinedLogicalNodeCore::schema(&n);
1979        assert_eq!(schema.fields().len(), 3);
1980        // Destination columns are qualified `var_1` and resolvable by name.
1981        let dst = TableReference::bare("var_1");
1982        assert!(schema.field_with_qualified_name(&dst, "node_id").is_ok());
1983        assert!(schema.field_with_qualified_name(&dst, "type_id").is_ok());
1984        // The edge-list column is qualified `var_2` (edge_var) and is a List.
1985        let edge = TableReference::bare("var_2");
1986        let edge_field = schema
1987            .field_with_qualified_name(&edge, VAR_LEN_EDGE_LIST_FIELD)
1988            .expect("edge-list column present");
1989        assert!(matches!(edge_field.data_type(), DataType::List(_)));
1990    }
1991
1992    #[test]
1993    fn optional_match_name() {
1994        let n = opt_match(empty_plan(), vec![], vec![]);
1995        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "OptionalMatch");
1996        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 2);
1997    }
1998
1999    #[test]
2000    fn optional_match_schema_appends_nullable_inner_fields() {
2001        // empty outer (0 fields) + an optional sub-plan with two `var_1` columns,
2002        // both kept, no join keys → 2 nullable output cols.
2003        let optional = scan_plan(
2004            1,
2005            vec![
2006                Field::new("node_id", DataType::UInt64, false),
2007                Field::new("type_id", DataType::UInt32, false),
2008            ],
2009        );
2010        let n = opt_match(optional, vec![], vec![0, 1]);
2011        let schema = UserDefinedLogicalNodeCore::schema(&n);
2012        assert_eq!(schema.fields().len(), 2);
2013        // Inner columns are made nullable for null-shaping.
2014        assert!(schema.field(0).is_nullable());
2015        assert!(schema.field(1).is_nullable());
2016        let var1 = TableReference::bare("var_1");
2017        assert!(schema.field_with_qualified_name(&var1, "node_id").is_ok());
2018    }
2019
2020    #[test]
2021    fn optional_match_schema_excludes_shared_var_columns() {
2022        // Regression (#718): a shared variable's columns are carried by the outer
2023        // side, so `inner_keep_idx` excludes them — only the genuinely new inner
2024        // column (here index 2) is appended, avoiding duplicate `var_0` fields.
2025        let optional = scan_plan(
2026            7,
2027            vec![
2028                Field::new("node_id", DataType::UInt64, false),
2029                Field::new("type_id", DataType::UInt32, false),
2030                Field::new("payload", DataType::UInt64, false),
2031            ],
2032        );
2033        let n = opt_match(optional, vec![(0, 0)], vec![2]);
2034        let schema = UserDefinedLogicalNodeCore::schema(&n);
2035        assert_eq!(schema.fields().len(), 1, "only the non-shared col is kept");
2036        assert_eq!(schema.field(0).name(), "payload");
2037        assert!(schema.field(0).is_nullable());
2038    }
2039
2040    #[test]
2041    fn path_unique_name() {
2042        let n = PathUniqueNode::new(empty_plan());
2043        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "PathUnique");
2044        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
2045    }
2046
2047    #[test]
2048    fn ontology_infer_name() {
2049        let n = OntologyInferNode::new(
2050            empty_plan(),
2051            "MANAGES",
2052            "transitive:MANAGES",
2053            "conservative_min",
2054        );
2055        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "OntologyInfer");
2056        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
2057    }
2058
2059    #[test]
2060    fn graph_merge_name() {
2061        let n = GraphMergeNode::new();
2062        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "GraphMerge");
2063        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 0);
2064    }
2065
2066    #[test]
2067    fn unwind_name() {
2068        use datafusion::logical_expr::lit;
2069        let n = unwind(lit(1i64), "x");
2070        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "Unwind");
2071        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
2072        assert_eq!(UserDefinedLogicalNodeCore::expressions(&n).len(), 1);
2073    }
2074
2075    #[test]
2076    fn unwind_schema_appends_element_column() {
2077        use datafusion::logical_expr::lit;
2078        // empty input (0 fields) + 1 element column → 1 nullable col named `x`.
2079        let n = unwind(lit(1i64), "x");
2080        let schema = UserDefinedLogicalNodeCore::schema(&n);
2081        assert_eq!(schema.fields().len(), 1);
2082        assert!(
2083            schema.field(0).is_nullable(),
2084            "unwound element column is nullable"
2085        );
2086        // The element column is named after the alias and UNQUALIFIED, so a bare
2087        // `RETURN x` (lowered to `col("x")`) resolves it.
2088        assert_eq!(schema.field(0).name(), "x");
2089        assert!(schema.field_with_unqualified_name("x").is_ok());
2090    }
2091
2092    #[test]
2093    fn graph_create_name_and_wrap() {
2094        let node = GraphCreateNode::new(
2095            empty_plan(),
2096            vec![ResolvedNodeSpec {
2097                var: 0,
2098                label_ids: vec![3],
2099                label_names: vec!["Person".to_owned()],
2100                properties: vec![("name".to_owned(), IrLiteral::Str("Alice".to_owned()))],
2101                computed_properties: vec![],
2102                is_reference: false,
2103            }],
2104            vec![],
2105            PathBuf::from("/tmp/gf"),
2106            OntologyMode::Strict,
2107        );
2108        assert_eq!(UserDefinedLogicalNodeCore::name(&node), "GraphCreate");
2109        // Input-driven now (one CREATE per input row).
2110        assert_eq!(UserDefinedLogicalNodeCore::inputs(&node).len(), 1);
2111        // Summary schema: nodes_created + edges_created + properties_set + labels_added.
2112        assert_eq!(UserDefinedLogicalNodeCore::schema(&node).fields().len(), 4);
2113        // Wraps as an Extension that DataFusion accepts.
2114        assert!(matches!(wrap(node), LogicalPlan::Extension(_)));
2115    }
2116
2117    #[test]
2118    fn graph_create_eq_and_hash_are_consistent() {
2119        use std::collections::hash_map::DefaultHasher;
2120
2121        let mk = || {
2122            GraphCreateNode::new(
2123                empty_plan(),
2124                vec![ResolvedNodeSpec {
2125                    var: 0,
2126                    label_ids: vec![],
2127                    label_names: vec![],
2128                    properties: vec![("score".to_owned(), IrLiteral::Float(1.5))],
2129                    computed_properties: vec![],
2130                    is_reference: false,
2131                }],
2132                vec![],
2133                PathBuf::from("/tmp/gf"),
2134                OntologyMode::Exploratory,
2135            )
2136        };
2137        let (a, b) = (mk(), mk());
2138        assert_eq!(a, b);
2139        let mut ha = DefaultHasher::new();
2140        let mut hb = DefaultHasher::new();
2141        a.hash(&mut ha);
2142        b.hash(&mut hb);
2143        assert_eq!(ha.finish(), hb.finish());
2144    }
2145
2146    fn set_node(value: Expr) -> GraphSetNode {
2147        GraphSetNode::new(
2148            empty_plan(),
2149            vec![SetTarget {
2150                var: 0,
2151                is_edge: false,
2152                prop_name: "age".to_owned(),
2153                value,
2154            }],
2155            HashMap::new(),
2156            PathBuf::from("/tmp/gf"),
2157            OntologyMode::Exploratory,
2158        )
2159    }
2160
2161    fn remove_node() -> GraphRemoveNode {
2162        GraphRemoveNode::new(
2163            empty_plan(),
2164            vec![RemoveTarget {
2165                var: 0,
2166                is_edge: false,
2167                prop_name: "age".to_owned(),
2168            }],
2169            HashMap::new(),
2170            PathBuf::from("/tmp/gf"),
2171            OntologyMode::Exploratory,
2172        )
2173    }
2174
2175    #[test]
2176    fn graph_set_name_schema_and_expr_surface() {
2177        use datafusion::logical_expr::lit;
2178        let n = set_node(lit(42i64));
2179        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "GraphSet");
2180        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
2181        // Summary schema: properties_set.
2182        assert_eq!(UserDefinedLogicalNodeCore::schema(&n).fields().len(), 1);
2183        // The value expr is surfaced so the optimizer sees its referenced cols.
2184        assert_eq!(UserDefinedLogicalNodeCore::expressions(&n).len(), 1);
2185        assert!(matches!(wrap(n), LogicalPlan::Extension(_)));
2186    }
2187
2188    #[test]
2189    fn graph_set_round_trips_value_expr() {
2190        use datafusion::logical_expr::{col, lit};
2191        // A runtime value expr (`var_0.age + 1`) must survive the expr round-trip
2192        // through `with_exprs_and_inputs` (used by the optimizer).
2193        let value = col("var_0.age") + lit(1i64);
2194        let n = set_node(value.clone());
2195        let exprs = UserDefinedLogicalNodeCore::expressions(&n);
2196        assert_eq!(exprs, vec![value.clone()]);
2197        let rebuilt = UserDefinedLogicalNodeCore::with_exprs_and_inputs(&n, exprs, vec![]).unwrap();
2198        assert_eq!(rebuilt.targets[0].value, value);
2199        assert_eq!(n, rebuilt);
2200    }
2201
2202    #[test]
2203    fn graph_set_eq_distinguishes_value_expr() {
2204        use datafusion::logical_expr::lit;
2205        // Two SET nodes differing only in the value expr must NOT compare equal
2206        // (the value is part of logical identity).
2207        assert_ne!(set_node(lit(1i64)), set_node(lit(2i64)));
2208    }
2209
2210    #[test]
2211    fn graph_remove_name_schema_and_no_exprs() {
2212        let n = remove_node();
2213        assert_eq!(UserDefinedLogicalNodeCore::name(&n), "GraphRemove");
2214        assert_eq!(UserDefinedLogicalNodeCore::inputs(&n).len(), 1);
2215        assert_eq!(UserDefinedLogicalNodeCore::schema(&n).fields().len(), 1);
2216        // REMOVE carries no value expression.
2217        assert!(UserDefinedLogicalNodeCore::expressions(&n).is_empty());
2218        assert!(matches!(wrap(n), LogicalPlan::Extension(_)));
2219    }
2220
2221    #[test]
2222    fn graph_remove_eq_and_hash_are_consistent() {
2223        use std::collections::hash_map::DefaultHasher;
2224        let (a, b) = (remove_node(), remove_node());
2225        assert_eq!(a, b);
2226        let mut ha = DefaultHasher::new();
2227        let mut hb = DefaultHasher::new();
2228        a.hash(&mut ha);
2229        b.hash(&mut hb);
2230        assert_eq!(ha.finish(), hb.finish());
2231    }
2232
2233    #[test]
2234    fn all_stubs_wrap_as_extension() {
2235        use datafusion::logical_expr::lit;
2236        let plans = vec![
2237            wrap(var_len("*", 1, None)),
2238            wrap(opt_match(empty_plan(), vec![], vec![])),
2239            wrap(PathUniqueNode::new(empty_plan())),
2240            wrap(OntologyInferNode::new(
2241                empty_plan(),
2242                "REL",
2243                "transitive:REL",
2244                "conservative_min",
2245            )),
2246            wrap(GraphMergeNode::new()),
2247            wrap(unwind(lit(1i64), "x")),
2248            wrap(set_node(lit(1i64))),
2249            wrap(remove_node()),
2250        ];
2251        for p in &plans {
2252            assert!(
2253                matches!(p, LogicalPlan::Extension(_)),
2254                "expected Extension, got {p:?}"
2255            );
2256        }
2257    }
2258
2259    #[test]
2260    fn fmt_for_explain_is_non_empty() {
2261        use datafusion::logical_expr::lit;
2262
2263        // Thin Display adapter that calls fmt_for_explain directly.
2264        struct ExplainWrapper<'a, T: UserDefinedLogicalNodeCore>(&'a T);
2265        impl<T: UserDefinedLogicalNodeCore> fmt::Display for ExplainWrapper<'_, T> {
2266            fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2267                self.0.fmt_for_explain(f)
2268            }
2269        }
2270        fn explain<T: UserDefinedLogicalNodeCore>(n: &T) -> String {
2271            ExplainWrapper(n).to_string()
2272        }
2273
2274        assert!(!explain(&var_len("KNOWS", 1, Some(3))).is_empty());
2275        assert!(!explain(&opt_match(empty_plan(), vec![], vec![])).is_empty());
2276        assert!(!explain(&PathUniqueNode::new(empty_plan())).is_empty());
2277        assert!(
2278            !explain(&OntologyInferNode::new(
2279                empty_plan(),
2280                "MANAGES",
2281                "transitive:MANAGES",
2282                "conservative_min"
2283            ))
2284            .is_empty()
2285        );
2286        assert!(!explain(&GraphMergeNode::new()).is_empty());
2287        assert!(!explain(&unwind(lit(1i64), "x")).is_empty());
2288
2289        // Spot-check actual explain output content.
2290        assert!(explain(&var_len("KNOWS", 1, Some(3))).contains("KNOWS"));
2291        assert!(
2292            explain(&OntologyInferNode::new(
2293                empty_plan(),
2294                "MANAGES",
2295                "transitive:MANAGES",
2296                "conservative_min"
2297            ))
2298            .contains("MANAGES")
2299        );
2300        assert!(explain(&unwind(lit(1i64), "x")).contains("x"));
2301    }
2302
2303    // -----------------------------------------------------------------------
2304    // ExpandNode (#763)
2305    // -----------------------------------------------------------------------
2306
2307    fn expand_node(edge_prop_fields: Vec<Arc<Field>>) -> ExpandNode {
2308        let edge_fields = vec![
2309            Arc::new(Field::new(
2310                "edge_uuid",
2311                DataType::FixedSizeBinary(16),
2312                false,
2313            )),
2314            Arc::new(Field::new("edge_id", DataType::UInt64, false)),
2315            Arc::new(Field::new("src_id", DataType::UInt64, false)),
2316            Arc::new(Field::new("dst_id", DataType::UInt64, false)),
2317        ];
2318        ExpandNode::new(
2319            empty_plan(),
2320            "KNOWS",
2321            0,
2322            2,
2323            1,
2324            Direction::Out,
2325            Some(7),
2326            PathBuf::from("/tmp/p"),
2327            OntologyMode::Strict,
2328            edge_fields,
2329            edge_prop_fields,
2330            dst_node_fields(),
2331        )
2332    }
2333
2334    #[test]
2335    fn expand_node_schema_order_qualifiers_and_prop_nullability() {
2336        use datafusion::common::TableReference;
2337
2338        // A non-nullable property field on disk must come out NULLABLE
2339        // (LEFT-join parity with the join path).
2340        let node = expand_node(vec![Arc::new(Field::new("since", DataType::Int64, false))]);
2341        let schema = UserDefinedLogicalNodeCore::schema(&node);
2342        // Order: (empty input) ++ var_1 edge topology ++ var_1 props ++ var_2 dst.
2343        let names: Vec<String> = schema
2344            .iter()
2345            .map(|(q, f)| {
2346                format!(
2347                    "{}.{}",
2348                    q.map(TableReference::table).unwrap_or("?"),
2349                    f.name()
2350                )
2351            })
2352            .collect();
2353        assert_eq!(
2354            names,
2355            [
2356                "var_1.edge_uuid",
2357                "var_1.edge_id",
2358                "var_1.src_id",
2359                "var_1.dst_id",
2360                "var_1.since",
2361                "var_2.node_id",
2362                "var_2.type_id",
2363            ]
2364        );
2365        let edge = TableReference::bare("var_1");
2366        let since = schema.field_with_qualified_name(&edge, "since").unwrap();
2367        assert!(since.is_nullable(), "prop columns are LEFT-join nullable");
2368        let uuid = schema
2369            .field_with_qualified_name(&edge, "edge_uuid")
2370            .unwrap();
2371        assert!(!uuid.is_nullable(), "topology columns keep non-null");
2372        assert_eq!(node.edge_prop_count, 1);
2373    }
2374
2375    #[test]
2376    fn expand_node_with_exprs_and_inputs_round_trips() {
2377        let node = expand_node(vec![Arc::new(Field::new("since", DataType::Int64, true))]);
2378        let rebuilt = UserDefinedLogicalNodeCore::with_exprs_and_inputs(
2379            &node,
2380            vec![],
2381            vec![(*empty_plan()).clone()],
2382        )
2383        .unwrap();
2384        assert_eq!(
2385            UserDefinedLogicalNodeCore::schema(&rebuilt).as_arrow(),
2386            UserDefinedLogicalNodeCore::schema(&node).as_arrow(),
2387            "schema survives optimizer-style input replacement"
2388        );
2389        assert_eq!(rebuilt.edge_prop_count, 1);
2390        assert_eq!(rebuilt.rel_type_name, "KNOWS");
2391    }
2392}