Skip to main content

core_rules/
engine.rs

1use crate::def::{evaluate, is_keymatch_rooted, NodeView, Predicate, RuleDef};
2use crate::index::{
3    candidate_spec, candidate_spec_approx_with_k, ivf_drift_rebuild_threshold, CandidateSpec,
4    RuleIndex,
5};
6use core_storage::{ColumnStore, EdgeProps, IdMap, Interner, Topology, Value};
7use std::collections::{BTreeMap, BTreeSet};
8
9/// A single derived-edge fire or retract captured during a commit.
10///
11/// Populated inside [`ProvSets::insert`] / [`ProvSets::remove`] while graph
12/// state is fully intact (before any tombstone step in `DeleteNode`).
13/// String keys are resolved at capture time so they remain valid even after
14/// node deletion.
15#[derive(Debug, Clone)]
16pub struct EngineEdgeDelta {
17    pub rule: String,
18    /// User-facing source node key.
19    pub src_key: String,
20    /// User-facing destination node key.
21    pub dst_key: String,
22    /// Edge-type string.
23    pub edge_type: String,
24    /// Internal edge-type symbol (for edge_props weight lookup in db.rs).
25    pub etype_sym: u32,
26    /// Internal source node id (for edge_props weight lookup in db.rs).
27    pub src_id: u32,
28    /// Internal destination node id (for edge_props weight lookup in db.rs).
29    pub dst_id: u32,
30    /// `true` = edge was fired (added to provenance); `false` = retracted.
31    pub fired: bool,
32}
33
34#[cfg(test)]
35pub use crate::index::{with_ivf_drift_rebuild, with_vector_dim_reject, with_vector_early_exit};
36
37/// Borrowed mutable view of graph state the engine writes derived edges into.
38pub struct GraphMut<'a> {
39    pub ids: &'a IdMap,
40    pub syms: &'a mut Interner,
41    pub labels: &'a [u32],
42    pub props: &'a ColumnStore,
43    pub topo: &'a mut Topology,
44    pub edge_props: &'a mut EdgeProps,
45}
46
47/// Used when `RuleDef.max_edges` is `None`.
48pub const DEFAULT_MAX_EDGES: u64 = 1_000_000;
49
50/// `(etype, src, dst)` as stored in `provenance` / `owned`.
51type Triple = (u32, u32, u32);
52/// Reverse-index entry: `(rule_id, etype, src, dst)`. `rule_id` is interned.
53type Touch = (u32, u32, u32, u32);
54
55/// IVF state for one index side, exported for V4 snapshot persistence:
56/// `(centroids, node→cluster assignments, drift_counter)`.
57pub type SideIvfExport = (Vec<Vec<f64>>, BTreeMap<u32, usize>, u64);
58/// IVF state for both sides (src, dst) of one approximate rule.
59pub type RuleIvfExport = (SideIvfExport, SideIvfExport);
60
61#[derive(Debug, Default)]
62pub struct RuleEngine {
63    rules: BTreeMap<String, RuleDef>,
64    indexes: BTreeMap<String, RuleIndex>,
65    provenance: BTreeMap<String, BTreeSet<Triple>>,
66    owned: BTreeSet<Triple>,
67    /// Derived reverse index: node → provenance triples that touch it.
68    /// Never serialized; rebuilt from `provenance` on persist-restore.
69    by_node: BTreeMap<u32, BTreeSet<Touch>>,
70    /// Intern table for rule names used as `Touch` rule_ids. Derived.
71    /// Additive-only: ids are never reused. `by_node` stores these ids, so
72    /// pruning-and-reusing a slot would alias leftover touches to a new name.
73    /// Bound: one slot per distinct rule name ever created in this process.
74    rule_intern: BTreeMap<String, u32>,
75    intern_rule: Vec<String>,
76    tripped: BTreeMap<String, bool>,
77    fires: BTreeMap<String, u64>,
78    /// Staging buffer for post-commit [`EngineEdgeDelta`] events.
79    ///
80    /// Populated by [`ProvSets::insert`] / [`ProvSets::remove`] during
81    /// `apply` (live writes AND WAL replay). Callers must drain via
82    /// [`RuleEngine::drain_deltas`] immediately after apply to consume live
83    /// events or discard replay noise.
84    pending_deltas: Vec<EngineEdgeDelta>,
85    /// Gate: whether to accumulate [`EngineEdgeDelta`] items during rule
86    /// application.
87    ///
88    /// **Safety invariant:** events are fire-and-forget live streams — a
89    /// subscriber that attaches *later* never receives past events by design.
90    /// Similarly, views call `backfill_view` at creation time (reading directly
91    /// from `topo`, not from pending deltas), so deltas accumulated before a
92    /// view is defined are not needed. Accumulation can therefore be skipped
93    /// whenever no subscriber and no view exists; the observable behaviour is
94    /// identical. Set to `true` by `set_emit_deltas` before the first
95    /// subscribe, create_view, or any operation that needs events; cleared when
96    /// the last listener is removed.
97    emit_deltas: bool,
98    /// Approximate rule names whose dst-side IVF drift exceeded
99    /// [`crate::IVF_DRIFT_REBUILD`] during the last index maintenance.
100    /// Drained by [`RuleEngine::take_rebuild_needed`] after apply.
101    rebuild_needed: BTreeSet<String>,
102}
103
104// ---------------------------------------------------------------------------
105// Private helpers (free functions, not methods, to avoid whole-struct borrows)
106// ---------------------------------------------------------------------------
107
108/// Rule-aware candidate spec: exact `ScanAll` for `approximate=false`, IVF
109/// `VectorClusters` for `approximate=true` (VectorSimilar-rooted predicates).
110fn candidate_spec_for(def: &RuleDef) -> CandidateSpec<'_> {
111    if def.approximate {
112        // k = max(max_edges, 64): return at least 64 candidates so evaluation
113        // has a meaningful pool; bounded by max_edges when set by the caller.
114        let k = def.max_edges.map(|me| me.max(64)).unwrap_or(64) as usize;
115        candidate_spec_approx_with_k(&def.predicate, k)
116    } else {
117        candidate_spec(&def.predicate)
118    }
119}
120
121/// Rule-aware src-side lookup spec. KeyMatch is still exact on the src side
122/// regardless of `approximate` (the approximation is on the dst candidate set).
123/// For KeyMatch-rooted predicates, src side is indexed as Scalar (FK field
124/// value → node bucket) so reverse lookup uses the dst key. Non-KeyMatch
125/// `All` uses the full [`candidate_spec_for`] Intersect (not `parts[0]`).
126fn src_lookup_spec_for(def: &RuleDef) -> CandidateSpec<'_> {
127    if is_keymatch_rooted(&def.predicate) {
128        let field =
129            keymatch_field(&def.predicate).expect("keymatch-rooted predicate has a KeyMatch field");
130        CandidateSpec::Scalar { field }
131    } else {
132        candidate_spec_for(def)
133    }
134}
135
136/// True when `p` contains a `VectorSimilar { field: f }` where `f == field`.
137fn predicate_covers_field(p: &Predicate, field: &str) -> bool {
138    match p {
139        Predicate::VectorSimilar { field: f, .. } => f == field,
140        Predicate::All(parts) | Predicate::Any(parts) => {
141            parts.iter().any(|q| predicate_covers_field(q, field))
142        }
143        _ => false,
144    }
145}
146
147/// Extract the FK field name from a KeyMatch (or All-leading-KeyMatch) predicate.
148fn keymatch_field(p: &Predicate) -> Option<&str> {
149    match p {
150        Predicate::KeyMatch { field } => Some(field),
151        Predicate::All(parts) => parts.first().and_then(keymatch_field),
152        Predicate::Any(_) => None,
153        _ => None,
154    }
155}
156
157/// Compute the set of desired (src, dst) → score edges involving node `n` on
158/// the given side.  Returns an empty map if `n`'s label doesn't match the rule.
159fn compute_desired(
160    def: &RuleDef,
161    index: &RuleIndex,
162    n: u32,
163    on_src_side: bool,
164    g: &GraphMut<'_>,
165) -> BTreeMap<(u32, u32), f64> {
166    let (my_label, other_label) = if on_src_side {
167        (&def.src_label, &def.dst_label)
168    } else {
169        (&def.dst_label, &def.src_label)
170    };
171
172    let Some(my_sym) = g.syms.get(my_label) else {
173        return BTreeMap::new();
174    };
175    if g.labels.get(n as usize).copied() != Some(my_sym) {
176        return BTreeMap::new();
177    }
178    let other_sym = g.syms.get(other_label);
179
180    let n_key = match g.ids.key_of(n) {
181        Some(k) => k,
182        None => return BTreeMap::new(),
183    };
184    let n_get = |f: &str| g.props.get(n, f).cloned();
185
186    let spec = candidate_spec_for(def);
187    let candidates: BTreeSet<u32> = if on_src_side {
188        if is_keymatch_rooted(&def.predicate) {
189            // KeyMatch src→dst: look up the dst node directly by FK field value.
190            // Covers `ByKey` and `All` whose first conjunct is KeyMatch
191            // (`Intersect([ByKey, …])`); evaluate() filters remaining conjuncts.
192            let field = keymatch_field(&def.predicate).expect("ByKey always comes from KeyMatch");
193            match n_get(field) {
194                Some(Value::Str(ref target_key)) => match g.ids.get(target_key) {
195                    Some(dst_id) => std::iter::once(dst_id).collect(),
196                    None => BTreeSet::new(),
197                },
198                _ => BTreeSet::new(),
199            }
200        } else {
201            index.dst_side.candidates(&spec, &n_get)
202        }
203    } else {
204        // n is dst: probe src_side to find src candidates.
205        let src_spec = src_lookup_spec_for(def);
206        if is_keymatch_rooted(&def.predicate) {
207            // Synthetic getter: returns n's key for the FK field so we find
208            // src nodes whose FK value points to n.
209            let key_getter = |_: &str| Some(Value::Str(n_key.to_string()));
210            index.src_side.candidates(&src_spec, &key_getter)
211        } else {
212            index.src_side.candidates(&src_spec, &n_get)
213        }
214    };
215
216    // Fast path: Cauchy-Schwarz suffix-norm early exit for exact VectorSimilar.
217    //
218    // Skipped for approximate rules (`def.approximate == true`): the IVF
219    // pre-filter already eliminates non-candidate nodes, and ScanAll metadata
220    // (vec_meta / vec_checkpoints) is not maintained for VectorClusters specs.
221    //
222    // Pre-fetch n's live vector ONCE outside the candidate loop so it is
223    // allocated only once per compute_desired call (not per candidate pair).
224    // m's vector is still fetched per pair — unavoidable without caching full
225    // vectors (which is the O(n·dim) trade-off the brief rules out).
226    //
227    // Freshness gate: `SideIndex::fresh_ckpts_for` returns `None` when the
228    // cached norm differs from the live norm, preventing stale checkpoints from
229    // producing a false reject (see doc comment on `fresh_ckpts_for`).
230    let n_early_exit_hint: Option<(Vec<f64>, f64, [f64; 8])> = if !def.approximate {
231        if let Predicate::VectorSimilar { field, .. } = &def.predicate {
232            if crate::index::vector_early_exit_enabled() {
233                let n_side = if on_src_side {
234                    &index.src_side
235                } else {
236                    &index.dst_side
237                };
238                if let Some(vn_v) = n_get(field) {
239                    if let Some(vn) = crate::index::as_numeric_list(&vn_v) {
240                        if let Some((norm_n, ckpts_n)) = n_side.fresh_ckpts_for(n, &vn) {
241                            Some((vn, norm_n, *ckpts_n))
242                        } else {
243                            None
244                        }
245                    } else {
246                        None
247                    }
248                } else {
249                    None
250                }
251            } else {
252                None
253            }
254        } else {
255            None
256        }
257    } else {
258        None
259    };
260
261    let mut out = BTreeMap::new();
262    for m in candidates {
263        if m == n {
264            continue; // never self-edges
265        }
266        if g.labels.get(m as usize).copied() != other_sym {
267            continue; // label filter
268        }
269        let m_key = match g.ids.key_of(m) {
270            Some(k) => k,
271            None => continue,
272        };
273        let m_get = |f: &str| g.props.get(m, f).cloned();
274        let (s_view, d_view, s_id, d_id) = if on_src_side {
275            (
276                NodeView {
277                    key: n_key,
278                    props: &n_get,
279                },
280                NodeView {
281                    key: m_key,
282                    props: &m_get,
283                },
284                n,
285                m,
286            )
287        } else {
288            (
289                NodeView {
290                    key: m_key,
291                    props: &m_get,
292                },
293                NodeView {
294                    key: n_key,
295                    props: &n_get,
296                },
297                m,
298                n,
299            )
300        };
301
302        // Use the pre-fetched n hint if available; fetch m per-pair.
303        if let (Some((ref vn, norm_n, ckpts_n)), Predicate::VectorSimilar { field, min }) =
304            (&n_early_exit_hint, &def.predicate)
305        {
306            let m_side = if on_src_side {
307                &index.dst_side
308            } else {
309                &index.src_side
310            };
311            if let Some(vm_v) = m_get(field) {
312                if let Some(vm) = crate::index::as_numeric_list(&vm_v) {
313                    if let Some((norm_m, ckpts_m)) = m_side.fresh_ckpts_for(m, &vm) {
314                        let (va, ckpts_a, na, vb, ckpts_b, nb) = if on_src_side {
315                            (
316                                vn.as_slice(),
317                                ckpts_n,
318                                *norm_n,
319                                vm.as_slice(),
320                                ckpts_m,
321                                norm_m,
322                            )
323                        } else {
324                            (
325                                vm.as_slice(),
326                                ckpts_m,
327                                norm_m,
328                                vn.as_slice(),
329                                ckpts_n,
330                                *norm_n,
331                            )
332                        };
333                        match crate::def::cosine_early_exit(va, vb, ckpts_a, ckpts_b, na, nb, *min)
334                        {
335                            None => continue, // exact reject
336                            Some(score) => {
337                                out.insert((s_id, d_id), score);
338                                continue; // full cosine already computed
339                            }
340                        }
341                    }
342                }
343            }
344        }
345
346        if let Some(score) = evaluate(&def.predicate, &s_view, &d_view) {
347            out.insert((s_id, d_id), score);
348        }
349    }
350    out
351}
352
353/// Compute the desired `(src, dst) → score` map for a **via-hop rule** where
354/// `def.via_label.is_some()`.
355///
356/// Semantics: src -[via_edge/via_dir]→ via(via_label), evaluate predicate
357/// between via and dst; fire src→dst if any via satisfies; score = max over via.
358///
359/// `anchor` selects which side of the computation to anchor:
360/// - `ViaAnchor::Src(src_id)`: expand from one specific src node.
361/// - `ViaAnchor::Dst(dst_id)`: n is a dst node; scan all src nodes and check
362///   if any via hop to them evaluates with n.
363///
364/// Always returns `(src, dst)` keyed pairs regardless of anchor.
365fn compute_desired_via(
366    def: &RuleDef,
367    anchor: ViaAnchor,
368    g: &GraphMut<'_>,
369) -> BTreeMap<(u32, u32), f64> {
370    let via_label = def.via_label.as_deref().unwrap();
371    let via_edge_str = def.via_edge.as_deref().unwrap();
372    let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
373
374    let src_sym = match g.syms.get(&def.src_label) {
375        Some(s) => s,
376        None => return BTreeMap::new(),
377    };
378    let via_sym = match g.syms.get(via_label) {
379        Some(s) => s,
380        None => return BTreeMap::new(),
381    };
382    let dst_sym = match g.syms.get(&def.dst_label) {
383        Some(s) => s,
384        None => return BTreeMap::new(),
385    };
386    let via_etype = match g.syms.get(via_edge_str) {
387        Some(e) => e,
388        None => return BTreeMap::new(),
389    };
390
391    // Determine which src ids to iterate over.
392    let srcs: Vec<u32> = match anchor {
393        ViaAnchor::Src(src_id) => {
394            if g.labels.get(src_id as usize).copied() == Some(src_sym) {
395                vec![src_id]
396            } else {
397                return BTreeMap::new();
398            }
399        }
400        ViaAnchor::Dst(_) => {
401            // Scan all src-label nodes.
402            (0..g.ids.len() as u32)
403                .filter(|&id| {
404                    matches!(
405                        g.labels.get(id as usize).copied(),
406                        Some(s) if s != u32::MAX && s == src_sym
407                    )
408                })
409                .collect()
410        }
411    };
412
413    // Collect dst candidates: all dst-label nodes (or just the anchored dst).
414    let anchored_dst: Option<u32> = match anchor {
415        ViaAnchor::Dst(dst_id) => {
416            if g.labels.get(dst_id as usize).copied() == Some(dst_sym) {
417                Some(dst_id)
418            } else {
419                return BTreeMap::new();
420            }
421        }
422        _ => None,
423    };
424
425    let mut out = BTreeMap::new();
426
427    for src in srcs {
428        let _src_key = match g.ids.key_of(src) {
429            Some(k) => k,
430            None => continue,
431        };
432        // Expand via hops from src.
433        let via_neighbors: Vec<u32> = g
434            .topo
435            .neighbors(via_etype, via_dir, src)
436            .iter()
437            .copied()
438            .filter(|&v| g.labels.get(v as usize).copied() == Some(via_sym))
439            .collect();
440
441        if via_neighbors.is_empty() {
442            continue;
443        }
444
445        // Collect dsts to evaluate: all dst-label nodes or just anchored one.
446        let dsts: Vec<u32> = if let Some(dst_id) = anchored_dst {
447            vec![dst_id]
448        } else {
449            (0..g.ids.len() as u32)
450                .filter(|&id| {
451                    id != src
452                        && matches!(
453                            g.labels.get(id as usize).copied(),
454                            Some(s) if s != u32::MAX && s == dst_sym
455                        )
456                })
457                .collect()
458        };
459
460        for dst in dsts {
461            if dst == src {
462                continue; // no self-edges
463            }
464            let dst_key = match g.ids.key_of(dst) {
465                Some(k) => k,
466                None => continue,
467            };
468            let dst_get = |f: &str| g.props.get(dst, f).cloned();
469            let dst_view = NodeView {
470                key: dst_key,
471                props: &dst_get,
472            };
473
474            // Score = max over via nodes that satisfy predicate(via, dst).
475            let mut best: Option<f64> = None;
476            for &via_id in &via_neighbors {
477                let via_key = match g.ids.key_of(via_id) {
478                    Some(k) => k,
479                    None => continue,
480                };
481                let via_get = |f: &str| g.props.get(via_id, f).cloned();
482                let via_view = NodeView {
483                    key: via_key,
484                    props: &via_get,
485                };
486                if let Some(score) = evaluate(&def.predicate, &via_view, &dst_view) {
487                    best = Some(match best {
488                        None => score,
489                        Some(prev) => prev.max(score),
490                    });
491                }
492            }
493
494            if let Some(score) = best {
495                out.insert((src, dst), score);
496            }
497        }
498    }
499
500    out
501}
502
503/// Anchor point for `compute_desired_via`.
504enum ViaAnchor {
505    /// Expand from one src node (src prop change or src insert).
506    Src(u32),
507    /// Re-evaluate all srcs that can reach some via satisfying predicate with
508    /// this dst (dst prop change).
509    Dst(u32),
510}
511
512fn edge_budget(def: &RuleDef) -> u64 {
513    // Only applies when max_edges is None (global-budget path).
514    // Some(k) rules use per-source top-k semantics, not this budget.
515    def.max_edges.unwrap_or(DEFAULT_MAX_EDGES)
516}
517
518/// Filter a per-source candidate map to the top-k destinations.
519///
520/// `per_src` must contain only pairs with the same source node (all
521/// `(src, dst)` keys share the same `src`).  Returns the top-`k` subset
522/// ordered by **(score DESC, dst_key ASC)** — higher scores win; ties are
523/// broken by the destination node's string key in ascending lexicographic
524/// order, giving a deterministic result independent of internal node IDs.
525///
526/// When `k` equals or exceeds the number of candidates, the input is
527/// returned unchanged (no allocation).
528///
529/// # Memory cost (per-source candidate ordering)
530///
531/// This function sorts and truncates a `Vec<((u32,u32), f64)>` of length
532/// equal to the number of matching candidates for one source.  That is
533/// O(M) per call, where M is the candidate count for this source.  Across
534/// a backfill sweep the peak additional memory is O(M_max) — the largest
535/// per-source candidate set — not the global total, because the Vec is
536/// dropped after each source.  No persistent per-source ordering is
537/// maintained beyond the materialized top-k provenance; backfill and
538/// rebuild recompute the ordering on demand from the live candidate index.
539pub(crate) fn filter_src_top_k(
540    per_src: BTreeMap<(u32, u32), f64>,
541    k: u64,
542    ids: &core_storage::IdMap,
543) -> BTreeMap<(u32, u32), f64> {
544    if per_src.len() as u64 <= k {
545        return per_src;
546    }
547    let mut candidates: Vec<((u32, u32), f64)> = per_src.into_iter().collect();
548    // Sort: score DESC (higher = better), then dst_key ASC as tiebreak.
549    candidates.sort_by(|&((_, da), sa), &((_, db), sb)| {
550        sb.total_cmp(&sa).then_with(|| {
551            let ka = ids.key_of(da).unwrap_or("");
552            let kb = ids.key_of(db).unwrap_or("");
553            ka.cmp(kb)
554        })
555    });
556    candidates.truncate(k as usize);
557    candidates.into_iter().collect()
558}
559
560/// Apply top-k derived-edge semantics for a single source node.
561///
562/// Retracts `(src, *)` provenance edges not in `desired_from_src`, then
563/// adds / refreshes weights for those that are.  Does **not** use the
564/// global tripped latch or budget check — top-k rules (`max_edges: Some(k)`)
565/// are self-capping by construction.
566fn apply_per_src_top_k(
567    def: &RuleDef,
568    src: u32,
569    desired_from_src: BTreeMap<(u32, u32), f64>,
570    prov: &mut ProvSets<'_>,
571    g: &mut GraphMut<'_>,
572) {
573    let et = g.syms.intern(&def.edge_type);
574
575    // Collect current (src, *) provenance triples for this rule.
576    // We filter to s == src so that (*, src) triples — where src is a dst
577    // for some other source — are not mistakenly retracted.
578    let current: Vec<Triple> = {
579        let rid = prov.rule_intern.get(&def.name).copied();
580        prov.by_node
581            .get(&src)
582            .into_iter()
583            .flatten()
584            .filter(|(r, t, s, _d)| Some(*r) == rid && *t == et && *s == src)
585            .map(|(_, t, s, d)| (*t, *s, *d))
586            .collect()
587    };
588
589    // Retract (src, dst) pairs no longer in the top-k.
590    for (t, s, d) in current {
591        if !desired_from_src.contains_key(&(s, d)) {
592            g.topo.remove_edge(t, s, d);
593            g.edge_props.remove_edge(t, s, d);
594            prov.remove(&def.name, (t, s, d), g.ids, g.syms);
595        }
596    }
597
598    // Insert new top-k pairs; refresh weights on already-owned pairs.
599    for ((s, d), score) in &desired_from_src {
600        let triple = (et, *s, *d);
601        let already = prov.contains(&triple);
602        if !already {
603            let newly = g.topo.add_edge(et, *s, *d);
604            if newly {
605                prov.insert(&def.name, triple, g.ids, g.syms);
606            }
607        }
608        let is_owned = already || prov.contains(&triple);
609        if is_owned {
610            if let Some(p) = &def.weight_prop {
611                g.edge_props.set(et, *s, *d, p, Value::Float(*score));
612            }
613        }
614    }
615}
616
617/// Never recycles ids. See `RuleEngine::rule_intern` for why.
618fn intern_rule(intern: &mut BTreeMap<String, u32>, names: &mut Vec<String>, rule: &str) -> u32 {
619    if let Some(&id) = intern.get(rule) {
620        return id;
621    }
622    let id = names.len() as u32;
623    intern.insert(rule.to_string(), id);
624    names.push(rule.to_string());
625    id
626}
627
628type ByNodeRebuild = (
629    BTreeMap<u32, BTreeSet<Touch>>,
630    BTreeMap<String, u32>,
631    Vec<String>,
632);
633
634fn rebuild_by_node(provenance: &BTreeMap<String, BTreeSet<Triple>>) -> ByNodeRebuild {
635    let mut by_node = BTreeMap::new();
636    let mut intern = BTreeMap::new();
637    let mut names = Vec::new();
638    for (rule, set) in provenance {
639        let rid = intern_rule(&mut intern, &mut names, rule);
640        for &triple in set {
641            touch_insert(&mut by_node, rid, triple);
642        }
643    }
644    (by_node, intern, names)
645}
646
647fn touch_insert(by_node: &mut BTreeMap<u32, BTreeSet<Touch>>, rid: u32, triple: Triple) {
648    let (t, s, d) = triple;
649    let entry = (rid, t, s, d);
650    by_node.entry(s).or_default().insert(entry);
651    if s != d {
652        by_node.entry(d).or_default().insert(entry);
653    }
654}
655
656fn touch_remove(by_node: &mut BTreeMap<u32, BTreeSet<Touch>>, rid: u32, triple: Triple) {
657    let (t, s, d) = triple;
658    let entry = (rid, t, s, d);
659    if let Some(set) = by_node.get_mut(&s) {
660        set.remove(&entry);
661        if set.is_empty() {
662            by_node.remove(&s);
663        }
664    }
665    if s != d {
666        if let Some(set) = by_node.get_mut(&d) {
667            set.remove(&entry);
668            if set.is_empty() {
669                by_node.remove(&d);
670            }
671        }
672    }
673}
674
675#[cfg(test)]
676fn resolve_by_node(
677    by_node: &BTreeMap<u32, BTreeSet<Touch>>,
678    names: &[String],
679) -> BTreeMap<u32, BTreeSet<(String, Triple)>> {
680    by_node
681        .iter()
682        .map(|(&n, set)| {
683            let resolved = set
684                .iter()
685                .map(|&(rid, t, s, d)| (names[rid as usize].clone(), (t, s, d)))
686                .collect();
687            (n, resolved)
688        })
689        .collect()
690}
691
692/// Mutable provenance + derived reverse index. Every insert/remove goes
693/// through [`ProvSets::insert`] / [`ProvSets::remove`].
694struct ProvSets<'a> {
695    set: &'a mut BTreeSet<Triple>,
696    owned: &'a mut BTreeSet<Triple>,
697    by_node: &'a mut BTreeMap<u32, BTreeSet<Touch>>,
698    rule_intern: &'a mut BTreeMap<String, u32>,
699    intern_rule: &'a mut Vec<String>,
700    /// Staging buffer for post-commit events. Keys are resolved at capture
701    /// time (before any tombstone step) so the strings remain valid after
702    /// node deletion.
703    deltas: &'a mut Vec<EngineEdgeDelta>,
704    /// Mirror of [`RuleEngine::emit_deltas`]: when `false`, pushes to
705    /// `deltas` are skipped entirely (no heap allocation, no String clone).
706    emit: bool,
707}
708
709impl ProvSets<'_> {
710    /// `ids` and `syms` are passed by the caller (not stored in ProvSets) to
711    /// avoid a conflicting borrow when callers also need `&mut g.syms` for
712    /// `intern` calls in the same function body.
713    fn insert(&mut self, rule: &str, triple: Triple, ids: &IdMap, syms: &Interner) -> bool {
714        if !self.set.insert(triple) {
715            return false;
716        }
717        self.owned.insert(triple);
718        let rid = intern_rule(self.rule_intern, self.intern_rule, rule);
719        touch_insert(self.by_node, rid, triple);
720        let (etype, src, dst) = triple;
721        if self.emit {
722            if let (Some(sk), Some(dk), Some(et)) =
723                (ids.key_of(src), ids.key_of(dst), syms.resolve(etype))
724            {
725                self.deltas.push(EngineEdgeDelta {
726                    rule: rule.to_string(),
727                    src_key: sk.to_string(),
728                    dst_key: dk.to_string(),
729                    edge_type: et.to_string(),
730                    etype_sym: etype,
731                    src_id: src,
732                    dst_id: dst,
733                    fired: true,
734                });
735            }
736        }
737        true
738    }
739
740    fn remove(&mut self, rule: &str, triple: Triple, ids: &IdMap, syms: &Interner) -> bool {
741        if !self.set.remove(&triple) {
742            return false;
743        }
744        self.owned.remove(&triple);
745        let rid = intern_rule(self.rule_intern, self.intern_rule, rule);
746        touch_remove(self.by_node, rid, triple);
747        let (etype, src, dst) = triple;
748        if self.emit {
749            if let (Some(sk), Some(dk), Some(et)) =
750                (ids.key_of(src), ids.key_of(dst), syms.resolve(etype))
751            {
752                self.deltas.push(EngineEdgeDelta {
753                    rule: rule.to_string(),
754                    src_key: sk.to_string(),
755                    dst_key: dk.to_string(),
756                    edge_type: et.to_string(),
757                    etype_sym: etype,
758                    src_id: src,
759                    dst_id: dst,
760                    fired: false,
761                });
762            }
763        }
764        true
765    }
766
767    fn contains(&self, triple: &Triple) -> bool {
768        self.set.contains(triple)
769    }
770
771    fn len(&self) -> usize {
772        self.set.len()
773    }
774}
775
776/// Diff-apply `desired` against provenance. `retract_touching = Some(n)` only
777/// retracts triples that involve `n` (incremental fire). `None` retracts any
778/// current provenance triple not in `desired` (backfill / rebuild).
779///
780/// `tripped` is a one-way latch: once set, no new provenance edges are added
781/// (gate on the flag itself, not `prov.len()`), even if retracts have brought
782/// the set below budget. Retracts and weight refreshes on already-owned edges
783/// still run. Crossing the budget on a not-yet-tripped rule sets the latch
784/// and skips that add and every later add in this call. Never an error.
785/// First-N (pre-trip) is BTree `(src, dst)` order of `desired` after retract.
786fn apply_desired(
787    def: &RuleDef,
788    desired: BTreeMap<(u32, u32), f64>,
789    retract_touching: Option<u32>,
790    prov: &mut ProvSets<'_>,
791    tripped: &mut bool,
792    g: &mut GraphMut<'_>,
793) {
794    let budget = edge_budget(def);
795    let et = g.syms.intern(&def.edge_type);
796
797    let current: Vec<Triple> = match retract_touching {
798        None => prov
799            .set
800            .iter()
801            .filter(|(t, _, _)| *t == et)
802            .copied()
803            .collect(),
804        Some(n) => {
805            let rid = prov.rule_intern.get(&def.name).copied();
806            prov.by_node
807                .get(&n)
808                .into_iter()
809                .flatten()
810                .filter(|(r, t, _, _)| Some(*r) == rid && *t == et)
811                .map(|(_, t, s, d)| (*t, *s, *d))
812                .collect()
813        }
814    };
815
816    for (t, s, d) in current {
817        if !desired.contains_key(&(s, d)) {
818            g.topo.remove_edge(t, s, d);
819            g.edge_props.remove_edge(t, s, d);
820            prov.remove(&def.name, (t, s, d), g.ids, g.syms);
821        }
822    }
823
824    for ((s, d), score) in desired {
825        let triple = (et, s, d);
826        let already = prov.contains(&triple);
827        if !already {
828            if *tripped || prov.len() as u64 >= budget {
829                *tripped = true;
830                continue;
831            }
832            let newly = g.topo.add_edge(et, s, d);
833            if newly {
834                prov.insert(&def.name, triple, g.ids, g.syms);
835            }
836        }
837        // Only set weight_prop on edges this rule owns (newly added now, or
838        // already in provenance). Pre-existing user edges are never owned, so
839        // writing a weight to them would leave a ghost property after deletion.
840        let is_owned_here = already || prov.contains(&triple);
841        if is_owned_here {
842            if let Some(p) = &def.weight_prop {
843                g.edge_props.set(et, s, d, p, Value::Float(score));
844            }
845        }
846    }
847}
848
849/// Union of `compute_desired(..., as_src)` over every live src-label node.
850/// BTree iteration of the result is the engine's deterministic first-N order
851/// for backfill and rebuild.
852///
853/// Kept for test reference comparators only.  Production paths use the
854/// streaming variants (`apply_streaming_create`, `apply_streaming_rebuild`)
855/// which never materialise the global map.
856#[cfg(test)]
857#[allow(dead_code)]
858fn compute_full_desired(
859    def: &RuleDef,
860    index: &RuleIndex,
861    g: &GraphMut<'_>,
862) -> BTreeMap<(u32, u32), f64> {
863    let mut desired = BTreeMap::new();
864    let src_sym = g.syms.get(&def.src_label);
865    for id in 0..g.ids.len() as u32 {
866        let label_sym = match g.labels.get(id as usize).copied() {
867            Some(s) if s != u32::MAX => s,
868            _ => continue,
869        };
870        if src_sym == Some(label_sym) {
871            desired.extend(compute_desired(def, index, id, true, g));
872        }
873    }
874    desired
875}
876
877/// Returns `true` if the `(s, d)` pair is still desired under `def` given the
878/// current graph state.  Calls `evaluate` directly — bypasses the candidate
879/// index.  Valid in `rebuild` after a full reindex because every pair that
880/// evaluates to `Some` is reachable via the freshly-built index; the direct
881/// call is therefore semantically equivalent to membership in
882/// `compute_desired(def, index, s, true, g)`.  O(eval) per call.
883fn pair_still_desired(def: &RuleDef, s: u32, d: u32, g: &GraphMut<'_>) -> bool {
884    let src_sym = match g.syms.get(&def.src_label) {
885        Some(sym) => sym,
886        None => return false,
887    };
888    let dst_sym = match g.syms.get(&def.dst_label) {
889        Some(sym) => sym,
890        None => return false,
891    };
892    if g.labels.get(s as usize).copied() != Some(src_sym) {
893        return false;
894    }
895    if g.labels.get(d as usize).copied() != Some(dst_sym) {
896        return false;
897    }
898    let s_key = match g.ids.key_of(s) {
899        Some(k) => k,
900        None => return false,
901    };
902    let d_key = match g.ids.key_of(d) {
903        Some(k) => k,
904        None => return false,
905    };
906    let s_get = |f: &str| g.props.get(s, f).cloned();
907    let d_get = |f: &str| g.props.get(d, f).cloned();
908    evaluate(
909        &def.predicate,
910        &NodeView {
911            key: s_key,
912            props: &s_get,
913        },
914        &NodeView {
915            key: d_key,
916            props: &d_get,
917        },
918    )
919    .is_some()
920}
921
922/// Count desired `(src, dst)` pairs up to `limit + 1`; returns as soon as the
923/// threshold is crossed.  Peak additional memory: O(max-candidates-per-src).
924/// Used in `rebuild` to detect the over-budget case without materialising the
925/// full desired map.
926fn count_desired_up_to(def: &RuleDef, index: &RuleIndex, limit: u64, g: &GraphMut<'_>) -> u64 {
927    let mut count = 0u64;
928    let src_sym = g.syms.get(&def.src_label);
929    for id in 0..g.ids.len() as u32 {
930        let label_sym = match g.labels.get(id as usize).copied() {
931            Some(s) if s != u32::MAX => s,
932            _ => continue,
933        };
934        if src_sym != Some(label_sym) {
935            continue;
936        }
937        count += compute_desired(def, index, id, true, g).len() as u64;
938        if count > limit {
939            return count;
940        }
941    }
942    count
943}
944
945/// Streaming backfill for `create_rule`.
946///
947/// Iterates src nodes in ascending id order.  For each src, computes and
948/// immediately applies the per-src desired edges (already dst-sorted within
949/// that src).  This traversal visits `(src, dst)` pairs in exactly the same
950/// order as iterating the result of `compute_full_desired` would — because the
951/// global BTree order on `(u32, u32)` keys is src-major ascending with
952/// dst-sorted-within, matching the per-src ascending-dst order emitted by
953/// `compute_desired`.
954///
955/// Consequently the first-N edges selected by the running cap are byte-identical
956/// to what the old full-map approach would have selected for EXACT rules
957/// (`def.approximate == false`).  For approximate rules the IVF candidate order
958/// is deterministic (replay-identical within the same fitted clusters) but is not
959/// equivalent to the old full-map path, which was never exercised for approximate
960/// rules.
961///
962/// Cap semantics: on `create_rule` there is no pre-existing provenance for
963/// this rule, so when the cap trips we can `break` immediately — no
964/// weight-refresh-on-existing path can be skipped, because every remaining
965/// `already == false` entry would have been `continue`-d by the old loop too.
966fn apply_streaming_create(
967    def: &RuleDef,
968    index: &RuleIndex,
969    prov: &mut ProvSets<'_>,
970    tripped: &mut bool,
971    g: &mut GraphMut<'_>,
972) {
973    let budget = edge_budget(def);
974    let et = g.syms.intern(&def.edge_type);
975    let src_sym = g.syms.get(&def.src_label);
976
977    'outer: for id in 0..g.ids.len() as u32 {
978        let label_sym = match g.labels.get(id as usize).copied() {
979            Some(s) if s != u32::MAX => s,
980            _ => continue,
981        };
982        if src_sym != Some(label_sym) {
983            continue;
984        }
985        let per_src = compute_desired(def, index, id, true, g);
986        for ((s, d), score) in per_src {
987            let triple = (et, s, d);
988            // On create_rule, this rule has no pre-existing provenance, so
989            // `already` is always false.  The path is kept for correctness
990            // if the same edge is co-owned by another rule (topo.add_edge
991            // returns false; we do not record it in our provenance).
992            let already = prov.contains(&triple);
993            if !already {
994                if *tripped || prov.len() as u64 >= budget {
995                    *tripped = true;
996                    break 'outer;
997                }
998                let newly = g.topo.add_edge(et, s, d);
999                if newly {
1000                    prov.insert(&def.name, triple, g.ids, g.syms);
1001                }
1002            }
1003            let is_owned_here = already || prov.contains(&triple);
1004            if is_owned_here {
1005                if let Some(p) = &def.weight_prop {
1006                    g.edge_props.set(et, s, d, p, Value::Float(score));
1007                }
1008            }
1009        }
1010    }
1011}
1012
1013/// Streaming backfill for `create_rule` with top-k per-source semantics.
1014///
1015/// Iterates src-label nodes in ascending id order. For each src, computes
1016/// the full candidate set, filters to the top-k destinations (score DESC,
1017/// dst_key ASC), and applies via `apply_per_src_top_k`.  No global budget or
1018/// tripped latch is used — the per-source cap is enforced by `filter_src_top_k`.
1019fn apply_streaming_create_top_k(
1020    def: &RuleDef,
1021    k: u64,
1022    index: &RuleIndex,
1023    prov: &mut ProvSets<'_>,
1024    g: &mut GraphMut<'_>,
1025) {
1026    let src_sym = g.syms.get(&def.src_label);
1027    for id in 0..g.ids.len() as u32 {
1028        let label_sym = match g.labels.get(id as usize).copied() {
1029            Some(s) if s != u32::MAX => s,
1030            _ => continue,
1031        };
1032        if src_sym != Some(label_sym) {
1033            continue;
1034        }
1035        let per_src = compute_desired(def, index, id, true, g);
1036        let top_k = filter_src_top_k(per_src, k, g.ids);
1037        apply_per_src_top_k(def, id, top_k, prov, g);
1038    }
1039}
1040
1041/// Streaming rebuild for top-k per-source rules.
1042///
1043/// Iterates all src-label nodes. For each src, computes fresh desired set,
1044/// filters to top-k, and applies via `apply_per_src_top_k` — which retracts
1045/// stale edges and inserts newly-ranked ones.  No budget counting, no tripped
1046/// latch.
1047fn apply_streaming_rebuild_top_k(
1048    def: &RuleDef,
1049    k: u64,
1050    index: &RuleIndex,
1051    prov: &mut ProvSets<'_>,
1052    g: &mut GraphMut<'_>,
1053) {
1054    let et = g.syms.intern(&def.edge_type);
1055
1056    // Collect all src nodes: those currently in provenance (may have lost their
1057    // label since last fire) + all live src-label nodes.
1058    let existing_srcs: BTreeSet<u32> = prov
1059        .set
1060        .iter()
1061        .filter(|(t, _, _)| *t == et)
1062        .map(|(_, s, _)| *s)
1063        .collect();
1064
1065    let src_sym = g.syms.get(&def.src_label);
1066    let mut all_srcs: BTreeSet<u32> = existing_srcs;
1067    for id in 0..g.ids.len() as u32 {
1068        let label_sym = match g.labels.get(id as usize).copied() {
1069            Some(s) if s != u32::MAX => s,
1070            _ => continue,
1071        };
1072        if src_sym == Some(label_sym) {
1073            all_srcs.insert(id);
1074        }
1075    }
1076
1077    for src in all_srcs {
1078        let desired_src = compute_desired(def, index, src, true, g);
1079        let top_k = filter_src_top_k(desired_src, k, g.ids);
1080        apply_per_src_top_k(def, src, top_k, prov, g);
1081    }
1082}
1083
1084/// Streaming rebuild for `rebuild`.
1085///
1086/// Replaces `compute_full_desired` + size-check + `apply_desired(None)`.
1087///
1088/// Over-budget path: counts desired pairs up to `budget + 1` (early exit);
1089/// if the total exceeds the budget, sets the tripped latch and returns without
1090/// touching provenance (identical to the current no-op rebuild behaviour).
1091///
1092/// Fits-budget path: un-trips, retracts each existing provenance triple whose
1093/// pair is no longer desired via direct `evaluate` — O(existing × eval) —
1094/// then streams-adds all desired edges in the same src-ascending / dst-within
1095/// order as `apply_streaming_create`.
1096fn apply_streaming_rebuild(
1097    def: &RuleDef,
1098    index: &RuleIndex,
1099    prov: &mut ProvSets<'_>,
1100    tripped: &mut bool,
1101    g: &mut GraphMut<'_>,
1102) {
1103    let budget = edge_budget(def);
1104    let et = g.syms.intern(&def.edge_type);
1105
1106    // 1. Count: early-exit as soon as the budget is exceeded.
1107    let total = count_desired_up_to(def, index, budget, g);
1108    if total > budget {
1109        *tripped = true;
1110        return; // provenance untouched; latch stays set.
1111    }
1112
1113    // 2. Un-trip — full desired fits.
1114    *tripped = false;
1115
1116    // 3. Retract existing edges that are no longer desired.
1117    //    O(existing × eval): evaluate each pair directly, no full-map build.
1118    let current: Vec<Triple> = prov
1119        .set
1120        .iter()
1121        .filter(|(t, _, _)| *t == et)
1122        .copied()
1123        .collect();
1124    for (t, s, d) in current {
1125        if !pair_still_desired(def, s, d, g) {
1126            g.topo.remove_edge(t, s, d);
1127            g.edge_props.remove_edge(t, s, d);
1128            prov.remove(&def.name, (t, s, d), g.ids, g.syms);
1129        }
1130    }
1131
1132    // 4. Stream-add desired edges; refresh weights on already-owned edges.
1133    //    count_desired_up_to verified total <= budget so no trip guard is needed.
1134    let src_sym = g.syms.get(&def.src_label);
1135    for id in 0..g.ids.len() as u32 {
1136        let label_sym = match g.labels.get(id as usize).copied() {
1137            Some(s) if s != u32::MAX => s,
1138            _ => continue,
1139        };
1140        if src_sym != Some(label_sym) {
1141            continue;
1142        }
1143        let per_src = compute_desired(def, index, id, true, g);
1144        for ((s, d), score) in per_src {
1145            let triple = (et, s, d);
1146            let already = prov.contains(&triple);
1147            if !already {
1148                let newly = g.topo.add_edge(et, s, d);
1149                if newly {
1150                    prov.insert(&def.name, triple, g.ids, g.syms);
1151                }
1152            }
1153            let is_owned_here = already || prov.contains(&triple);
1154            if is_owned_here {
1155                if let Some(p) = &def.weight_prop {
1156                    g.edge_props.set(et, s, d, p, Value::Float(score));
1157                }
1158            }
1159        }
1160    }
1161}
1162
1163/// Increment `fires` once per live node whose label matches either side.
1164/// Backfill / rebuild counting: one tick per participating node evaluated.
1165fn bump_fires_for_participants(def: &RuleDef, g: &GraphMut<'_>, fires: &mut u64) {
1166    let src_sym = g.syms.get(&def.src_label);
1167    let dst_sym = g.syms.get(&def.dst_label);
1168    for id in 0..g.ids.len() as u32 {
1169        let label_sym = match g.labels.get(id as usize).copied() {
1170            Some(s) if s != u32::MAX => s,
1171            _ => continue,
1172        };
1173        if src_sym == Some(label_sym) || dst_sym == Some(label_sym) {
1174            *fires += 1;
1175        }
1176    }
1177}
1178
1179/// Insert node `id` into the rule's src/dst indexes where its label matches.
1180fn index_node_for_rule(
1181    id: u32,
1182    label_sym: u32,
1183    def: &RuleDef,
1184    index: &mut RuleIndex,
1185    syms: &Interner,
1186    props: &ColumnStore,
1187) {
1188    let get = |f: &str| props.get(id, f).cloned();
1189    if syms.get(&def.src_label) == Some(label_sym) {
1190        let spec = src_lookup_spec_for(def);
1191        index.src_side.insert(&spec, id, &get);
1192    }
1193    if syms.get(&def.dst_label) == Some(label_sym) {
1194        let spec = candidate_spec_for(def);
1195        index.dst_side.insert(&spec, id, &get);
1196    }
1197}
1198
1199// ---------------------------------------------------------------------------
1200// RuleEngine
1201// ---------------------------------------------------------------------------
1202
1203impl RuleEngine {
1204    pub fn new() -> Self {
1205        Self::default()
1206    }
1207
1208    pub fn rules(&self) -> impl Iterator<Item = &RuleDef> {
1209        self.rules.values()
1210    }
1211
1212    pub fn is_owned(&self, etype: u32, src: u32, dst: u32) -> bool {
1213        self.owned.contains(&(etype, src, dst))
1214    }
1215
1216    /// Read-only view of the provenance map: rule name → set of (etype_sym, src, dst).
1217    pub fn provenance(&self) -> &BTreeMap<String, BTreeSet<(u32, u32, u32)>> {
1218        &self.provenance
1219    }
1220
1221    /// O(degree) reverse-index lookup: every provenance triple that touches `node`.
1222    pub fn provenance_touching(
1223        &self,
1224        node: u32,
1225    ) -> impl Iterator<Item = (&str, u32, u32, u32)> + '_ {
1226        self.by_node
1227            .get(&node)
1228            .into_iter()
1229            .flatten()
1230            .map(|&(rid, t, s, d)| (self.intern_rule[rid as usize].as_str(), t, s, d))
1231    }
1232
1233    /// Number of provenance triples incident on `node`.
1234    pub fn provenance_touching_len(&self, node: u32) -> usize {
1235        self.by_node.get(&node).map_or(0, BTreeSet::len)
1236    }
1237
1238    /// One-way latch: `true` after a budget breach until [`Self::rebuild`]
1239    /// is the only exit (and only if the full desired set then fits).
1240    pub fn is_tripped(&self, name: &str) -> bool {
1241        self.tripped.get(name).copied().unwrap_or(false)
1242    }
1243
1244    /// Evaluations of this rule: one tick per `on_node_changed` fire, and
1245    /// one tick per participating node on backfill **and rebuild** (even
1246    /// when rebuild is a provenance no-op).
1247    pub fn fire_count(&self, name: &str) -> u64 {
1248        self.fires.get(name).copied().unwrap_or(0)
1249    }
1250
1251    /// Drain and return all pending edge-fire / retract deltas since the last
1252    /// call.  Callers (`db.rs` `log_then_apply_with`) invoke this after a
1253    /// successful WAL commit + apply to build [`DbEvent`]s for live
1254    /// subscriptions.  [`GraphDb::open_with`] drains and discards after WAL
1255    /// replay so replay noise never leaks to subscribers.
1256    ///
1257    /// # T2 note (as-of replay)
1258    ///
1259    /// When Plan-15 T2 adds as-of replay for subscribers, that path should
1260    /// call apply-only (no `log_then_apply_with`) and then call
1261    /// `drain_deltas()` to feed those events to the replaying subscriber.
1262    /// The suppression is already in place: `apply` accumulates but never
1263    /// emits; `drain_deltas` is the only emission gate.
1264    pub fn drain_deltas(&mut self) -> Vec<EngineEdgeDelta> {
1265        std::mem::take(&mut self.pending_deltas)
1266    }
1267
1268    /// Number of accumulated deltas not yet drained.  Used by
1269    /// `debug_assert` in `log_then_apply_with` to catch stale-delta bugs.
1270    pub fn pending_delta_count(&self) -> usize {
1271        self.pending_deltas.len()
1272    }
1273
1274    /// Borrow the slice of deltas accumulated since `cursor` without
1275    /// consuming them.  `cursor` should be the value returned by
1276    /// `pending_delta_count()` before an engine call.
1277    ///
1278    /// The returned slice is valid until the next call to `drain_deltas()`.
1279    /// T1's drain discipline is preserved: these deltas are still in the
1280    /// buffer and will be drained by `log_then_apply_with` after `apply`
1281    /// returns.
1282    pub fn pending_deltas_since(&self, cursor: usize) -> &[EngineEdgeDelta] {
1283        &self.pending_deltas[cursor..]
1284    }
1285
1286    /// Snapshot support: definitions + provenance + tripped/fires. Candidate
1287    /// indexes and the `by_node` reverse index are NOT included (derived:
1288    /// `reindex_all` / `rebuild_by_node` on open).
1289    #[allow(clippy::type_complexity)]
1290    pub fn to_persist(
1291        &self,
1292    ) -> (
1293        Vec<RuleDef>,
1294        BTreeMap<String, BTreeSet<(u32, u32, u32)>>,
1295        BTreeMap<String, bool>,
1296        BTreeMap<String, u64>,
1297    ) {
1298        (
1299            self.rules.values().cloned().collect(),
1300            self.provenance.clone(),
1301            self.tripped.clone(),
1302            self.fires.clone(),
1303        )
1304    }
1305
1306    /// Reconstruct engine from a snapshot.  Caller must call `reindex_all` after.
1307    pub fn from_persist(
1308        rules: Vec<RuleDef>,
1309        prov: BTreeMap<String, BTreeSet<(u32, u32, u32)>>,
1310        tripped: BTreeMap<String, bool>,
1311        fires: BTreeMap<String, u64>,
1312    ) -> Self {
1313        let mut owned = BTreeSet::new();
1314        for set in prov.values() {
1315            owned.extend(set.iter().copied());
1316        }
1317        let indexes = rules
1318            .iter()
1319            .map(|r| (r.name.clone(), RuleIndex::default()))
1320            .collect();
1321        let rules: BTreeMap<String, RuleDef> =
1322            rules.into_iter().map(|r| (r.name.clone(), r)).collect();
1323        // Fill any missing keys so live rules always have entries.
1324        let mut tripped = tripped;
1325        let mut fires = fires;
1326        for name in rules.keys() {
1327            tripped.entry(name.clone()).or_insert(false);
1328            fires.entry(name.clone()).or_insert(0);
1329        }
1330        let (by_node, rule_intern, intern_rule) = rebuild_by_node(&prov);
1331        Self {
1332            rules,
1333            indexes,
1334            provenance: prov,
1335            owned,
1336            by_node,
1337            rule_intern,
1338            intern_rule,
1339            tripped,
1340            fires,
1341            pending_deltas: Vec::new(),
1342            emit_deltas: false,
1343            rebuild_needed: BTreeSet::new(),
1344        }
1345    }
1346
1347    /// Enable or disable delta accumulation.
1348    ///
1349    /// Set to `true` before the first subscriber or view is added.
1350    /// Set to `false` when the last subscriber and last view are removed.
1351    /// See the `emit_deltas` field doc for the safety invariant.
1352    pub fn set_emit_deltas(&mut self, emit: bool) {
1353        self.emit_deltas = emit;
1354    }
1355
1356    /// Whether delta accumulation is currently enabled.
1357    pub fn emit_deltas(&self) -> bool {
1358        self.emit_deltas
1359    }
1360
1361    /// Drain rule names that exceeded the IVF dst-drift rebuild threshold
1362    /// during the most recent `on_node_changed` / `on_node_removed`.
1363    pub fn take_rebuild_needed(&mut self) -> Vec<String> {
1364        std::mem::take(&mut self.rebuild_needed)
1365            .into_iter()
1366            .collect()
1367    }
1368
1369    /// Re-queue `name` so a later write can issue `RebuildRule`.
1370    ///
1371    /// Used when auto-rebuild WAL IO fails after a durable user op.
1372    pub fn queue_rebuild_needed(&mut self, name: String) {
1373        self.rebuild_needed.insert(name);
1374    }
1375
1376    fn maybe_queue_ivf_rebuild(&mut self, rule_name: &str, def: &RuleDef) {
1377        if !def.approximate {
1378            return;
1379        }
1380        let Some(idx) = self.indexes.get(rule_name) else {
1381            return;
1382        };
1383        if idx.dst_side.ivf_drift > ivf_drift_rebuild_threshold() {
1384            self.rebuild_needed.insert(rule_name.to_string());
1385        }
1386    }
1387
1388    /// Export IVF state for all approximate rules.  Passed to `snapshot()` in
1389    /// `core-api` and stored in the V4 snapshot so `open()` can restore cluster
1390    /// assignments without re-fitting k-means.
1391    pub fn export_ivf_state(&self) -> BTreeMap<String, RuleIvfExport> {
1392        let mut out = BTreeMap::new();
1393        for (name, def) in &self.rules {
1394            if def.approximate {
1395                if let Some(idx) = self.indexes.get(name) {
1396                    out.insert(
1397                        name.clone(),
1398                        (
1399                            idx.src_side.export_ivf_state(),
1400                            idx.dst_side.export_ivf_state(),
1401                        ),
1402                    );
1403                }
1404            }
1405        }
1406        out
1407    }
1408
1409    /// Rebuild all candidate indexes by scanning every node.  Call on open.
1410    pub fn reindex_all(
1411        &mut self,
1412        ids: &IdMap,
1413        syms: &Interner,
1414        labels: &[u32],
1415        props: &ColumnStore,
1416    ) {
1417        for idx in self.indexes.values_mut() {
1418            *idx = RuleIndex::default();
1419        }
1420        // Collect rule names once outside the per-node loop to avoid repeated
1421        // allocation and to satisfy the borrow checker without cloning inside.
1422        let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1423
1424        // Init HNSW for approximate rules before inserting nodes.
1425        for name in &rule_names {
1426            if self.rules[name].approximate {
1427                let idx = self.indexes.get_mut(name).unwrap();
1428                idx.src_side.init_hnsw(name);
1429                idx.dst_side.init_hnsw(name);
1430            }
1431        }
1432
1433        for id in 0..ids.len() as u32 {
1434            let label_sym = match labels.get(id as usize).copied() {
1435                Some(s) if s != u32::MAX => s,
1436                _ => continue,
1437            };
1438            for name in &rule_names {
1439                let def = self.rules[name].clone();
1440                let idx = self.indexes.get_mut(name).unwrap();
1441                index_node_for_rule(id, label_sym, &def, idx, syms, props);
1442            }
1443        }
1444        // After all nodes are indexed, fit IVF clusters for approximate rules.
1445        // HNSW was built incrementally; IVF kept as legacy fallback.
1446        for name in &rule_names {
1447            if self.rules[name].approximate {
1448                let idx = self.indexes.get_mut(name).unwrap();
1449                idx.src_side.fit_ivf_clusters(name);
1450                idx.dst_side.fit_ivf_clusters(name);
1451            }
1452        }
1453    }
1454
1455    /// Like `reindex_all` but LOADS persisted IVF state for approximate rules
1456    /// instead of re-fitting k-means.  This eliminates the cold-start re-fit
1457    /// cost when opening a V4 snapshot.
1458    ///
1459    /// `ivf_state`: map from rule name to `(src_export, dst_export)` as
1460    /// produced by `export_ivf_state` / stored in the V4 snapshot.
1461    ///
1462    /// For approximate rules absent from `ivf_state` (e.g. a rule added
1463    /// after the snapshot), falls back to `fit_ivf_clusters`.
1464    pub fn reindex_all_load_ivf(
1465        &mut self,
1466        ids: &IdMap,
1467        syms: &Interner,
1468        labels: &[u32],
1469        props: &ColumnStore,
1470        ivf_state: BTreeMap<String, RuleIvfExport>,
1471    ) {
1472        for idx in self.indexes.values_mut() {
1473            *idx = RuleIndex::default();
1474        }
1475        let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1476
1477        // Init HNSW for approximate rules before inserting nodes so each
1478        // insert also populates the HNSW graph incrementally.
1479        for name in &rule_names {
1480            if self.rules[name].approximate {
1481                let idx = self.indexes.get_mut(name).unwrap();
1482                idx.src_side.init_hnsw(name);
1483                idx.dst_side.init_hnsw(name);
1484            }
1485        }
1486
1487        for id in 0..ids.len() as u32 {
1488            let label_sym = match labels.get(id as usize).copied() {
1489                Some(s) if s != u32::MAX => s,
1490                _ => continue,
1491            };
1492            for name in &rule_names {
1493                let def = self.rules[name].clone();
1494                let idx = self.indexes.get_mut(name).unwrap();
1495                index_node_for_rule(id, label_sym, &def, idx, syms, props);
1496            }
1497        }
1498        // For approximate rules: restore persisted IVF state (no re-fit).
1499        // HNSW was built incrementally; it will be replaced by `load_hnsw_state`
1500        // in `restore_snapshot_state` when a snapshot blob is available.
1501        for name in &rule_names {
1502            if !self.rules[name].approximate {
1503                continue;
1504            }
1505            let idx = self.indexes.get_mut(name).unwrap();
1506            if let Some(((sc, sa, sd), (dc, da, dd))) = ivf_state.get(name) {
1507                idx.src_side.load_ivf_state(sc.clone(), sa.clone(), *sd);
1508                idx.dst_side.load_ivf_state(dc.clone(), da.clone(), *dd);
1509            } else {
1510                // No persisted state for this rule: fall back to full re-fit.
1511                idx.src_side.fit_ivf_clusters(name);
1512                idx.dst_side.fit_ivf_clusters(name);
1513            }
1514        }
1515    }
1516
1517    /// Export HNSW graphs for all approximate rules as opaque bincoded blobs.
1518    ///
1519    /// Returns a map from rule name to `(src_blob, dst_blob)`.  An empty `Vec`
1520    /// means the corresponding side has no initialized HNSW graph.
1521    pub fn export_hnsw_state(&self) -> BTreeMap<String, (Vec<u8>, Vec<u8>)> {
1522        let mut out = BTreeMap::new();
1523        for (name, def) in &self.rules {
1524            if def.approximate {
1525                if let Some(idx) = self.indexes.get(name) {
1526                    out.insert(
1527                        name.clone(),
1528                        (
1529                            idx.src_side.export_hnsw_blob(),
1530                            idx.dst_side.export_hnsw_blob(),
1531                        ),
1532                    );
1533                }
1534            }
1535        }
1536        out
1537    }
1538
1539    /// Restore HNSW graphs from bincoded blobs (overrides any incrementally built
1540    /// graphs produced during `reindex_all_load_ivf`).
1541    ///
1542    /// Called from `restore_snapshot_state` in db.rs after reindex.
1543    pub fn load_hnsw_state(&mut self, blobs: BTreeMap<String, (Vec<u8>, Vec<u8>)>) {
1544        for (name, (src_blob, dst_blob)) in blobs {
1545            if let Some(idx) = self.indexes.get_mut(&name) {
1546                if !src_blob.is_empty() {
1547                    idx.src_side.load_hnsw_blob(&src_blob);
1548                }
1549                if !dst_blob.is_empty() {
1550                    idx.dst_side.load_hnsw_blob(&dst_blob);
1551                }
1552            }
1553        }
1554    }
1555
1556    /// Find approximate nearest-neighbor ids on the dst side of the first
1557    /// approximate VectorSimilar rule covering `(dst_label, field)`.
1558    ///
1559    /// Returns `None` when no matching rule or HNSW index exists.
1560    pub fn hnsw_search_dst(
1561        &self,
1562        field: &str,
1563        dst_label: &str,
1564        q: &[f64],
1565        k: usize,
1566    ) -> Option<Vec<(u32, f64)>> {
1567        for (name, def) in &self.rules {
1568            if !def.approximate || def.dst_label != dst_label {
1569                continue;
1570            }
1571            // Check that the predicate covers this vector field.
1572            if !predicate_covers_field(&def.predicate, field) {
1573                continue;
1574            }
1575            if let Some(idx) = self.indexes.get(name) {
1576                if let Some(h) = idx.dst_side.hnsw_ref() {
1577                    if !h.is_empty() {
1578                        return Some(h.search(q, k));
1579                    }
1580                }
1581            }
1582        }
1583        None
1584    }
1585
1586    /// Register a rule and backfill existing nodes.
1587    /// Returns Err on failed validate() or duplicate name.
1588    pub fn create_rule(&mut self, def: RuleDef, g: &mut GraphMut<'_>) -> Result<(), String> {
1589        def.validate()?;
1590        if self.rules.contains_key(&def.name) {
1591            return Err(format!("rule {:?} already exists", def.name));
1592        }
1593        let name = def.name.clone();
1594        self.rules.insert(name.clone(), def);
1595        self.indexes.insert(name.clone(), RuleIndex::default());
1596        self.provenance.entry(name.clone()).or_default();
1597        self.tripped.insert(name.clone(), false);
1598        self.fires.insert(name.clone(), 0);
1599
1600        // Phase 1: index all existing nodes for this rule.
1601        let n_total = g.ids.len() as u32;
1602        let def = self.rules[&name].clone();
1603
1604        // Phase 1a: init HNSW for approximate rules before inserting nodes so
1605        // each insert also populates the HNSW graph incrementally.
1606        if def.approximate {
1607            let idx = self.indexes.get_mut(&name).unwrap();
1608            idx.src_side.init_hnsw(&name);
1609            idx.dst_side.init_hnsw(&name);
1610        }
1611
1612        for id in 0..n_total {
1613            let label_sym = match g.labels.get(id as usize).copied() {
1614                Some(s) if s != u32::MAX => s,
1615                _ => continue,
1616            };
1617            let idx = self.indexes.get_mut(&name).unwrap();
1618            index_node_for_rule(id, label_sym, &def, idx, g.syms, g.props);
1619        }
1620
1621        // Phase 1b: fit IVF clusters for approximate rules (after all nodes indexed).
1622        // HNSW was built incrementally above; IVF is kept as legacy fallback.
1623        if def.approximate {
1624            let idx = self.indexes.get_mut(&name).unwrap();
1625            idx.src_side.fit_ivf_clusters(&name);
1626            idx.dst_side.fit_ivf_clusters(&name);
1627        }
1628
1629        // Phase 2: streaming backfill.
1630        // Branches on max_edges semantics:
1631        //   None    → global-budget path (tripped latch, first-N in BTree order)
1632        //   Some(k) → per-source top-k path (no tripped latch, score-ordered)
1633        let mut prov = ProvSets {
1634            set: self.provenance.get_mut(&name).unwrap(),
1635            owned: &mut self.owned,
1636            by_node: &mut self.by_node,
1637            rule_intern: &mut self.rule_intern,
1638            intern_rule: &mut self.intern_rule,
1639            deltas: &mut self.pending_deltas,
1640            emit: self.emit_deltas,
1641        };
1642        if def.via_label.is_some() {
1643            // Via-hop backfill: bypass the candidate index; use compute_desired_via.
1644            let budget = edge_budget(&def);
1645            let et = g.syms.intern(&def.edge_type);
1646            let src_sym = g.syms.get(&def.src_label);
1647            let tripped = self.tripped.get_mut(&name).unwrap();
1648            'via_outer: for id in 0..g.ids.len() as u32 {
1649                let label_sym = match g.labels.get(id as usize).copied() {
1650                    Some(s) if s != u32::MAX => s,
1651                    _ => continue,
1652                };
1653                if src_sym != Some(label_sym) {
1654                    continue;
1655                }
1656                let per_src = compute_desired_via(&def, ViaAnchor::Src(id), g);
1657                if let Some(k) = def.max_edges {
1658                    let top_k = filter_src_top_k(per_src, k, g.ids);
1659                    apply_per_src_top_k(&def, id, top_k, &mut prov, g);
1660                } else {
1661                    for ((s, d), score) in per_src {
1662                        let triple = (et, s, d);
1663                        let already = prov.contains(&triple);
1664                        if !already {
1665                            if *tripped || prov.len() as u64 >= budget {
1666                                *tripped = true;
1667                                break 'via_outer;
1668                            }
1669                            let newly = g.topo.add_edge(et, s, d);
1670                            if newly {
1671                                prov.insert(&name, triple, g.ids, g.syms);
1672                            }
1673                        }
1674                        let is_owned_here = already || prov.contains(&triple);
1675                        if is_owned_here {
1676                            if let Some(p) = &def.weight_prop {
1677                                g.edge_props.set(et, s, d, p, Value::Float(score));
1678                            }
1679                        }
1680                    }
1681                }
1682            }
1683        } else if let Some(k) = def.max_edges {
1684            apply_streaming_create_top_k(&def, k, &self.indexes[&name], &mut prov, g);
1685        } else {
1686            let tripped = self.tripped.get_mut(&name).unwrap();
1687            apply_streaming_create(&def, &self.indexes[&name], &mut prov, tripped, g);
1688        }
1689        // Fires: one tick per participating node evaluated (same unit as
1690        // on_node_changed). Empty-graph create_rule therefore leaves fires=0.
1691        let fires = self.fires.get_mut(&name).unwrap();
1692        bump_fires_for_participants(&def, g, fires);
1693
1694        Ok(())
1695    }
1696
1697    /// Remove the rule and exactly its owned edges.  Returns Err if unknown.
1698    pub fn delete_rule(&mut self, name: &str, g: &mut GraphMut<'_>) -> Result<(), String> {
1699        if !self.rules.contains_key(name) {
1700            return Err(format!("rule {:?} not found", name));
1701        }
1702        let def = self.rules.remove(name).unwrap();
1703        self.indexes.remove(name);
1704        self.tripped.remove(name);
1705        self.fires.remove(name);
1706        let mut leftover = self.provenance.remove(name).unwrap_or_default();
1707        // intern so the symbol exists; edge_type was already interned at create time.
1708        let _et = g.syms.intern(&def.edge_type);
1709        let triples: Vec<Triple> = leftover.iter().copied().collect();
1710        let mut sets = ProvSets {
1711            set: &mut leftover,
1712            owned: &mut self.owned,
1713            by_node: &mut self.by_node,
1714            rule_intern: &mut self.rule_intern,
1715            intern_rule: &mut self.intern_rule,
1716            deltas: &mut self.pending_deltas,
1717            emit: self.emit_deltas,
1718        };
1719        for triple in triples {
1720            let (t, s, d) = triple;
1721            g.topo.remove_edge(t, s, d);
1722            g.edge_props.remove_edge(t, s, d);
1723            sets.remove(name, triple, g.ids, g.syms);
1724        }
1725        // Surviving rules that share the same edge_type may derive edges that
1726        // were previously blocked (add_edge returned false because the deleted
1727        // rule already owned them, so their provenance never recorded them).
1728        // Rebuilding each such rule lets it claim those edges now that the
1729        // deleted rule's entries have been removed from the topology.
1730        let same_etype_survivors: Vec<String> = self
1731            .rules
1732            .values()
1733            .filter(|r| r.edge_type == def.edge_type)
1734            .map(|r| r.name.clone())
1735            .collect();
1736        for survivor in same_etype_survivors {
1737            // rebuild returns Err only for unknown rules; survivor is live.
1738            let _ = self.rebuild(&survivor, g);
1739        }
1740        Ok(())
1741    }
1742
1743    /// Called when node `n` is inserted (changed=None) or a field is updated.
1744    /// - None: all rules where n's label matches either side fire; index gains n.
1745    /// - Some((field, old_value)): only rules watching `field` fire; index is
1746    ///   updated using old_value for removal so stale buckets are cleaned.
1747    ///
1748    /// For via-hop rules (`def.via_label.is_some()`), also fires when n carries
1749    /// the via-label: finds all srcs that route through n and recomputes their
1750    /// derived edges. Via-hop rules bypass the candidate index and use
1751    /// `compute_desired_via` instead.
1752    pub fn on_node_changed(
1753        &mut self,
1754        n: u32,
1755        changed: Option<(&str, Option<Value>)>,
1756        g: &mut GraphMut<'_>,
1757    ) {
1758        let n_label = g.labels.get(n as usize).copied();
1759        let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1760
1761        for rule_name in rule_names {
1762            let def = self.rules[&rule_name].clone();
1763
1764            if def.via_label.is_some() {
1765                // --- Via-hop rule path ---
1766                self.on_node_changed_via(&rule_name, &def, n, n_label, changed.clone(), g);
1767            } else {
1768                // --- Standard 2-node rule path ---
1769                let src_sym = g.syms.get(&def.src_label);
1770                let dst_sym = g.syms.get(&def.dst_label);
1771                let as_src = src_sym.is_some() && n_label == src_sym;
1772                let as_dst = dst_sym.is_some() && n_label == dst_sym;
1773
1774                let fires = match changed {
1775                    None => as_src || as_dst,
1776                    Some((field, _)) => def.watched_fields().contains(field) && (as_src || as_dst),
1777                };
1778                if !fires {
1779                    continue;
1780                }
1781                *self.fires.entry(rule_name.clone()).or_default() += 1;
1782
1783                // --- Index maintenance ---
1784                if let Some((field, ref old_val)) = changed {
1785                    let old_val_cloned = old_val.clone();
1786                    let old_getter = |f: &str| {
1787                        if f == field {
1788                            old_val_cloned.clone()
1789                        } else {
1790                            g.props.get(n, f).cloned()
1791                        }
1792                    };
1793                    let idx = self.indexes.get_mut(&rule_name).unwrap();
1794                    if as_src {
1795                        let spec = src_lookup_spec_for(&def);
1796                        idx.src_side.remove(&spec, n, &old_getter);
1797                    }
1798                    if as_dst {
1799                        let spec = candidate_spec_for(&def);
1800                        idx.dst_side.remove(&spec, n, &old_getter);
1801                    }
1802                }
1803
1804                {
1805                    let cur_getter = |f: &str| g.props.get(n, f).cloned();
1806                    let idx = self.indexes.get_mut(&rule_name).unwrap();
1807                    if as_src {
1808                        let spec = src_lookup_spec_for(&def);
1809                        idx.src_side.insert(&spec, n, &cur_getter);
1810                    }
1811                    if as_dst {
1812                        let spec = candidate_spec_for(&def);
1813                        idx.dst_side.insert(&spec, n, &cur_getter);
1814                    }
1815                }
1816
1817                self.maybe_queue_ivf_rebuild(&rule_name, &def);
1818
1819                // --- Desired set + diff-apply ---
1820                if let Some(k) = def.max_edges {
1821                    let et = g.syms.intern(&def.edge_type);
1822                    let affected_srcs_for_n_dst: BTreeSet<u32> = if as_dst {
1823                        let rid = self.rule_intern.get(&def.name).copied();
1824                        self.by_node
1825                            .get(&n)
1826                            .into_iter()
1827                            .flatten()
1828                            .filter(|(r, t, _s, d)| Some(*r) == rid && *t == et && *d == n)
1829                            .map(|(_, _, s, _)| *s)
1830                            .collect()
1831                    } else {
1832                        BTreeSet::new()
1833                    };
1834
1835                    let mut prov = ProvSets {
1836                        set: self.provenance.entry(rule_name.clone()).or_default(),
1837                        owned: &mut self.owned,
1838                        by_node: &mut self.by_node,
1839                        rule_intern: &mut self.rule_intern,
1840                        intern_rule: &mut self.intern_rule,
1841                        deltas: &mut self.pending_deltas,
1842                        emit: self.emit_deltas,
1843                    };
1844
1845                    if as_src {
1846                        let desired_n_src =
1847                            compute_desired(&def, &self.indexes[&rule_name], n, true, g);
1848                        let top_k = filter_src_top_k(desired_n_src, k, g.ids);
1849                        apply_per_src_top_k(&def, n, top_k, &mut prov, g);
1850                    }
1851
1852                    if as_dst {
1853                        let new_desired =
1854                            compute_desired(&def, &self.indexes[&rule_name], n, false, g);
1855                        let new_srcs: BTreeSet<u32> = new_desired.keys().map(|(s, _)| *s).collect();
1856                        let affected_srcs: BTreeSet<u32> =
1857                            affected_srcs_for_n_dst.union(&new_srcs).copied().collect();
1858                        for src in affected_srcs {
1859                            if src == n {
1860                                continue;
1861                            }
1862                            let desired_src =
1863                                compute_desired(&def, &self.indexes[&rule_name], src, true, g);
1864                            let top_k = filter_src_top_k(desired_src, k, g.ids);
1865                            apply_per_src_top_k(&def, src, top_k, &mut prov, g);
1866                        }
1867                    }
1868                } else {
1869                    let mut desired = BTreeMap::new();
1870                    if as_src {
1871                        desired.extend(compute_desired(
1872                            &def,
1873                            &self.indexes[&rule_name],
1874                            n,
1875                            true,
1876                            g,
1877                        ));
1878                    }
1879                    if as_dst {
1880                        desired.extend(compute_desired(
1881                            &def,
1882                            &self.indexes[&rule_name],
1883                            n,
1884                            false,
1885                            g,
1886                        ));
1887                    }
1888                    let tripped = self.tripped.entry(rule_name.clone()).or_default();
1889                    apply_desired(
1890                        &def,
1891                        desired,
1892                        Some(n),
1893                        &mut ProvSets {
1894                            set: self.provenance.entry(rule_name).or_default(),
1895                            owned: &mut self.owned,
1896                            by_node: &mut self.by_node,
1897                            rule_intern: &mut self.rule_intern,
1898                            intern_rule: &mut self.intern_rule,
1899                            deltas: &mut self.pending_deltas,
1900                            emit: self.emit_deltas,
1901                        },
1902                        tripped,
1903                        g,
1904                    );
1905                }
1906            }
1907        }
1908    }
1909
1910    /// Inner handler for `on_node_changed` when the rule is a via-hop rule.
1911    ///
1912    /// For each role n can play (src, via, dst), computes and applies the
1913    /// desired edge set using `compute_desired_via`. Via-hop rules bypass the
1914    /// candidate index; no index maintenance is performed here.
1915    ///
1916    /// Incremental correctness by change class:
1917    /// - **src prop / insert** (`as_src`): re-expand via from n, recompute all
1918    ///   (n, dst) pairs. `apply_via_for_srcs([n])`.
1919    /// - **via-node prop change** (`as_via`, field in watched_fields): find
1920    ///   srcs that hop to n via `via_edge`, recompute their (src, dst) pairs.
1921    ///   `apply_via_for_srcs(reverse_via_neighbors(n))`.
1922    /// - **dst prop / insert** (`as_dst`): anchor on n, compute desired for all
1923    ///   srcs. `apply_via_for_srcs(all_src_label_nodes)`.
1924    fn on_node_changed_via(
1925        &mut self,
1926        rule_name: &str,
1927        def: &RuleDef,
1928        n: u32,
1929        n_label: Option<u32>,
1930        changed: Option<(&str, Option<Value>)>,
1931        g: &mut GraphMut<'_>,
1932    ) {
1933        let src_sym = g.syms.get(&def.src_label);
1934        let dst_sym = g.syms.get(&def.dst_label);
1935        let via_sym = def.via_label.as_deref().and_then(|l| g.syms.get(l));
1936
1937        let as_src = src_sym.is_some() && n_label == src_sym;
1938        let as_dst = dst_sym.is_some() && n_label == dst_sym;
1939        let as_via = via_sym.is_some() && n_label == via_sym;
1940
1941        // Via-hop predicates are evaluated between via and dst, so watched_fields
1942        // cover both via-node and dst-node fields (predicate fields come from the
1943        // via→dst evaluation). A via-node prop change fires if its field is watched.
1944        let fires = match changed {
1945            None => as_src || as_via || as_dst,
1946            Some((field, _)) => {
1947                let wf = def.watched_fields();
1948                (wf.contains(field)) && (as_src || as_via || as_dst)
1949            }
1950        };
1951        if !fires {
1952            return;
1953        }
1954        *self.fires.entry(rule_name.to_string()).or_default() += 1;
1955
1956        // Collect affected srcs: union of srcs identified from each role.
1957        let mut affected_srcs: BTreeSet<u32> = BTreeSet::new();
1958        if as_src {
1959            affected_srcs.insert(n);
1960        }
1961        if as_via {
1962            // Srcs that hop to this via-node via via_edge (reverse direction).
1963            let via_edge_str = def.via_edge.as_deref().unwrap();
1964            let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
1965            let rev_dir = match via_dir {
1966                core_storage::Direction::Out => core_storage::Direction::In,
1967                core_storage::Direction::In => core_storage::Direction::Out,
1968            };
1969            if let (Some(via_etype), Some(s_sym)) = (g.syms.get(via_edge_str), src_sym) {
1970                for &src in g.topo.neighbors(via_etype, rev_dir, n).as_ref() {
1971                    if g.labels.get(src as usize).copied() == Some(s_sym) {
1972                        affected_srcs.insert(src);
1973                    }
1974                }
1975            }
1976        }
1977        if as_dst {
1978            // Recompute all srcs whose via-hops might produce edges to n.
1979            let desired_touching_n = compute_desired_via(def, ViaAnchor::Dst(n), g);
1980            for (src, _dst) in desired_touching_n.keys() {
1981                affected_srcs.insert(*src);
1982            }
1983            // Also include any srcs that currently have provenance pointing to n.
1984            let et = g.syms.intern(&def.edge_type);
1985            let rid = self.rule_intern.get(rule_name).copied();
1986            let old_srcs: Vec<u32> = self
1987                .by_node
1988                .get(&n)
1989                .into_iter()
1990                .flatten()
1991                .filter(|(r, t, _s, d)| Some(*r) == rid && *t == et && *d == n)
1992                .map(|(_, _, s, _)| *s)
1993                .collect();
1994            affected_srcs.extend(old_srcs);
1995        }
1996
1997        // For each affected src, compute desired_via(Src) and apply.
1998        let affected_srcs: Vec<u32> = affected_srcs.into_iter().collect();
1999
2000        if let Some(k) = def.max_edges {
2001            let mut prov = ProvSets {
2002                set: self.provenance.entry(rule_name.to_string()).or_default(),
2003                owned: &mut self.owned,
2004                by_node: &mut self.by_node,
2005                rule_intern: &mut self.rule_intern,
2006                intern_rule: &mut self.intern_rule,
2007                deltas: &mut self.pending_deltas,
2008                emit: self.emit_deltas,
2009            };
2010            for src in affected_srcs {
2011                let desired_src = compute_desired_via(def, ViaAnchor::Src(src), g);
2012                let top_k = filter_src_top_k(desired_src, k, g.ids);
2013                apply_per_src_top_k(def, src, top_k, &mut prov, g);
2014            }
2015        } else {
2016            let tripped = self.tripped.entry(rule_name.to_string()).or_default();
2017            let budget = edge_budget(def);
2018            // Apply per-src so each affected src retracts its stale edges and
2019            // adds its new desired edges independently.
2020            for src in affected_srcs {
2021                let desired_src = compute_desired_via(def, ViaAnchor::Src(src), g);
2022                if !*tripped {
2023                    let mut prov = ProvSets {
2024                        set: self.provenance.entry(rule_name.to_string()).or_default(),
2025                        owned: &mut self.owned,
2026                        by_node: &mut self.by_node,
2027                        rule_intern: &mut self.rule_intern,
2028                        intern_rule: &mut self.intern_rule,
2029                        deltas: &mut self.pending_deltas,
2030                        emit: self.emit_deltas,
2031                    };
2032                    apply_desired(def, desired_src, Some(src), &mut prov, tripped, g);
2033                }
2034                // If budget was just tripped inside apply_desired, stop adding
2035                // but continue retracting stale edges for already-processed srcs
2036                // (apply_desired handles retracts even when tripped).
2037                let _ = budget;
2038            }
2039        }
2040    }
2041
2042    /// Called when a user edge `(etype_str, src_id, dst_id)` is inserted or
2043    /// deleted (not a derived edge — those are managed by provenance, not here).
2044    ///
2045    /// For any via-hop rule where `via_edge == etype_str` and src_id carries
2046    /// `src_label`, the src_id's desired derived-edge set may have changed:
2047    /// a new WORKS_AT edge makes a new Org reachable as a via-node, and a
2048    /// deleted WORKS_AT removes a previously reachable Org.
2049    ///
2050    /// This is the only hook the engine exposes for topology changes. It is
2051    /// called from `db.rs` on `WalRecord::InsertEdge` and `WalRecord::DeleteEdge`
2052    /// immediately after the topo is updated (so `g.topo` already reflects the
2053    /// new state).
2054    pub fn on_edge_changed(
2055        &mut self,
2056        etype_str: &str,
2057        src_id: u32,
2058        dst_id: u32,
2059        g: &mut GraphMut<'_>,
2060    ) {
2061        let rule_names: Vec<String> = self.rules.keys().cloned().collect();
2062        for rule_name in rule_names {
2063            let def = self.rules[&rule_name].clone();
2064            let Some(ref via_edge) = def.via_edge else {
2065                continue; // not a via-hop rule
2066            };
2067            if via_edge != etype_str {
2068                continue; // edge type doesn't match this rule's via_edge
2069            }
2070
2071            // Check that src_id carries src_label and dst_id carries via_label.
2072            let src_sym = match g.syms.get(&def.src_label) {
2073                Some(s) => s,
2074                None => continue,
2075            };
2076            let via_sym = match def.via_label.as_deref().and_then(|l| g.syms.get(l)) {
2077                Some(s) => s,
2078                None => continue,
2079            };
2080            // via_dir == Out → the edge goes src_id → dst_id (src-label node to via-label node)
2081            // via_dir == In  → the edge goes dst_id ← src_id, i.e., src_id is the via-label
2082            //                  end and dst_id is the src-label end. Adjust accordingly.
2083            let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
2084            let (rule_src, rule_via) = match via_dir {
2085                core_storage::Direction::Out => (src_id, dst_id),
2086                core_storage::Direction::In => (dst_id, src_id),
2087            };
2088
2089            if g.labels.get(rule_src as usize).copied() != Some(src_sym) {
2090                continue;
2091            }
2092            if g.labels.get(rule_via as usize).copied() != Some(via_sym) {
2093                continue;
2094            }
2095
2096            // Recompute derived edges for rule_src — its via-hop set just changed.
2097            *self.fires.entry(rule_name.clone()).or_default() += 1;
2098            let desired_src = compute_desired_via(&def, ViaAnchor::Src(rule_src), g);
2099
2100            if let Some(k) = def.max_edges {
2101                let mut prov = ProvSets {
2102                    set: self.provenance.entry(rule_name).or_default(),
2103                    owned: &mut self.owned,
2104                    by_node: &mut self.by_node,
2105                    rule_intern: &mut self.rule_intern,
2106                    intern_rule: &mut self.intern_rule,
2107                    deltas: &mut self.pending_deltas,
2108                    emit: self.emit_deltas,
2109                };
2110                let top_k = filter_src_top_k(desired_src, k, g.ids);
2111                apply_per_src_top_k(&def, rule_src, top_k, &mut prov, g);
2112            } else {
2113                let tripped = self.tripped.entry(rule_name.clone()).or_default();
2114                let mut prov = ProvSets {
2115                    set: self.provenance.entry(rule_name).or_default(),
2116                    owned: &mut self.owned,
2117                    by_node: &mut self.by_node,
2118                    rule_intern: &mut self.rule_intern,
2119                    intern_rule: &mut self.intern_rule,
2120                    deltas: &mut self.pending_deltas,
2121                    emit: self.emit_deltas,
2122                };
2123                apply_desired(&def, desired_src, Some(rule_src), &mut prov, tripped, g);
2124            }
2125        }
2126    }
2127
2128    /// Retract every provenance edge touching `n` across all rules and drop
2129    /// `n` from every rule index using its *current* props.
2130    ///
2131    /// Caller must invoke this while labels/props are still intact (before
2132    /// tombstone). Rules are walked in BTree name order; touching edges in
2133    /// BTree triple order. A second call on an already-retracted node is a
2134    /// no-op (crash-window replay / absent state).
2135    pub fn on_node_removed(&mut self, n: u32, g: &mut GraphMut<'_>) {
2136        let n_label = g.labels.get(n as usize).copied();
2137        let rule_names: Vec<String> = self.rules.keys().cloned().collect();
2138
2139        for rule_name in rule_names {
2140            let def = self.rules[&rule_name].clone();
2141            let src_sym = g.syms.get(&def.src_label);
2142            let dst_sym = g.syms.get(&def.dst_label);
2143            let as_src = src_sym.is_some() && n_label == src_sym;
2144            let as_dst = dst_sym.is_some() && n_label == dst_sym;
2145
2146            {
2147                let cur_getter = |f: &str| g.props.get(n, f).cloned();
2148                let idx = self.indexes.get_mut(&rule_name).unwrap();
2149                if as_src {
2150                    let spec = src_lookup_spec_for(&def);
2151                    idx.src_side.remove(&spec, n, &cur_getter);
2152                }
2153                if as_dst {
2154                    let spec = candidate_spec_for(&def);
2155                    idx.dst_side.remove(&spec, n, &cur_getter);
2156                }
2157            }
2158
2159            self.maybe_queue_ivf_rebuild(&rule_name, &def);
2160        }
2161
2162        let touching: Vec<(String, Triple)> = self
2163            .by_node
2164            .get(&n)
2165            .into_iter()
2166            .flatten()
2167            .map(|&(rid, t, s, d)| (self.intern_rule[rid as usize].clone(), (t, s, d)))
2168            .collect();
2169
2170        // Collect srcs that need top-k backfill BEFORE retracting provenance.
2171        // For top-k rules: when n is a dst, the src loses one from its top-k
2172        // and needs the next-best candidate added.
2173        let topk_backfill: Vec<(String, u32)> = touching
2174            .iter()
2175            .filter_map(|(rule_name, triple)| {
2176                let &(_, s, d) = triple;
2177                let def = self.rules.get(rule_name)?;
2178                def.max_edges?; // only top-k rules need backfill
2179                if d == n && s != n {
2180                    Some((rule_name.clone(), s))
2181                } else {
2182                    None
2183                }
2184            })
2185            .collect();
2186
2187        for (rule_name, triple) in touching {
2188            let (t, s, d) = triple;
2189            g.topo.remove_edge(t, s, d);
2190            g.edge_props.remove_edge(t, s, d);
2191            if let Some(set) = self.provenance.get_mut(&rule_name) {
2192                ProvSets {
2193                    set,
2194                    owned: &mut self.owned,
2195                    by_node: &mut self.by_node,
2196                    rule_intern: &mut self.rule_intern,
2197                    intern_rule: &mut self.intern_rule,
2198                    deltas: &mut self.pending_deltas,
2199                    emit: self.emit_deltas,
2200                }
2201                .remove(&rule_name, triple, g.ids, g.syms);
2202            }
2203        }
2204
2205        // Backfill top-k srcs whose dst was removed.
2206        // By now n is removed from the dst index (done in the first loop above),
2207        // so compute_desired(src, true) will not include n in candidates — the
2208        // resulting top-k automatically promotes the next-best candidate.
2209        for (rule_name, src) in topk_backfill {
2210            let def = self.rules[&rule_name].clone();
2211            let k = def.max_edges.unwrap(); // guarded by filter above
2212            let desired_src = compute_desired(&def, &self.indexes[&rule_name], src, true, g);
2213            let top_k = filter_src_top_k(desired_src, k, g.ids);
2214            let mut prov = ProvSets {
2215                set: self.provenance.entry(rule_name.clone()).or_default(),
2216                owned: &mut self.owned,
2217                by_node: &mut self.by_node,
2218                rule_intern: &mut self.rule_intern,
2219                intern_rule: &mut self.intern_rule,
2220                deltas: &mut self.pending_deltas,
2221                emit: self.emit_deltas,
2222            };
2223            apply_per_src_top_k(&def, src, top_k, &mut prov, g);
2224        }
2225    }
2226
2227    /// Recompute one rule from scratch. Only exit from the tripped latch.
2228    ///
2229    /// If the full desired set fits in the budget, it is applied completely
2230    /// and `tripped` is cleared. If it still exceeds the budget, existing
2231    /// provenance is left completely untouched and `tripped` stays true
2232    /// (rebuild-is-noop for at/over-cap rules). Always counts as a fire
2233    /// evaluation per participating node. Returns Err if unknown.
2234    pub fn rebuild(&mut self, name: &str, g: &mut GraphMut<'_>) -> Result<(), String> {
2235        if !self.rules.contains_key(name) {
2236            return Err(format!("rule {:?} not found", name));
2237        }
2238        self.rebuild_needed.remove(name);
2239        let def = self.rules[name].clone();
2240
2241        // Reindex this rule from scratch (indexes only).
2242        *self.indexes.get_mut(name).unwrap() = RuleIndex::default();
2243
2244        // Init HNSW before indexing so inserts populate the graph incrementally.
2245        if def.approximate {
2246            let idx = self.indexes.get_mut(name).unwrap();
2247            idx.src_side.init_hnsw(name);
2248            idx.dst_side.init_hnsw(name);
2249        }
2250
2251        let n_total = g.ids.len() as u32;
2252        for id in 0..n_total {
2253            let label_sym = match g.labels.get(id as usize).copied() {
2254                Some(s) if s != u32::MAX => s,
2255                _ => continue,
2256            };
2257            let idx = self.indexes.get_mut(name).unwrap();
2258            index_node_for_rule(id, label_sym, &def, idx, g.syms, g.props);
2259        }
2260
2261        // Fit IVF clusters for approximate rules after reindex (drift reset).
2262        // HNSW was built incrementally; IVF kept as legacy fallback.
2263        if def.approximate {
2264            let idx = self.indexes.get_mut(name).unwrap();
2265            idx.src_side.fit_ivf_clusters(name);
2266            idx.dst_side.fit_ivf_clusters(name);
2267        }
2268
2269        // Streaming rebuild: branches on max_edges semantics.
2270        //   None    → global-budget path (may no-op if still over budget)
2271        //   Some(k) → per-source top-k rebuild (always converges; no tripped latch)
2272        let mut prov = ProvSets {
2273            set: self.provenance.get_mut(name).unwrap(),
2274            owned: &mut self.owned,
2275            by_node: &mut self.by_node,
2276            rule_intern: &mut self.rule_intern,
2277            intern_rule: &mut self.intern_rule,
2278            deltas: &mut self.pending_deltas,
2279            emit: self.emit_deltas,
2280        };
2281        if let Some(k) = def.max_edges {
2282            apply_streaming_rebuild_top_k(&def, k, &self.indexes[name], &mut prov, g);
2283        } else {
2284            let tripped = self.tripped.get_mut(name).unwrap();
2285            apply_streaming_rebuild(&def, &self.indexes[name], &mut prov, tripped, g);
2286        }
2287        let fires = self.fires.entry(name.to_string()).or_default();
2288        bump_fires_for_participants(&def, g, fires);
2289
2290        Ok(())
2291    }
2292
2293    #[cfg(test)]
2294    fn by_node_consistent(&self) -> bool {
2295        let (rebuilt, intern, names) = rebuild_by_node(&self.provenance);
2296        resolve_by_node(&self.by_node, &self.intern_rule) == resolve_by_node(&rebuilt, &names)
2297            && intern.len() == names.len()
2298    }
2299}
2300
2301// ---------------------------------------------------------------------------
2302// Tests
2303// ---------------------------------------------------------------------------
2304
2305#[cfg(test)]
2306mod tests {
2307    use super::*;
2308    use crate::def::{evaluate, NodeView, Predicate, RuleDef};
2309    use core_storage::{ColumnStore, Direction, EdgeProps, IdMap, Interner, Topology, Value};
2310
2311    struct Fx {
2312        ids: IdMap,
2313        syms: Interner,
2314        labels: Vec<u32>,
2315        props: ColumnStore,
2316        topo: Topology,
2317        eprops: EdgeProps,
2318    }
2319    impl Fx {
2320        fn new() -> Self {
2321            Fx {
2322                ids: IdMap::new(),
2323                syms: Interner::new(),
2324                labels: vec![],
2325                props: ColumnStore::new(),
2326                topo: Topology::new(),
2327                eprops: EdgeProps::new(),
2328            }
2329        }
2330        fn add(&mut self, label: &str, key: &str, props: Vec<(&str, Value)>) -> u32 {
2331            let id = self.ids.get_or_insert(key);
2332            let sym = self.syms.intern(label);
2333            self.labels.resize(id as usize + 1, u32::MAX);
2334            self.labels[id as usize] = sym;
2335            for (f, v) in props {
2336                self.props.set(id, f, v);
2337            }
2338            id
2339        }
2340        fn g(&mut self) -> GraphMut<'_> {
2341            GraphMut {
2342                ids: &self.ids,
2343                syms: &mut self.syms,
2344                labels: &self.labels,
2345                props: &self.props,
2346                topo: &mut self.topo,
2347                edge_props: &mut self.eprops,
2348            }
2349        }
2350    }
2351
2352    fn tags(items: &[&str]) -> Value {
2353        Value::List(items.iter().map(|s| Value::Str((*s).into())).collect())
2354    }
2355
2356    fn overlap_rule() -> RuleDef {
2357        RuleDef {
2358            name: "rel".into(),
2359            src_label: "A".into(),
2360            dst_label: "A".into(),
2361            predicate: Predicate::Overlap {
2362                field: "tags".into(),
2363                min: 0.4,
2364            },
2365            edge_type: "REL".into(),
2366            weight_prop: Some("score".into()),
2367            max_edges: None,
2368            approximate: false,
2369            via_label: None,
2370            via_edge: None,
2371            via_dir: None,
2372        }
2373    }
2374
2375    fn emb(xs: &[f64]) -> Value {
2376        Value::List(xs.iter().copied().map(Value::Float).collect())
2377    }
2378
2379    fn approx_vec_rule() -> RuleDef {
2380        RuleDef {
2381            name: "sim".into(),
2382            src_label: "V".into(),
2383            dst_label: "V".into(),
2384            predicate: Predicate::VectorSimilar {
2385                field: "emb".into(),
2386                min: 0.5,
2387            },
2388            edge_type: "SIM".into(),
2389            weight_prop: None,
2390            max_edges: None,
2391            approximate: true,
2392            via_label: None,
2393            via_edge: None,
2394            via_dir: None,
2395        }
2396    }
2397
2398    #[test]
2399    fn approximate_rule_rebuilds_after_drift_threshold() {
2400        with_ivf_drift_rebuild(1, || {
2401            let mut fx = Fx::new();
2402            let mut ids = Vec::new();
2403            for i in 0..6 {
2404                let x = i as f64 * 0.2;
2405                ids.push(fx.add("V", &format!("v{i}"), vec![("emb", emb(&[x, 1.0 - x]))]));
2406            }
2407            let mut eng = RuleEngine::new();
2408            {
2409                let mut g = fx.g();
2410                eng.create_rule(approx_vec_rule(), &mut g).unwrap();
2411            }
2412            assert!(eng.take_rebuild_needed().is_empty());
2413            {
2414                let mut g = fx.g();
2415                eng.on_node_removed(ids[0], &mut g);
2416            }
2417            assert!(
2418                eng.take_rebuild_needed().is_empty(),
2419                "drift=1 is not > threshold 1"
2420            );
2421            {
2422                let mut g = fx.g();
2423                eng.on_node_removed(ids[1], &mut g);
2424            }
2425            assert_eq!(eng.take_rebuild_needed(), vec!["sim".to_string()]);
2426            {
2427                let mut g = fx.g();
2428                eng.rebuild("sim", &mut g).unwrap();
2429            }
2430            assert!(
2431                eng.take_rebuild_needed().is_empty(),
2432                "rebuild must reset drift and not re-queue itself"
2433            );
2434            let drift = eng
2435                .export_ivf_state()
2436                .get("sim")
2437                .map(|(_, dst)| dst.2)
2438                .unwrap();
2439            assert_eq!(drift, 0, "rebuild resets dst-side IVF drift");
2440        });
2441    }
2442
2443    #[test]
2444    fn backfill_creates_edges_with_scores_and_delete_removes_exactly_them() {
2445        let mut fx = Fx::new();
2446        let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2447        let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2448        let _c = fx.add("A", "c", vec![("tags", tags(&["q"]))]);
2449        // pre-existing user edge with same type: must survive rule delete
2450        let et = fx.syms.intern("REL");
2451        fx.topo.add_edge(et, a, b);
2452        let mut eng = RuleEngine::new();
2453        let mut g = fx.g();
2454        eng.create_rule(overlap_rule(), &mut g).unwrap();
2455        // a↔b jaccard 1.0 both directions; user edge a→b pre-existed so only b→a is owned
2456        assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2457        assert_eq!(
2458            g.edge_props.get(et, b, a, "score"),
2459            Some(&Value::Float(1.0))
2460        );
2461        assert!(!eng.is_owned(et, a, b));
2462        assert!(eng.is_owned(et, b, a));
2463        eng.delete_rule("rel", &mut g).unwrap();
2464        assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b)); // user edge kept
2465        assert!(!g.topo.neighbors(et, Direction::Out, b).contains(&a)); // derived removed
2466        assert_eq!(g.edge_props.get(et, b, a, "score"), None);
2467    }
2468
2469    #[test]
2470    fn incremental_update_adds_and_removes_edges() {
2471        let mut fx = Fx::new();
2472        let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2473        let b = fx.add("A", "b", vec![("tags", tags(&["y", "z"]))]);
2474        let et = fx.syms.intern("REL");
2475        let mut eng = RuleEngine::new();
2476        {
2477            let mut g = fx.g();
2478            eng.create_rule(overlap_rule(), &mut g).unwrap(); // jaccard 1/3 < 0.4 → no edges
2479            assert_eq!(g.topo.edge_count(), 0);
2480        }
2481        // b's tags change to overlap strongly
2482        let old = fx.props.get(b, "tags").cloned();
2483        fx.props.set(b, "tags", tags(&["x", "y"]));
2484        {
2485            let mut g = fx.g();
2486            eng.on_node_changed(b, Some(("tags", old)), &mut g);
2487            assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2488            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2489        }
2490        // and change away again → edges retract
2491        let old = fx.props.get(b, "tags").cloned();
2492        fx.props.set(b, "tags", tags(&["qqq"]));
2493        let mut g = fx.g();
2494        eng.on_node_changed(b, Some(("tags", old)), &mut g);
2495        assert_eq!(g.topo.edge_count(), 0);
2496        assert_eq!(g.edge_props.get(et, a, b, "score"), None);
2497    }
2498
2499    #[test]
2500    fn key_match_new_node_links_and_rebuild_is_noop() {
2501        let mut fx = Fx::new();
2502        fx.add("C", "c1", vec![]);
2503        let mut eng = RuleEngine::new();
2504        {
2505            let mut g = fx.g();
2506            eng.create_rule(
2507                RuleDef {
2508                    name: "fk".into(),
2509                    src_label: "T".into(),
2510                    dst_label: "C".into(),
2511                    predicate: Predicate::KeyMatch {
2512                        field: "cid".into(),
2513                    },
2514                    edge_type: "AT".into(),
2515                    weight_prop: None,
2516                    max_edges: None,
2517                    approximate: false,
2518                    via_label: None,
2519                    via_edge: None,
2520                    via_dir: None,
2521                },
2522                &mut g,
2523            )
2524            .unwrap();
2525        }
2526        let t = fx.add("T", "t1", vec![("cid", Value::Str("c1".into()))]);
2527        let (at, c1, count_before) = {
2528            let mut g = fx.g();
2529            eng.on_node_changed(t, None, &mut g);
2530            let at = g.syms.get("AT").unwrap();
2531            let c1 = g.ids.get("c1").unwrap();
2532            assert!(g.topo.neighbors(at, Direction::Out, t).contains(&c1));
2533            (at, c1, g.topo.edge_count())
2534        };
2535        let mut g = fx.g();
2536        eng.rebuild("fk", &mut g).unwrap();
2537        assert_eq!(g.topo.edge_count(), count_before); // rebuild is a no-op on consistent state
2538        assert!(g.topo.neighbors(at, Direction::Out, t).contains(&c1));
2539    }
2540
2541    #[test]
2542    fn score_refresh_on_persisting_owned_edge() {
2543        // Pins: weight set unconditionally even when add_edge returns false (edge persists).
2544        // jaccard({x,y,z},{x,y,q}) = |{x,y}|/|{x,y,z,q}| = 2/4 = 0.5 ≥ 0.2 → edges both ways.
2545        let mut fx = Fx::new();
2546        let a = fx.add("A", "a", vec![("tags", tags(&["x", "y", "z"]))]);
2547        let b = fx.add("A", "b", vec![("tags", tags(&["x", "y", "q"]))]);
2548        let et = fx.syms.intern("SIM");
2549        let mut eng = RuleEngine::new();
2550        {
2551            let mut g = fx.g();
2552            eng.create_rule(
2553                RuleDef {
2554                    name: "sim".into(),
2555                    src_label: "A".into(),
2556                    dst_label: "A".into(),
2557                    predicate: Predicate::Overlap {
2558                        field: "tags".into(),
2559                        min: 0.2,
2560                    },
2561                    edge_type: "SIM".into(),
2562                    weight_prop: Some("score".into()),
2563                    max_edges: None,
2564                    approximate: false,
2565                    via_label: None,
2566                    via_edge: None,
2567                    via_dir: None,
2568                },
2569                &mut g,
2570            )
2571            .unwrap();
2572            // Both directions present and owned with score ≈ 0.5.
2573            assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2574            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2575            assert!(eng.is_owned(et, a, b) || eng.is_owned(et, b, a));
2576            let check = |v: Option<&Value>| {
2577                if let Some(Value::Float(f)) = v {
2578                    assert!(
2579                        (f - 0.5).abs() < 1e-9,
2580                        "initial score should be 0.5, got {f}"
2581                    );
2582                }
2583            };
2584            check(g.edge_props.get(et, a, b, "score"));
2585            check(g.edge_props.get(et, b, a, "score"));
2586        }
2587        // Change b's tags to match a exactly → jaccard = 1.0.
2588        let old = fx.props.get(b, "tags").cloned();
2589        fx.props.set(b, "tags", tags(&["x", "y", "z"]));
2590        {
2591            let mut g = fx.g();
2592            eng.on_node_changed(b, Some(("tags", old)), &mut g);
2593            // Both directions still present.
2594            assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2595            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2596            // Scores must now be 1.0 on both directions.
2597            assert_eq!(
2598                g.edge_props.get(et, a, b, "score"),
2599                Some(&Value::Float(1.0)),
2600                "score on a→b must refresh to 1.0"
2601            );
2602            assert_eq!(
2603                g.edge_props.get(et, b, a, "score"),
2604                Some(&Value::Float(1.0)),
2605                "score on b→a must refresh to 1.0"
2606            );
2607        }
2608    }
2609
2610    #[test]
2611    fn dst_side_keymatch_links_when_c_node_inserted_after_t() {
2612        // Exercises the synthetic key-probe on src_side Scalar index (dst-side KeyMatch path).
2613        let mut fx = Fx::new();
2614        // Insert T node first with cid="c9" — no C node yet → no edge.
2615        let t = fx.add("T", "t1", vec![("cid", Value::Str("c9".into()))]);
2616        let mut eng = RuleEngine::new();
2617        {
2618            let mut g = fx.g();
2619            eng.create_rule(
2620                RuleDef {
2621                    name: "fk".into(),
2622                    src_label: "T".into(),
2623                    dst_label: "C".into(),
2624                    predicate: Predicate::KeyMatch {
2625                        field: "cid".into(),
2626                    },
2627                    edge_type: "AT".into(),
2628                    weight_prop: None,
2629                    max_edges: None,
2630                    approximate: false,
2631                    via_label: None,
2632                    via_edge: None,
2633                    via_dir: None,
2634                },
2635                &mut g,
2636            )
2637            .unwrap();
2638            // No C node → no edge.
2639            let at = g.syms.intern("AT");
2640            assert_eq!(g.topo.edge_count(), 0, "no C node yet → no edge");
2641            // t is indexed in src_side with Scalar{cid}="c9"
2642            let _ = at;
2643        }
2644        // Now insert C node "c9" and notify the engine.
2645        let c9 = fx.add("C", "c9", vec![]);
2646        {
2647            let mut g = fx.g();
2648            eng.on_node_changed(c9, None, &mut g);
2649            let at = g.syms.get("AT").unwrap();
2650            // The dst-side path must have probed src_side with key="c9" and found t.
2651            assert!(
2652                g.topo.neighbors(at, Direction::Out, t).contains(&c9),
2653                "T→C edge must appear when C node is inserted"
2654            );
2655            assert!(eng.is_owned(at, t, c9));
2656        }
2657    }
2658
2659    #[test]
2660    fn on_node_removed_retracts_both_sides_and_deindexes() {
2661        let mut fx = Fx::new();
2662        let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2663        let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2664        let et = fx.syms.intern("REL");
2665        let mut eng = RuleEngine::new();
2666        {
2667            let mut g = fx.g();
2668            eng.create_rule(overlap_rule(), &mut g).unwrap();
2669            assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2670            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2671        }
2672        {
2673            let mut g = fx.g();
2674            eng.on_node_removed(a, &mut g);
2675            assert!(!g.topo.neighbors(et, Direction::Out, a).contains(&b));
2676            assert!(!g.topo.neighbors(et, Direction::Out, b).contains(&a));
2677            assert_eq!(g.edge_props.get(et, a, b, "score"), None);
2678            assert_eq!(g.edge_props.get(et, b, a, "score"), None);
2679            assert!(!eng.is_owned(et, a, b));
2680            assert!(!eng.is_owned(et, b, a));
2681        }
2682        // Partner re-links to a NEW matching node; de-indexed a is not a candidate.
2683        let c = fx.add("A", "c", vec![("tags", tags(&["x", "y"]))]);
2684        {
2685            let mut g = fx.g();
2686            eng.on_node_changed(c, None, &mut g);
2687            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&c));
2688            assert!(g.topo.neighbors(et, Direction::Out, c).contains(&b));
2689            assert!(!g.topo.neighbors(et, Direction::Out, c).contains(&a));
2690            assert!(!g.topo.neighbors(et, Direction::Out, a).contains(&c));
2691        }
2692        // Second remove is a no-op (crash-window / already-retracted).
2693        {
2694            let mut g = fx.g();
2695            eng.on_node_removed(a, &mut g);
2696            assert!(g.topo.neighbors(et, Direction::Out, b).contains(&c));
2697        }
2698    }
2699
2700    #[test]
2701    fn duplicate_name_and_unknown_delete_error() {
2702        let mut fx = Fx::new();
2703        let mut eng = RuleEngine::new();
2704        let mut g = fx.g();
2705        eng.create_rule(overlap_rule(), &mut g).unwrap();
2706        assert!(eng.create_rule(overlap_rule(), &mut g).is_err());
2707        assert!(eng.delete_rule("nope", &mut g).is_err());
2708    }
2709
2710    /// C1: two rules sharing the same edge_type both match a pair of nodes.
2711    /// During backfill of R2, add_edge returns false for edges R1 already owns,
2712    /// so R2's provenance lacks them.  Deleting R1 removes those edges from the
2713    /// topology — but the rebuild-survivors step must then re-run R2 so it claims
2714    /// them.  Deleting R2 afterward must actually remove the edge.
2715    #[test]
2716    fn coowned_edge_type_survives_first_delete_gone_after_second() {
2717        let mut fx = Fx::new();
2718        let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2719        let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2720        let mut eng = RuleEngine::new();
2721        {
2722            let mut g = fx.g();
2723            // R1: Overlap min=0.1 — derives a↔b (jaccard 1.0 ≥ 0.1).
2724            eng.create_rule(
2725                RuleDef {
2726                    name: "r1".into(),
2727                    src_label: "A".into(),
2728                    dst_label: "A".into(),
2729                    predicate: Predicate::Overlap {
2730                        field: "tags".into(),
2731                        min: 0.1,
2732                    },
2733                    edge_type: "REL2".into(),
2734                    weight_prop: None,
2735                    max_edges: None,
2736                    approximate: false,
2737                    via_label: None,
2738                    via_edge: None,
2739                    via_dir: None,
2740                },
2741                &mut g,
2742            )
2743            .unwrap();
2744            // R2: same edge_type, Overlap min=0.2 — also derives a↔b.
2745            eng.create_rule(
2746                RuleDef {
2747                    name: "r2".into(),
2748                    src_label: "A".into(),
2749                    dst_label: "A".into(),
2750                    predicate: Predicate::Overlap {
2751                        field: "tags".into(),
2752                        min: 0.2,
2753                    },
2754                    edge_type: "REL2".into(),
2755                    weight_prop: None,
2756                    max_edges: None,
2757                    approximate: false,
2758                    via_label: None,
2759                    via_edge: None,
2760                    via_dir: None,
2761                },
2762                &mut g,
2763            )
2764            .unwrap();
2765
2766            let et = g.syms.intern("REL2");
2767            // Both directions must exist (either rule claims them).
2768            assert!(
2769                g.topo.neighbors(et, Direction::Out, a).contains(&b),
2770                "a→b must exist after both rules created"
2771            );
2772            assert!(
2773                g.topo.neighbors(et, Direction::Out, b).contains(&a),
2774                "b→a must exist after both rules created"
2775            );
2776
2777            // Delete R1 — rebuild-survivors re-runs R2 which must reclaim the edges.
2778            eng.delete_rule("r1", &mut g).unwrap();
2779            assert!(
2780                g.topo.neighbors(et, Direction::Out, a).contains(&b),
2781                "a→b must survive R1 deletion (R2 rebuilds and claims it)"
2782            );
2783            assert!(
2784                g.topo.neighbors(et, Direction::Out, b).contains(&a),
2785                "b→a must survive R1 deletion (R2 rebuilds and claims it)"
2786            );
2787            // R2 now owns both directions.
2788            assert!(
2789                eng.is_owned(et, a, b),
2790                "a→b must be owned by R2 after rebuild"
2791            );
2792            assert!(
2793                eng.is_owned(et, b, a),
2794                "b→a must be owned by R2 after rebuild"
2795            );
2796
2797            // Delete R2 — no survivor left, edges must be gone.
2798            eng.delete_rule("r2", &mut g).unwrap();
2799            assert!(
2800                !g.topo.neighbors(et, Direction::Out, a).contains(&b),
2801                "a→b must be gone after both rules deleted"
2802            );
2803            assert!(
2804                !g.topo.neighbors(et, Direction::Out, b).contains(&a),
2805                "b→a must be gone after both rules deleted"
2806            );
2807        }
2808    }
2809
2810    /// Helper: FieldEqual rule with top-k per-source cap.
2811    fn topk_eq_rule(k: u64) -> RuleDef {
2812        RuleDef {
2813            name: "eq".into(),
2814            src_label: "N".into(),
2815            dst_label: "N".into(),
2816            predicate: Predicate::FieldEqual { field: "k".into() },
2817            edge_type: "EQ".into(),
2818            weight_prop: None,
2819            max_edges: Some(k),
2820            approximate: false,
2821            via_label: None,
2822            via_edge: None,
2823            via_dir: None,
2824        }
2825    }
2826
2827    fn prov_pairs(eng: &RuleEngine, name: &str) -> BTreeSet<(u32, u32)> {
2828        eng.provenance()
2829            .get(name)
2830            .map(|s| s.iter().map(|&(_, a, b)| (a, b)).collect())
2831            .unwrap_or_default()
2832    }
2833
2834    /// k=1: each src gets its single best-scored dst (score DESC, key ASC
2835    /// tiebreak).  FieldEqual has uniform score 1.0, so the winner is the dst
2836    /// with the lexicographically smallest key that is not the src itself.
2837    #[test]
2838    fn topk_k1_keeps_best_scored_dst() {
2839        let mut fx = Fx::new();
2840        let mut eng = RuleEngine::new();
2841        {
2842            let mut g = fx.g();
2843            eng.create_rule(topk_eq_rule(1), &mut g).unwrap();
2844        }
2845        // Insert 4 nodes all sharing k="const".  Keys: n0 < n1 < n2 < n3.
2846        let mut ids = Vec::new();
2847        for i in 0..4usize {
2848            let id = fx.add(
2849                "N",
2850                &format!("n{i}"),
2851                vec![("k", Value::Str("const".into()))],
2852            );
2853            ids.push(id);
2854            let mut g = fx.g();
2855            eng.on_node_changed(id, None, &mut g);
2856        }
2857        let et = fx.syms.get("EQ").unwrap();
2858        // Each src's single allowed dst must be the smallest key ≠ self.
2859        // n0 → n1 (smallest other)
2860        // n1 → n0 (n0 < n1)
2861        // n2 → n0
2862        // n3 → n0
2863        let expected_dsts = [ids[1], ids[0], ids[0], ids[0]];
2864        for (i, (&src, &expected_dst)) in ids.iter().zip(expected_dsts.iter()).enumerate() {
2865            let out: Vec<u32> = fx.topo.neighbors(et, Direction::Out, src).to_vec();
2866            assert_eq!(
2867                out,
2868                vec![expected_dst],
2869                "src n{i} should point only to the best dst"
2870            );
2871        }
2872        assert_eq!(eng.provenance()["eq"].len(), 4);
2873        assert!(!eng.is_tripped("eq"), "top-k rules never trip");
2874    }
2875
2876    /// k=2 insert-evict: adding a better dst evicts the worst of the current k.
2877    /// Uses NumericWithin (scored) so scores differ across dsts.
2878    #[test]
2879    fn topk_insert_evict() {
2880        // Rule: S→D with VectorSimilar-alike (we use NumericWithin for simplicity).
2881        // 3 src nodes, numeric field "v"; tolerance 10.0 so score = 1-|Δ|/10.
2882        // k=1 per source.
2883        let mut fx = Fx::new();
2884        let rule = RuleDef {
2885            name: "nw".into(),
2886            src_label: "S".into(),
2887            dst_label: "D".into(),
2888            predicate: Predicate::NumericWithin {
2889                field: "v".into(),
2890                tolerance: 10.0,
2891            },
2892            edge_type: "NEAR".into(),
2893            weight_prop: Some("score".into()),
2894            max_edges: Some(1),
2895            approximate: false,
2896            via_label: None,
2897            via_edge: None,
2898            via_dir: None,
2899        };
2900        let mut eng = RuleEngine::new();
2901        {
2902            let mut g = fx.g();
2903            eng.create_rule(rule, &mut g).unwrap();
2904        }
2905
2906        // src s0 with v=0.0
2907        let s0 = fx.add("S", "s0", vec![("v", Value::Float(0.0))]);
2908        // dst d_far with v=9.0 → score=0.1 (worst)
2909        let d_far = fx.add("D", "d_far", vec![("v", Value::Float(9.0))]);
2910        {
2911            let mut g = fx.g();
2912            eng.on_node_changed(s0, None, &mut g);
2913            eng.on_node_changed(d_far, None, &mut g);
2914        }
2915        let et = fx.syms.get("NEAR").unwrap();
2916        // s0 → d_far (only candidate)
2917        assert!(fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2918        assert_eq!(eng.provenance()["nw"].len(), 1);
2919
2920        // Insert d_close with v=1.0 → score=0.9 (better than d_far).
2921        let d_close = fx.add("D", "d_close", vec![("v", Value::Float(1.0))]);
2922        {
2923            let mut g = fx.g();
2924            eng.on_node_changed(d_close, None, &mut g);
2925        }
2926        // s0 should now point to d_close (evicting d_far).
2927        let out: Vec<u32> = fx.topo.neighbors(et, Direction::Out, s0).to_vec();
2928        assert_eq!(out, vec![d_close], "d_close should evict d_far");
2929        assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2930        assert_eq!(eng.provenance()["nw"].len(), 1);
2931        assert!(eng.by_node_consistent());
2932    }
2933
2934    /// Retract-backfill: removing the best dst causes the next-best to fill in.
2935    #[test]
2936    fn topk_retract_backfill() {
2937        let mut fx = Fx::new();
2938        let rule = RuleDef {
2939            name: "nw".into(),
2940            src_label: "S".into(),
2941            dst_label: "D".into(),
2942            predicate: Predicate::NumericWithin {
2943                field: "v".into(),
2944                tolerance: 10.0,
2945            },
2946            edge_type: "NEAR".into(),
2947            weight_prop: Some("score".into()),
2948            max_edges: Some(1),
2949            approximate: false,
2950            via_label: None,
2951            via_edge: None,
2952            via_dir: None,
2953        };
2954        let mut eng = RuleEngine::new();
2955
2956        let s0 = fx.add("S", "s0", vec![("v", Value::Float(0.0))]);
2957        let d_close = fx.add("D", "d_close", vec![("v", Value::Float(1.0))]); // score=0.9
2958        let d_far = fx.add("D", "d_far", vec![("v", Value::Float(8.0))]); // score=0.2
2959        {
2960            let mut g = fx.g();
2961            eng.create_rule(rule, &mut g).unwrap();
2962        }
2963        let et = fx.syms.get("NEAR").unwrap();
2964        // d_close is the top-1 dst.
2965        assert!(fx.topo.neighbors(et, Direction::Out, s0).contains(&d_close));
2966        assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2967        assert_eq!(eng.provenance()["nw"].len(), 1);
2968
2969        // Break d_close's match by pushing its v out of tolerance.
2970        let old = fx.props.get(d_close, "v").cloned();
2971        fx.props.set(d_close, "v", Value::Float(50.0));
2972        {
2973            let mut g = fx.g();
2974            eng.on_node_changed(d_close, Some(("v", old)), &mut g);
2975        }
2976        // d_far should backfill.
2977        assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_close));
2978        assert!(
2979            fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far),
2980            "d_far should backfill after d_close retracted"
2981        );
2982        assert_eq!(eng.provenance()["nw"].len(), 1);
2983        assert!(eng.by_node_consistent());
2984    }
2985
2986    /// Tie-breaking: equal scores → dst_key ASC wins.
2987    #[test]
2988    fn topk_tie_broken_by_dst_key() {
2989        // FieldEqual: all dsts have score 1.0 → tiebreak by key.
2990        let mut fx = Fx::new();
2991        let mut eng = RuleEngine::new();
2992        {
2993            let mut g = fx.g();
2994            eng.create_rule(topk_eq_rule(2), &mut g).unwrap();
2995        }
2996        // 5 nodes all with k="x" → each src matches 4 others; top-2 by key.
2997        // Keys: a, b, c, d, e (alphabetical).
2998        for name in ["a", "b", "c", "d", "e"] {
2999            let id = fx.add("N", name, vec![("k", Value::Str("x".into()))]);
3000            let mut g = fx.g();
3001            eng.on_node_changed(id, None, &mut g);
3002        }
3003        let et = fx.syms.get("EQ").unwrap();
3004        let get_id = |key: &str| fx.ids.get(key).unwrap();
3005        // Node "a" should point to the two smallest keys that aren't "a": b, c.
3006        let a = get_id("a");
3007        let b = get_id("b");
3008        let c = get_id("c");
3009        let out_a: BTreeSet<u32> = fx
3010            .topo
3011            .neighbors(et, Direction::Out, a)
3012            .iter()
3013            .copied()
3014            .collect();
3015        assert!(out_a.contains(&b), "a→b (b is best key after a)");
3016        assert!(out_a.contains(&c), "a→c (c is 2nd best key)");
3017        assert_eq!(out_a.len(), 2);
3018        // Node "e" should point to "a" and "b" (two smallest keys ≠ "e").
3019        let e = get_id("e");
3020        let out_e: BTreeSet<u32> = fx
3021            .topo
3022            .neighbors(et, Direction::Out, e)
3023            .iter()
3024            .copied()
3025            .collect();
3026        assert!(out_e.contains(&a), "e→a");
3027        assert!(out_e.contains(&b), "e→b");
3028        assert_eq!(out_e.len(), 2);
3029        assert!(eng.by_node_consistent());
3030    }
3031
3032    /// When k >= candidate count, all candidates are included (no truncation).
3033    #[test]
3034    fn topk_k_larger_than_candidate_count() {
3035        let mut fx = Fx::new();
3036        let mut eng = RuleEngine::new();
3037        {
3038            let mut g = fx.g();
3039            // k=100 but only 3 other nodes → all 3 included.
3040            eng.create_rule(topk_eq_rule(100), &mut g).unwrap();
3041        }
3042        for i in 0..4usize {
3043            let id = fx.add("N", &format!("n{i}"), vec![("k", Value::Str("c".into()))]);
3044            let mut g = fx.g();
3045            eng.on_node_changed(id, None, &mut g);
3046        }
3047        // 4 nodes × 3 matches each = 12 directed edges.
3048        assert_eq!(eng.provenance()["eq"].len(), 12);
3049        assert!(!eng.is_tripped("eq"));
3050    }
3051
3052    /// rebuild() with top-k rule re-converges to the correct per-source top-k
3053    /// after externally removing a node's field.
3054    #[test]
3055    fn topk_rebuild_exact() {
3056        let mut fx = Fx::new();
3057        let mut eng = RuleEngine::new();
3058        {
3059            let mut g = fx.g();
3060            eng.create_rule(topk_eq_rule(1), &mut g).unwrap();
3061        }
3062        // 3 nodes with k="x" → each gets 1 dst (smallest key ≠ self).
3063        let _a = fx.add("N", "a", vec![("k", Value::Str("x".into()))]);
3064        let _b = fx.add("N", "b", vec![("k", Value::Str("x".into()))]);
3065        let _c = fx.add("N", "c", vec![("k", Value::Str("x".into()))]);
3066        {
3067            let mut g = fx.g();
3068            eng.on_node_changed(_a, None, &mut g);
3069            eng.on_node_changed(_b, None, &mut g);
3070            eng.on_node_changed(_c, None, &mut g);
3071        }
3072        assert_eq!(eng.provenance()["eq"].len(), 3);
3073
3074        // rebuild should produce the same result.
3075        {
3076            let mut g = fx.g();
3077            eng.rebuild("eq", &mut g).unwrap();
3078        }
3079        assert_eq!(eng.provenance()["eq"].len(), 3);
3080        assert!(!eng.is_tripped("eq"));
3081        assert!(eng.by_node_consistent());
3082    }
3083
3084    /// by_node index stays consistent across top-k inserts, evictions and rebuild.
3085    #[test]
3086    fn topk_by_node_consistent() {
3087        let mut fx = Fx::new();
3088        let mut eng = RuleEngine::new();
3089        {
3090            let mut g = fx.g();
3091            eng.create_rule(topk_eq_rule(2), &mut g).unwrap();
3092        }
3093        for i in 0..5usize {
3094            let id = fx.add(
3095                "N",
3096                &format!("n{i}"),
3097                vec![("k", Value::Str("const".into()))],
3098            );
3099            let mut g = fx.g();
3100            eng.on_node_changed(id, None, &mut g);
3101        }
3102        assert!(eng.by_node_consistent(), "consistent after insertions");
3103
3104        // Evict by changing a prop.
3105        let id2 = fx.ids.get("n2").unwrap();
3106        let old = fx.props.get(id2, "k").cloned();
3107        fx.props.set(id2, "k", Value::Str("other".into()));
3108        {
3109            let mut g = fx.g();
3110            eng.on_node_changed(id2, Some(("k", old)), &mut g);
3111        }
3112        assert!(eng.by_node_consistent(), "consistent after eviction");
3113
3114        {
3115            let mut g = fx.g();
3116            eng.rebuild("eq", &mut g).unwrap();
3117        }
3118        assert!(eng.by_node_consistent(), "consistent after rebuild");
3119    }
3120
3121    fn numeric_rule() -> RuleDef {
3122        RuleDef {
3123            name: "nw".into(),
3124            src_label: "C".into(),
3125            dst_label: "C".into(),
3126            predicate: Predicate::NumericWithin {
3127                field: "year".into(),
3128                tolerance: 2.0,
3129            },
3130            edge_type: "NEAR".into(),
3131            weight_prop: Some("score".into()),
3132            max_edges: None,
3133            approximate: false,
3134            via_label: None,
3135            via_edge: None,
3136            via_dir: None,
3137        }
3138    }
3139
3140    fn geo_rule() -> RuleDef {
3141        RuleDef {
3142            name: "geo".into(),
3143            src_label: "City".into(),
3144            dst_label: "City".into(),
3145            predicate: Predicate::GeoRadius {
3146                field: "loc".into(),
3147                km: 400.0,
3148            },
3149            edge_type: "NEAR_GEO".into(),
3150            weight_prop: Some("score".into()),
3151            max_edges: None,
3152            approximate: false,
3153            via_label: None,
3154            via_edge: None,
3155            via_dir: None,
3156        }
3157    }
3158
3159    fn vec_rule() -> RuleDef {
3160        RuleDef {
3161            name: "vec".into(),
3162            src_label: "Doc".into(),
3163            dst_label: "Doc".into(),
3164            predicate: Predicate::VectorSimilar {
3165                field: "emb".into(),
3166                min: 0.9,
3167            },
3168            edge_type: "SIM".into(),
3169            weight_prop: Some("score".into()),
3170            max_edges: None,
3171            approximate: false,
3172            via_label: None,
3173            via_edge: None,
3174            via_dir: None,
3175        }
3176    }
3177
3178    fn pair_edges(topo: &Topology, et: u32, a: u32, b: u32) -> bool {
3179        topo.neighbors(et, Direction::Out, a).contains(&b)
3180            && topo.neighbors(et, Direction::Out, b).contains(&a)
3181    }
3182
3183    #[test]
3184    fn numeric_within_incremental_crosses_bucket_and_clears_old_index() {
3185        let mut fx = Fx::new();
3186        let a = fx.add("C", "a", vec![("year", Value::Float(10.0))]);
3187        let b = fx.add("C", "b", vec![("year", Value::Float(12.0))]);
3188        let et = fx.syms.intern("NEAR");
3189        let mut eng = RuleEngine::new();
3190        {
3191            let mut g = fx.g();
3192            eng.create_rule(numeric_rule(), &mut g).unwrap();
3193            // |12−10| = 2 ≤ 2 → score 0.0 both ways
3194            assert!(pair_edges(g.topo, et, a, b));
3195        }
3196
3197        // 12.0 (bucket 6) → 16.1 (bucket 8): two buckets away, so the old
3198        // value's ±1 probe no longer reaches b. Match breaks.
3199        let old = fx.props.get(b, "year").cloned();
3200        fx.props.set(b, "year", Value::Float(16.1));
3201        {
3202            let mut g = fx.g();
3203            eng.on_node_changed(b, Some(("year", old)), &mut g);
3204            assert!(!pair_edges(g.topo, et, a, b));
3205            assert_eq!(g.topo.edge_count(), 0);
3206        }
3207        let def = numeric_rule();
3208        let spec = candidate_spec_for(&def);
3209        let old_map: std::collections::HashMap<_, _> =
3210            [("year".to_string(), Value::Float(12.0))].into();
3211        let old_get = |f: &str| old_map.get(f).cloned();
3212        let src_hits = eng.indexes["nw"].src_side.candidates(&spec, &old_get);
3213        let dst_hits = eng.indexes["nw"].dst_side.candidates(&spec, &old_get);
3214        assert!(!src_hits.contains(&b), "old src bucket must drop b");
3215        assert!(!dst_hits.contains(&b), "old dst bucket must drop b");
3216        assert!(src_hits.contains(&a));
3217
3218        // 16.1 → 11.9 (bucket 5): match returns.
3219        let old = fx.props.get(b, "year").cloned();
3220        fx.props.set(b, "year", Value::Float(11.9));
3221        let mut g = fx.g();
3222        eng.on_node_changed(b, Some(("year", old)), &mut g);
3223        assert!(pair_edges(g.topo, et, a, b));
3224    }
3225
3226    fn loc_val(lat: f64, lon: f64) -> Value {
3227        Value::List(vec![Value::Float(lat), Value::Float(lon)])
3228    }
3229
3230    fn emb_val(vals: &[f64]) -> Value {
3231        Value::List(vals.iter().copied().map(Value::Float).collect())
3232    }
3233
3234    #[test]
3235    fn rebuild_is_noop_for_numeric_geo_and_vector() {
3236        let mut fx = Fx::new();
3237        let ca = fx.add("C", "ca", vec![("year", Value::Int(1998))]);
3238        let cb = fx.add("C", "cb", vec![("year", Value::Float(2000.0))]);
3239        let pa = fx.add("City", "paris", vec![("loc", loc_val(48.8566, 2.3522))]);
3240        let lo = fx.add("City", "london", vec![("loc", loc_val(51.5074, -0.1278))]);
3241        let da = fx.add("Doc", "d1", vec![("emb", emb_val(&[1.0, 0.0]))]);
3242        let db = fx.add("Doc", "d2", vec![("emb", emb_val(&[1.0, 0.0]))]);
3243
3244        let mut eng = RuleEngine::new();
3245        {
3246            let mut g = fx.g();
3247            eng.create_rule(numeric_rule(), &mut g).unwrap();
3248            eng.create_rule(geo_rule(), &mut g).unwrap();
3249            eng.create_rule(vec_rule(), &mut g).unwrap();
3250        }
3251
3252        let (near, ngeo, sim) = (
3253            fx.syms.get("NEAR").unwrap(),
3254            fx.syms.get("NEAR_GEO").unwrap(),
3255            fx.syms.get("SIM").unwrap(),
3256        );
3257        assert!(pair_edges(&fx.topo, near, ca, cb));
3258        assert!(pair_edges(&fx.topo, ngeo, pa, lo));
3259        assert!(pair_edges(&fx.topo, sim, da, db));
3260        let before = fx.topo.edge_count();
3261
3262        {
3263            let mut g = fx.g();
3264            eng.rebuild("nw", &mut g).unwrap();
3265            eng.rebuild("geo", &mut g).unwrap();
3266            eng.rebuild("vec", &mut g).unwrap();
3267        }
3268        assert_eq!(fx.topo.edge_count(), before);
3269        assert!(pair_edges(&fx.topo, near, ca, cb));
3270        assert!(pair_edges(&fx.topo, ngeo, pa, lo));
3271        assert!(pair_edges(&fx.topo, sim, da, db));
3272    }
3273
3274    fn fk_rule() -> RuleDef {
3275        RuleDef {
3276            name: "works_at".into(),
3277            src_label: "T".into(),
3278            dst_label: "C".into(),
3279            predicate: Predicate::KeyMatch {
3280                field: "cid".into(),
3281            },
3282            edge_type: "AT".into(),
3283            weight_prop: None,
3284            max_edges: None,
3285            approximate: false,
3286            via_label: None,
3287            via_edge: None,
3288            via_dir: None,
3289        }
3290    }
3291
3292    #[test]
3293    fn by_node_matches_rebuild_after_mutation_storm() {
3294        let mut fx = Fx::new();
3295        let hub = fx.add("C", "hub", vec![]);
3296        let other = fx.add("C", "other", vec![]);
3297        let mut people = Vec::new();
3298        for i in 0..40 {
3299            let cid = if i < 30 { "hub" } else { "other" };
3300            people.push(fx.add(
3301                "T",
3302                &format!("t{i}"),
3303                vec![("cid", Value::Str(cid.into())), ("tags", tags(&["x", "y"]))],
3304            ));
3305        }
3306        let mut overlap = overlap_rule();
3307        overlap.src_label = "T".into();
3308        overlap.dst_label = "T".into();
3309        let mut eng = RuleEngine::new();
3310        {
3311            let mut g = fx.g();
3312            eng.create_rule(fk_rule(), &mut g).unwrap();
3313            eng.create_rule(overlap, &mut g).unwrap();
3314        }
3315        assert!(eng.by_node_consistent());
3316        assert_eq!(eng.provenance_touching_len(hub), 30);
3317
3318        // Incremental: re-home half the hub people, flip tags, then restore.
3319        for (i, &id) in people.iter().enumerate().take(15) {
3320            let old = fx.props.get(id, "cid").cloned();
3321            fx.props.set(id, "cid", Value::Str("other".into()));
3322            let mut g = fx.g();
3323            eng.on_node_changed(id, Some(("cid", old)), &mut g);
3324            assert!(
3325                eng.by_node_consistent(),
3326                "inconsistent after cid update {i}"
3327            );
3328        }
3329        for &id in people.iter().take(8) {
3330            let old = fx.props.get(id, "tags").cloned();
3331            fx.props.set(id, "tags", tags(&["q"]));
3332            let mut g = fx.g();
3333            eng.on_node_changed(id, Some(("tags", old)), &mut g);
3334        }
3335        assert!(eng.by_node_consistent());
3336
3337        // Delete-node cleanup uses the reverse index.
3338        {
3339            let mut g = fx.g();
3340            eng.on_node_removed(people[0], &mut g);
3341        }
3342        fx.labels[people[0] as usize] = u32::MAX;
3343        assert!(eng.by_node_consistent());
3344        assert_eq!(eng.provenance_touching_len(people[0]), 0);
3345
3346        {
3347            let mut g = fx.g();
3348            eng.rebuild("works_at", &mut g).unwrap();
3349            eng.rebuild("rel", &mut g).unwrap();
3350        }
3351        assert!(eng.by_node_consistent());
3352
3353        {
3354            let mut g = fx.g();
3355            eng.delete_rule("rel", &mut g).unwrap();
3356        }
3357        assert!(eng.by_node_consistent());
3358        assert_eq!(eng.provenance_touching(people[1]).count(), 1);
3359
3360        // Persist-restore rebuilds the reverse index from provenance.
3361        let (defs, prov, tripped, fires) = eng.to_persist();
3362        let restored = RuleEngine::from_persist(defs, prov, tripped, fires);
3363        assert!(restored.by_node_consistent());
3364        assert_eq!(
3365            restored.provenance_touching_len(hub),
3366            eng.provenance_touching_len(hub)
3367        );
3368        assert_eq!(
3369            restored.provenance_touching_len(other),
3370            eng.provenance_touching_len(other)
3371        );
3372    }
3373
3374    #[test]
3375    fn provenance_touching_high_degree_hub() {
3376        let mut fx = Fx::new();
3377        let hub = fx.add("C", "hub", vec![]);
3378        let mut first = None;
3379        for i in 0..256 {
3380            let id = fx.add(
3381                "T",
3382                &format!("t{i}"),
3383                vec![("cid", Value::Str("hub".into()))],
3384            );
3385            if first.is_none() {
3386                first = Some(id);
3387            }
3388        }
3389        let first = first.unwrap();
3390        let mut eng = RuleEngine::new();
3391        {
3392            let mut g = fx.g();
3393            eng.create_rule(fk_rule(), &mut g).unwrap();
3394        }
3395        assert!(eng.by_node_consistent());
3396        assert_eq!(eng.provenance_touching_len(hub), 256);
3397        assert_eq!(eng.provenance_touching_len(first), 1);
3398        let hits: Vec<_> = eng.provenance_touching(first).collect();
3399        assert_eq!(hits.len(), 1);
3400        assert_eq!(hits[0].0, "works_at");
3401        assert_eq!(hits[0].2, first);
3402        assert_eq!(hits[0].3, hub);
3403    }
3404
3405    /// by_node index stays consistent across global-budget trip and rebuild
3406    /// (max_edges: None path — DEFAULT_MAX_EDGES = 1_000_000).
3407    ///
3408    /// Uses a tiny budget via a special rule with `max_edges: None` but many
3409    /// nodes to naturally exceed the default; instead we directly test the
3410    /// None-path by verifying that the by_node index is consistent at each
3411    /// step of normal insertions and rebuilds.
3412    #[test]
3413    fn by_node_consistent_across_inserts_and_rebuild() {
3414        let mut fx = Fx::new();
3415        let mut eng = RuleEngine::new();
3416        let rule = RuleDef {
3417            name: "eq".into(),
3418            src_label: "N".into(),
3419            dst_label: "N".into(),
3420            predicate: Predicate::FieldEqual { field: "k".into() },
3421            edge_type: "EQ".into(),
3422            weight_prop: None,
3423            max_edges: None, // global-budget path, DEFAULT_MAX_EDGES = 1_000_000
3424            approximate: false,
3425            via_label: None,
3426            via_edge: None,
3427            via_dir: None,
3428        };
3429        {
3430            let mut g = fx.g();
3431            eng.create_rule(rule, &mut g).unwrap();
3432        }
3433        let mut ids = Vec::new();
3434        for i in 0..6 {
3435            let id = fx.add(
3436                "N",
3437                &format!("n{i}"),
3438                vec![("k", Value::Str("const".into()))],
3439            );
3440            ids.push(id);
3441            let mut g = fx.g();
3442            eng.on_node_changed(id, None, &mut g);
3443        }
3444        // 6 nodes × 5 matches each = 30 directed edges (well under 1M budget).
3445        assert_eq!(eng.provenance()["eq"].len(), 30);
3446        assert!(!eng.is_tripped("eq"));
3447        assert!(eng.by_node_consistent(), "consistent after insertions");
3448
3449        // Change one node's field — triggers retract + backfill on that src.
3450        let old = fx.props.get(ids[3], "k").cloned();
3451        fx.props.set(ids[3], "k", Value::Str("other".into()));
3452        {
3453            let mut g = fx.g();
3454            eng.on_node_changed(ids[3], Some(("k", old)), &mut g);
3455        }
3456        assert!(eng.by_node_consistent(), "consistent after property change");
3457
3458        {
3459            let mut g = fx.g();
3460            eng.rebuild("eq", &mut g).unwrap();
3461        }
3462        assert!(!eng.is_tripped("eq"));
3463        assert!(eng.by_node_consistent(), "consistent after rebuild");
3464    }
3465
3466    fn mix64(mut x: u64) -> u64 {
3467        x = x.wrapping_add(0x9E3779B97F4A7C15);
3468        x = (x ^ (x >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
3469        x = (x ^ (x >> 27)).wrapping_mul(0x94D049BB133111EB);
3470        x ^ (x >> 31)
3471    }
3472
3473    fn rand_emb(seed: u64, i: u32, dim: usize) -> Value {
3474        let vals: Vec<f64> = (0..dim)
3475            .map(|d| {
3476                let bits = mix64(seed ^ ((i as u64 + 1).wrapping_mul(0x100000001)) ^ (d as u64));
3477                let mut f = (bits as f64) / (u64::MAX as f64) * 2.0 - 1.0;
3478                if f == 0.0 {
3479                    f = 1.0;
3480                }
3481                f
3482            })
3483            .collect();
3484        emb_val(&vals)
3485    }
3486
3487    fn seed_docs(n: u32, seed: u64) -> (Fx, Vec<u32>) {
3488        let dims = [2usize, 3, 4, 8];
3489        let mut fx = Fx::new();
3490        let mut ids = Vec::new();
3491        for i in 0..n {
3492            let dim = dims[(i as usize) % dims.len()];
3493            ids.push(fx.add(
3494                "Doc",
3495                &format!("d{i}"),
3496                vec![("emb", rand_emb(seed, i, dim))],
3497            ));
3498        }
3499        (fx, ids)
3500    }
3501
3502    /// Identity proof: 500 mixed-dim vectors, derived edges with the dim
3503    /// reject on vs forced off (and vs brute-force evaluate) are identical.
3504    #[test]
3505    fn vector_dim_reject_matches_unfiltered_and_oracle() {
3506        const N: u32 = 500;
3507        const SEED: u64 = 0xC0FF_EE00_D15C;
3508        let def = vec_rule();
3509
3510        let (mut fx_on, ids) = seed_docs(N, SEED);
3511        let mut eng_on = RuleEngine::new();
3512        {
3513            let mut g = fx_on.g();
3514            eng_on.create_rule(def.clone(), &mut g).unwrap();
3515        }
3516        let on = prov_pairs(&eng_on, "vec");
3517        assert!(!on.is_empty(), "seeded set must produce some edges");
3518
3519        let (mut fx_off, _) = seed_docs(N, SEED);
3520        let mut eng_off = RuleEngine::new();
3521        {
3522            let mut g = fx_off.g();
3523            with_vector_dim_reject(false, || {
3524                eng_off.create_rule(def.clone(), &mut g).unwrap();
3525            });
3526        }
3527        assert_eq!(on, prov_pairs(&eng_off, "vec"), "filter vs no-filter");
3528
3529        let mut brute = BTreeSet::new();
3530        for &s in &ids {
3531            for &d in &ids {
3532                if s == d {
3533                    continue;
3534                }
3535                let skey = fx_on.ids.key_of(s).unwrap();
3536                let dkey = fx_on.ids.key_of(d).unwrap();
3537                let sget = |f: &str| fx_on.props.get(s, f).cloned();
3538                let dget = |f: &str| fx_on.props.get(d, f).cloned();
3539                if evaluate(
3540                    &def.predicate,
3541                    &NodeView {
3542                        key: skey,
3543                        props: &sget,
3544                    },
3545                    &NodeView {
3546                        key: dkey,
3547                        props: &dget,
3548                    },
3549                )
3550                .is_some()
3551                {
3552                    brute.insert((s, d));
3553                }
3554            }
3555        }
3556        assert_eq!(on, brute, "filter vs brute-force evaluate");
3557    }
3558
3559    /// Dim change must flow through remove(old)+insert(new); edges match a
3560    /// fresh engine built from the post-update props.
3561    #[test]
3562    fn vector_dim_change_updates_cache_and_matches_fresh_build() {
3563        let mut fx = Fx::new();
3564        let a = fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0]))]);
3565        let b = fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0]))]);
3566        let c = fx.add("Doc", "c", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3567        let mut eng = RuleEngine::new();
3568        {
3569            let mut g = fx.g();
3570            eng.create_rule(vec_rule(), &mut g).unwrap();
3571        }
3572        assert_eq!(eng.indexes["vec"].src_side.vec_dim(a), Some(2));
3573        assert_eq!(eng.indexes["vec"].src_side.vec_dim(c), Some(3));
3574        assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(a, b), (b, a)]));
3575
3576        let old = fx.props.get(b, "emb").cloned();
3577        fx.props.set(b, "emb", emb_val(&[1.0, 0.0, 0.0]));
3578        {
3579            let mut g = fx.g();
3580            eng.on_node_changed(b, Some(("emb", old)), &mut g);
3581        }
3582        assert_eq!(eng.indexes["vec"].src_side.vec_dim(b), Some(3));
3583        assert_eq!(eng.indexes["vec"].dst_side.vec_dim(b), Some(3));
3584        let after = prov_pairs(&eng, "vec");
3585        assert_eq!(after, BTreeSet::from([(b, c), (c, b)]));
3586
3587        // Separate graph: first engine already owns the b↔c edges in `fx.topo`.
3588        let mut fresh_fx = Fx::new();
3589        let fa = fresh_fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0]))]);
3590        let fb = fresh_fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3591        let fc = fresh_fx.add("Doc", "c", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3592        let mut fresh = RuleEngine::new();
3593        {
3594            let mut g = fresh_fx.g();
3595            fresh.create_rule(vec_rule(), &mut g).unwrap();
3596        }
3597        assert_eq!(
3598            prov_pairs(&fresh, "vec"),
3599            BTreeSet::from([(fb, fc), (fc, fb)])
3600        );
3601        assert_eq!(fresh.indexes["vec"].src_side.vec_dim(fb), Some(3));
3602        assert_eq!(fresh.indexes["vec"].src_side.vec_dim(fa), Some(2));
3603    }
3604
3605    // -----------------------------------------------------------------------
3606    // Streaming backfill — order-identity and memory-bound tests (Plan 11 M1)
3607    // -----------------------------------------------------------------------
3608
3609    /// Top-k order-identity property test.
3610    ///
3611    /// For rules with `max_edges: Some(k)` (top-k per-source semantics),
3612    /// verifies that `create_rule` streaming backfill produces the same
3613    /// per-source top-k set as an independent brute-force reference.
3614    ///
3615    /// The reference is intentionally independent of `filter_src_top_k`:
3616    /// it sorts candidates inline (score DESC, dst-key ASC, take k) so a
3617    /// comparator bug cannot self-agree between reference and actual.
3618    ///
3619    /// Covers all four `CandidateSpec` paths through `compute_desired`:
3620    /// - `FieldEqual` → `CandidateSpec::Scalar` (uniform score=1.0, tiebreak by key)
3621    /// - `NumericWithin` → `CandidateSpec::NumericBucket` (scored, variable top-k)
3622    /// - `KeyMatch` → `CandidateSpec::ByKey` (FK probe, at most 1 dst per src)
3623    /// - `VectorSimilar` / `approximate=false` → `CandidateSpec::ScanAll` (scored)
3624    #[test]
3625    fn streaming_topk_order_identity_property_test() {
3626        // Reference: build index, compute_desired per src, then brute-force
3627        // sort (score DESC, dst-key ASC, take k) — independent of filter_src_top_k.
3628        fn reference_topk(rule: &RuleDef, k: u64, fx: &mut Fx) -> BTreeSet<(u32, u32)> {
3629            let mut idx = RuleIndex::default();
3630            for id in 0..fx.ids.len() as u32 {
3631                let label_sym = match fx.labels.get(id as usize).copied() {
3632                    Some(s) if s != u32::MAX => s,
3633                    _ => continue,
3634                };
3635                index_node_for_rule(id, label_sym, rule, &mut idx, &fx.syms, &fx.props);
3636            }
3637            let src_sym = fx.syms.get(&rule.src_label);
3638            let mut out = BTreeSet::new();
3639            let ids_snap: Vec<u32> = (0..fx.ids.len() as u32).collect();
3640            for id in ids_snap {
3641                let label_sym = match fx.labels.get(id as usize).copied() {
3642                    Some(s) if s != u32::MAX => s,
3643                    _ => continue,
3644                };
3645                if src_sym != Some(label_sym) {
3646                    continue;
3647                }
3648                let g = GraphMut {
3649                    ids: &fx.ids,
3650                    syms: &mut fx.syms,
3651                    labels: &fx.labels,
3652                    props: &fx.props,
3653                    topo: &mut fx.topo,
3654                    edge_props: &mut fx.eprops,
3655                };
3656                let per_src = compute_desired(rule, &idx, id, true, &g);
3657                // Independent brute-force sort: score DESC, dst-key ASC, take k.
3658                let mut candidates: Vec<((u32, u32), f64)> = per_src.into_iter().collect();
3659                candidates.sort_by(|&((_, da), sa), &((_, db), sb)| {
3660                    sb.total_cmp(&sa).then_with(|| {
3661                        let ka = fx.ids.key_of(da).unwrap_or("");
3662                        let kb = fx.ids.key_of(db).unwrap_or("");
3663                        ka.cmp(kb)
3664                    })
3665                });
3666                candidates.truncate(k as usize);
3667                out.extend(candidates.into_iter().map(|(k, _)| k));
3668            }
3669            out
3670        }
3671
3672        // Helper: run create_rule and return provenance (src,dst) pairs.
3673        fn streaming_pairs(rule: RuleDef, fx: &mut Fx) -> BTreeSet<(u32, u32)> {
3674            let name = rule.name.clone();
3675            let mut eng = RuleEngine::new();
3676            eng.create_rule(rule, &mut fx.g()).unwrap();
3677            eng.provenance()
3678                .get(&name)
3679                .map(|s| s.iter().map(|&(_, a, b)| (a, b)).collect())
3680                .unwrap_or_default()
3681        }
3682
3683        // ----------------------------------------------------------------
3684        // Case 1: FieldEqual (uniform score=1.0, tiebreak by key ASC)
3685        // N→N, 3-value "k" field; top-k filters per src by key.
3686        // ----------------------------------------------------------------
3687        for seed in [0u64, 1, 42, 0xDEAD_BEEF, 0x1234_5678, 99, 12_648_430, 7] {
3688            for k in [1u64, 2, 3, 5] {
3689                let rule = RuleDef {
3690                    name: "eq".into(),
3691                    src_label: "N".into(),
3692                    dst_label: "N".into(),
3693                    predicate: Predicate::FieldEqual { field: "k".into() },
3694                    edge_type: "EQ".into(),
3695                    weight_prop: None,
3696                    max_edges: Some(k),
3697                    approximate: false,
3698                    via_label: None,
3699                    via_edge: None,
3700                    via_dir: None,
3701                };
3702
3703                let build = || {
3704                    let mut fx = Fx::new();
3705                    for i in 0..12u32 {
3706                        let h = mix64(seed ^ (i as u64 + 1));
3707                        let val = match h % 3 {
3708                            0 => "a",
3709                            1 => "b",
3710                            _ => "c",
3711                        };
3712                        fx.add(
3713                            "N",
3714                            &format!("n{i:02}"),
3715                            vec![("k", Value::Str(val.into()))],
3716                        );
3717                    }
3718                    fx
3719                };
3720
3721                let expected = reference_topk(&rule, k, &mut build());
3722                let actual = streaming_pairs(rule, &mut build());
3723
3724                assert_eq!(
3725                    expected, actual,
3726                    "FieldEqual seed={seed} k={k}: streaming top-k must match brute-force top-k"
3727                );
3728            }
3729        }
3730
3731        // ----------------------------------------------------------------
3732        // Case 2: NumericWithin (scored — top-k filters by score DESC, key ASC)
3733        // S→D, numeric field "v", tolerance 10.0.
3734        // ----------------------------------------------------------------
3735        for seed in [0u64, 1, 42, 7] {
3736            for k in [1u64, 2, 4] {
3737                let rule = RuleDef {
3738                    name: "nw".into(),
3739                    src_label: "S".into(),
3740                    dst_label: "D".into(),
3741                    predicate: Predicate::NumericWithin {
3742                        field: "v".into(),
3743                        tolerance: 10.0,
3744                    },
3745                    edge_type: "NEAR".into(),
3746                    weight_prop: Some("score".into()),
3747                    max_edges: Some(k),
3748                    approximate: false,
3749                    via_label: None,
3750                    via_edge: None,
3751                    via_dir: None,
3752                };
3753
3754                let build = || {
3755                    let mut fx = Fx::new();
3756                    for i in 0..6u32 {
3757                        let h = mix64(seed ^ (i as u64 + 1));
3758                        let v = (h % 20) as f64;
3759                        fx.add("S", &format!("s{i}"), vec![("v", Value::Float(v))]);
3760                    }
3761                    for i in 0..8u32 {
3762                        let h = mix64(seed ^ (i as u64 + 101));
3763                        let v = (h % 20) as f64;
3764                        fx.add("D", &format!("d{i}"), vec![("v", Value::Float(v))]);
3765                    }
3766                    fx
3767                };
3768
3769                let expected = reference_topk(&rule, k, &mut build());
3770                let actual = streaming_pairs(rule, &mut build());
3771
3772                assert_eq!(
3773                    expected, actual,
3774                    "NumericWithin seed={seed} k={k}: streaming top-k must match brute-force top-k"
3775                );
3776            }
3777        }
3778
3779        // ----------------------------------------------------------------
3780        // Case 3: KeyMatch (CandidateSpec::ByKey)
3781        // T→C FK rule: each T has a "cid" field whose value is the key of
3782        // a C node.  Each src has at most 1 candidate, so filter_src_top_k
3783        // is the identity — but the ByKey candidate path must be exercised.
3784        // ----------------------------------------------------------------
3785        for seed in [0u64, 1, 42, 7] {
3786            for k in [1u64, 2] {
3787                let rule = RuleDef {
3788                    name: "fk".into(),
3789                    src_label: "T".into(),
3790                    dst_label: "C".into(),
3791                    predicate: Predicate::KeyMatch {
3792                        field: "cid".into(),
3793                    },
3794                    edge_type: "AT".into(),
3795                    weight_prop: None,
3796                    max_edges: Some(k),
3797                    approximate: false,
3798                    via_label: None,
3799                    via_edge: None,
3800                    via_dir: None,
3801                };
3802
3803                let build = || {
3804                    let mut fx = Fx::new();
3805                    // 4 C nodes.
3806                    for i in 0..4u32 {
3807                        fx.add("C", &format!("c{i}"), vec![]);
3808                    }
3809                    // 8 T nodes, each pointing at a C node determined by hash.
3810                    for i in 0..8u32 {
3811                        let h = mix64(seed ^ (i as u64 + 1));
3812                        let cid = format!("c{}", h % 4);
3813                        fx.add("T", &format!("t{i}"), vec![("cid", Value::Str(cid))]);
3814                    }
3815                    fx
3816                };
3817
3818                let expected = reference_topk(&rule, k, &mut build());
3819                let actual = streaming_pairs(rule, &mut build());
3820
3821                assert_eq!(
3822                    expected, actual,
3823                    "KeyMatch seed={seed} k={k}: streaming top-k must match brute-force top-k"
3824                );
3825            }
3826        }
3827
3828        // ----------------------------------------------------------------
3829        // Case 4: VectorSimilar approximate=false (CandidateSpec::ScanAll)
3830        // V→V cosine-sim rule.  6 nodes in 2 clusters of 3; min=0.9 so only
3831        // within-cluster pairs qualify.  top-k=2 filters the 2 best in cluster.
3832        // ----------------------------------------------------------------
3833        {
3834            // cluster A: unit vectors near [1,0]; cluster B: near [0,1].
3835            let cluster_a: &[(&str, f64, f64)] = &[
3836                ("va0", 1.0_f64, 0.0_f64),
3837                ("va1", 0.98_f64, 0.199_f64), // cos(~11.5°) ≈ 0.98
3838                ("va2", 0.97_f64, 0.243_f64), // cos(~14°) ≈ 0.97
3839            ];
3840            let cluster_b: &[(&str, f64, f64)] = &[
3841                ("vb0", 0.0_f64, 1.0_f64),
3842                ("vb1", 0.1_f64, 0.995_f64),
3843                ("vb2", 0.05_f64, 0.999_f64),
3844            ];
3845            for k in [1u64, 2] {
3846                let rule = RuleDef {
3847                    name: "vsim".into(),
3848                    src_label: "V".into(),
3849                    dst_label: "V".into(),
3850                    predicate: Predicate::VectorSimilar {
3851                        field: "emb".into(),
3852                        min: 0.9,
3853                    },
3854                    edge_type: "VSIM".into(),
3855                    weight_prop: Some("score".into()),
3856                    max_edges: Some(k),
3857                    approximate: false,
3858                    via_label: None,
3859                    via_edge: None,
3860                    via_dir: None,
3861                };
3862
3863                let build = || {
3864                    let mut fx = Fx::new();
3865                    let mut add_v = |key: &str, x: f64, y: f64| {
3866                        let norm = (x * x + y * y).sqrt();
3867                        let v = Value::List(vec![Value::Float(x / norm), Value::Float(y / norm)]);
3868                        fx.add("V", key, vec![("emb", v)]);
3869                    };
3870                    for &(k, x, y) in cluster_a.iter().chain(cluster_b.iter()) {
3871                        add_v(k, x, y);
3872                    }
3873                    fx
3874                };
3875
3876                let expected = reference_topk(&rule, k, &mut build());
3877                let actual = streaming_pairs(rule, &mut build());
3878
3879                assert_eq!(
3880                    expected, actual,
3881                    "VectorSimilar/ScanAll k={k}: streaming top-k must match brute-force top-k"
3882                );
3883            }
3884        }
3885    }
3886
3887    /// Streaming peak-transient allocation bound.
3888    ///
3889    /// Measures the PEAK process RSS *during* `create_rule` by polling from a
3890    /// background sampler thread at ~1 ms intervals.  Unlike a before/after
3891    /// snapshot this captures transient allocations freed before the call
3892    /// returns.
3893    ///
3894    /// **Why the OLD code would fail this test:**
3895    /// The old `compute_full_desired` built a global `BTreeMap<(u32,u32),f64>`
3896    /// for ALL 250 000 desired pairs (500 Talent × 500 Company, same field
3897    /// value, FieldEqual) before applying the cap.  At ~26 bytes per BTree
3898    /// entry (amortised node overhead on aarch64) that is ≈6.5 MiB transient
3899    /// — held for the entire duration of `apply_desired`.  The peak sampler
3900    /// would observe this spike; the 3 MiB threshold would be exceeded.
3901    ///
3902    /// **Why the NEW code passes:**
3903    /// `apply_streaming_create` caps after ~1 000 evaluations (one pass over
3904    /// the first few src nodes).  The largest in-flight allocation is one
3905    /// per-src `BTreeMap` of ≤ 500 entries ≈ 13 KiB — never materialising
3906    /// the full 250 000-pair map.  Peak transient delta is sub-100 KiB.
3907    ///
3908    /// Threshold 3 MiB: old ≈ 6.5 MiB (FAILS); new ≈ 13 KiB (PASSES).
3909    ///
3910    /// Marked `#[ignore]` (forks `ps`, environment-dependent).
3911    /// Run: `cargo test -p core-rules streaming_peak_transient_bound -- --ignored --test-threads=1`
3912    #[test]
3913    #[ignore]
3914    fn streaming_peak_transient_bound() {
3915        use std::sync::{
3916            atomic::{AtomicBool, AtomicU64, Ordering},
3917            Arc,
3918        };
3919
3920        // Sample process RSS every ~1 ms from a background thread.
3921        // Returns the peak RSS observed while `f` executes.
3922        fn peak_rss_during<F: FnOnce()>(f: F) -> u64 {
3923            let done = Arc::new(AtomicBool::new(false));
3924            let peak = Arc::new(AtomicU64::new(0));
3925            let done2 = done.clone();
3926            let peak2 = peak.clone();
3927            let pid = std::process::id().to_string();
3928
3929            let handle = std::thread::spawn(move || {
3930                while !done2.load(Ordering::Relaxed) {
3931                    let rss = std::process::Command::new("ps")
3932                        .args(["-o", "rss=", "-p", &pid])
3933                        .output()
3934                        .ok()
3935                        .and_then(|o| String::from_utf8(o.stdout).ok())
3936                        .and_then(|s| s.trim().parse::<u64>().ok())
3937                        .unwrap_or(0)
3938                        * 1024;
3939                    peak2.fetch_max(rss, Ordering::Relaxed);
3940                    std::thread::sleep(std::time::Duration::from_millis(1));
3941                }
3942            });
3943
3944            f();
3945
3946            done.store(true, Ordering::Relaxed);
3947            let _ = handle.join();
3948            peak.load(Ordering::Relaxed)
3949        }
3950
3951        // 500 Talent × 500 Company, all FieldEqual on k="same"
3952        // → 250 000 desired pairs, top-k = 2 per source (max_edges: Some(2)).
3953        // Peak transient: one per-src BTreeMap of ≤ 500 entries ≈ 13 KiB.
3954        let mut fx = Fx::new();
3955        for i in 0..500u32 {
3956            fx.add(
3957                "Talent",
3958                &format!("t{i}"),
3959                vec![("k", Value::Str("same".into()))],
3960            );
3961        }
3962        for i in 0..500u32 {
3963            fx.add(
3964                "Company",
3965                &format!("c{i}"),
3966                vec![("k", Value::Str("same".into()))],
3967            );
3968        }
3969        let rule = RuleDef {
3970            name: "eq_tc".into(),
3971            src_label: "Talent".into(),
3972            dst_label: "Company".into(),
3973            predicate: Predicate::FieldEqual { field: "k".into() },
3974            edge_type: "EQ".into(),
3975            weight_prop: None,
3976            max_edges: Some(2), // top-k=2 per source; 500 * 2 = 1000 total edges
3977            approximate: false,
3978            via_label: None,
3979            via_edge: None,
3980            via_dir: None,
3981        };
3982
3983        // Baseline: RSS before any create_rule allocation.
3984        let pid = std::process::id().to_string();
3985        let baseline = std::process::Command::new("ps")
3986            .args(["-o", "rss=", "-p", &pid])
3987            .output()
3988            .ok()
3989            .and_then(|o| String::from_utf8(o.stdout).ok())
3990            .and_then(|s| s.trim().parse::<u64>().ok())
3991            .unwrap_or(0)
3992            * 1024;
3993
3994        let mut eng = RuleEngine::new();
3995        let peak = peak_rss_during(|| {
3996            eng.create_rule(rule, &mut fx.g()).unwrap();
3997        });
3998
3999        let peak_delta = peak.saturating_sub(baseline);
4000
4001        // Threshold 3 MiB.  Old O(pairs) path: 250k entries × ~26 bytes ≈ 6.5 MiB
4002        // transient; would exceed threshold.  New streaming path: single per-src
4003        // BTreeMap ≤ 500 entries ≈ 13 KiB; never approaches threshold.
4004        assert!(
4005            peak_delta < 3 * 1024 * 1024,
4006            "peak transient delta {} bytes ({} KiB) exceeded 3 MiB; \
4007             streaming path may be building the full pairs map",
4008            peak_delta,
4009            peak_delta / 1024
4010        );
4011        assert_eq!(eng.provenance()["eq_tc"].len(), 1_000); // 500 Talent × top-k 2 = 1000
4012        assert!(!eng.is_tripped("eq_tc")); // top-k rules never trip
4013        eprintln!(
4014            "streaming_peak_transient_bound: baseline={baseline} peak={peak} \
4015             delta={peak_delta} bytes ({} KiB)",
4016            peak_delta / 1024
4017        );
4018    }
4019
4020    // -----------------------------------------------------------------------
4021    // Task 3 (Plan 11): Checkpointed Cauchy-Schwarz suffix-norm early exit
4022    // -----------------------------------------------------------------------
4023
4024    /// Helper: a near-threshold vector pair. Returns (a, b) where cos(a,b) is
4025    /// just above the provided threshold (so the pair SHOULD match).
4026    fn near_threshold_pair(dim: usize, min: f64) -> (Vec<f64>, Vec<f64>) {
4027        // Construct b = cos_target * a + epsilon * perp, then normalise both.
4028        // For simplicity: a = [1, 0, ..., 0], b = [cos_target, sin_small, 0, ...]
4029        let cos_target = min + 1e-6; // just above min
4030        let sin_small = (1.0 - cos_target * cos_target).sqrt();
4031        let mut a = vec![0.0f64; dim];
4032        a[0] = 1.0;
4033        let mut b = vec![0.0f64; dim];
4034        b[0] = cos_target;
4035        if dim > 1 {
4036            b[1] = sin_small;
4037        }
4038        (a, b)
4039    }
4040
4041    fn emb_val2(xs: &[f64]) -> Value {
4042        Value::List(xs.iter().copied().map(Value::Float).collect())
4043    }
4044
4045    /// Build an identical test fixture twice so ON/OFF/oracle comparisons all
4046    /// operate on the same graph topology.  Uses dims [2,4,8,16] with a
4047    /// near-threshold pair at dim=8 to exercise the checkpoint boundaries.
4048    fn make_early_exit_fixture(seed: u64, min: f64) -> (Fx, Vec<u32>, usize, usize) {
4049        let dims = [2usize, 4, 8, 16];
4050        let n = 100u32;
4051        let mut fx = Fx::new();
4052        let mut ids = Vec::new();
4053        for i in 0..n {
4054            let dim = dims[(i as usize) % dims.len()];
4055            let emb = rand_emb(seed, i, dim);
4056            ids.push(fx.add("Doc", &format!("d{i}"), vec![("emb", emb)]));
4057        }
4058        // Near-threshold pair at dim=8, cos just above min → must match.
4059        let (va, vb) = near_threshold_pair(8, min);
4060        let nt_a = fx.add("Doc", "nt_a", vec![("emb", emb_val2(&va))]);
4061        let nt_b = fx.add("Doc", "nt_b", vec![("emb", emb_val2(&vb))]);
4062        ids.push(nt_a);
4063        ids.push(nt_b);
4064        (fx, ids, nt_a as usize, nt_b as usize)
4065    }
4066
4067    /// Identity proof: derived edges are identical with early-exit ON, OFF,
4068    /// and vs the brute-force oracle.  Tests mixed dims (2, 4, 8, 16) with
4069    /// near-threshold cosines (cos ≈ min ± epsilon) to exercise exact rejects.
4070    #[test]
4071    fn vector_early_exit_identity_proof() {
4072        const SEED: u64 = 0xEA_4E_5A;
4073        const MIN: f64 = 0.85;
4074
4075        let def = RuleDef {
4076            name: "vec".into(),
4077            src_label: "Doc".into(),
4078            dst_label: "Doc".into(),
4079            predicate: Predicate::VectorSimilar {
4080                field: "emb".into(),
4081                min: MIN,
4082            },
4083            edge_type: "SIM".into(),
4084            weight_prop: Some("score".into()),
4085            max_edges: None,
4086            approximate: false,
4087            via_label: None,
4088            via_edge: None,
4089            via_dir: None,
4090        };
4091
4092        // Build three identical fixtures (independent topo state, same data).
4093        let (mut fx_on, ids, nt_a, nt_b) = make_early_exit_fixture(SEED, MIN);
4094        let (mut fx_off, _, _, _) = make_early_exit_fixture(SEED, MIN);
4095        let (fx_oracle, _, _, _) = make_early_exit_fixture(SEED, MIN);
4096
4097        let nt_a = nt_a as u32;
4098        let nt_b = nt_b as u32;
4099
4100        // Run with early-exit ON (default).
4101        let mut eng_on = RuleEngine::new();
4102        {
4103            let mut g = fx_on.g();
4104            eng_on.create_rule(def.clone(), &mut g).unwrap();
4105        }
4106        let edges_on = prov_pairs(&eng_on, "vec");
4107        assert!(!edges_on.is_empty(), "should produce some edges");
4108
4109        // Near-threshold pair must appear with early-exit ON.
4110        assert!(
4111            edges_on.contains(&(nt_a, nt_b)),
4112            "near-threshold pair nt_a→nt_b must match with early-exit ON"
4113        );
4114        assert!(
4115            edges_on.contains(&(nt_b, nt_a)),
4116            "near-threshold pair nt_b→nt_a must match with early-exit ON"
4117        );
4118
4119        // Run with early-exit OFF; must produce identical edge set.
4120        let mut eng_off = RuleEngine::new();
4121        {
4122            let mut g = fx_off.g();
4123            with_vector_early_exit(false, || {
4124                eng_off.create_rule(def.clone(), &mut g).unwrap();
4125            });
4126        }
4127        let edges_off = prov_pairs(&eng_off, "vec");
4128        assert_eq!(
4129            edges_on, edges_off,
4130            "early-exit ON vs OFF must produce identical edges"
4131        );
4132
4133        // Brute-force oracle: evaluate() on all (s,d) pairs.
4134        let mut oracle = BTreeSet::new();
4135        for &s in &ids {
4136            for &d in &ids {
4137                if s == d {
4138                    continue;
4139                }
4140                let skey = fx_oracle.ids.key_of(s).unwrap();
4141                let dkey = fx_oracle.ids.key_of(d).unwrap();
4142                let sg = |f: &str| fx_oracle.props.get(s, f).cloned();
4143                let dg = |f: &str| fx_oracle.props.get(d, f).cloned();
4144                if evaluate(
4145                    &def.predicate,
4146                    &NodeView {
4147                        key: skey,
4148                        props: &sg,
4149                    },
4150                    &NodeView {
4151                        key: dkey,
4152                        props: &dg,
4153                    },
4154                )
4155                .is_some()
4156                {
4157                    oracle.insert((s, d));
4158                }
4159            }
4160        }
4161        assert_eq!(
4162            edges_on, oracle,
4163            "early-exit ON vs brute-force oracle must be identical"
4164        );
4165    }
4166
4167    /// Coherence: checkpoints are rebuilt through the insert/remove choke-points
4168    /// when a vector prop is updated.  Dim change, freshness gate exercised.
4169    #[test]
4170    fn vector_early_exit_checkpoint_coherence() {
4171        let mut fx = Fx::new();
4172        // Two dim=4 nodes that match under VectorSimilar min=0.9.
4173        let a = fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0]))]);
4174        let b = fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0]))]);
4175        // dim=6 node that should NOT match dim=4 nodes.
4176        let c = fx.add(
4177            "Doc",
4178            "c",
4179            vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0, 0.0, 0.0]))],
4180        );
4181        let def = RuleDef {
4182            name: "vec".into(),
4183            src_label: "Doc".into(),
4184            dst_label: "Doc".into(),
4185            predicate: Predicate::VectorSimilar {
4186                field: "emb".into(),
4187                min: 0.9,
4188            },
4189            edge_type: "SIM".into(),
4190            weight_prop: None,
4191            max_edges: None,
4192            approximate: false,
4193            via_label: None,
4194            via_edge: None,
4195            via_dir: None,
4196        };
4197
4198        let mut eng = RuleEngine::new();
4199        {
4200            let mut g = fx.g();
4201            eng.create_rule(def.clone(), &mut g).unwrap();
4202        }
4203
4204        // Checkpoints must be populated for all three nodes.
4205        assert!(
4206            eng.indexes["vec"].src_side.vec_ckpts(a).is_some(),
4207            "a must have src checkpoints"
4208        );
4209        assert!(
4210            eng.indexes["vec"].dst_side.vec_ckpts(b).is_some(),
4211            "b must have dst checkpoints"
4212        );
4213        assert!(
4214            eng.indexes["vec"].src_side.vec_ckpts(c).is_some(),
4215            "c must have src checkpoints (dim=6)"
4216        );
4217
4218        // ckpts[0] must equal the full L2 norm.
4219        let ckpts_a = *eng.indexes["vec"].src_side.vec_ckpts(a).unwrap();
4220        let norm_a = eng.indexes["vec"].src_side.vec_meta(a).unwrap().1;
4221        assert!(
4222            (ckpts_a[0] - norm_a).abs() < 1e-12,
4223            "ckpts[0] must equal the full L2 norm"
4224        );
4225
4226        // Initial edges: a↔b only (c is different dim).
4227        assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(a, b), (b, a)]));
4228
4229        // Update b to dim=6 (same as c) — choke-points must rebuild checkpoints.
4230        let old_b = fx.props.get(b, "emb").cloned();
4231        fx.props
4232            .set(b, "emb", emb_val(&[1.0, 0.0, 0.0, 0.0, 0.0, 0.0]));
4233        {
4234            let mut g = fx.g();
4235            eng.on_node_changed(b, Some(("emb", old_b)), &mut g);
4236        }
4237        // b's dim must now be 6 in both sides.
4238        assert_eq!(eng.indexes["vec"].src_side.vec_dim(b), Some(6));
4239        assert_eq!(eng.indexes["vec"].dst_side.vec_dim(b), Some(6));
4240        // b must have new checkpoints for dim=6.
4241        assert!(eng.indexes["vec"].src_side.vec_ckpts(b).is_some());
4242        // Edges must now be b↔c (both dim=6, cos=1.0 > 0.9).
4243        assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(b, c), (c, b)]));
4244
4245        // Freshness gate: fresh_ckpts_for returns None when live vector differs.
4246        // Simulate by passing a different live vector to fresh_ckpts_for.
4247        let wrong_live = vec![2.0f64, 0.0, 0.0, 0.0, 0.0, 0.0]; // same dim, different norm
4248        let gate_result = eng.indexes["vec"].src_side.fresh_ckpts_for(b, &wrong_live);
4249        assert!(
4250            gate_result.is_none(),
4251            "freshness gate must reject a mismatched-norm live vector"
4252        );
4253
4254        // fresh_ckpts_for must succeed with the correct live vector.
4255        let correct_live = vec![1.0f64, 0.0, 0.0, 0.0, 0.0, 0.0];
4256        let gate_result = eng.indexes["vec"]
4257            .src_side
4258            .fresh_ckpts_for(b, &correct_live);
4259        assert!(
4260            gate_result.is_some(),
4261            "freshness gate must accept the matching live vector"
4262        );
4263    }
4264
4265    /// Razor test: dim=1536 pair with true cosine within 1e-12 of `min`.
4266    ///
4267    /// Purpose: with energy spread uniformly across all 1536 elements, each
4268    /// checkpoint boundary contributes a tiny slice of dot product.  Float
4269    /// rounding of suffix-norm accumulation can shift `cos_max` by O(dim × ε)
4270    /// ≈ 3.4 × 10⁻¹³ at dim=1536, inside the 1e-12 margin tested here.  The
4271    /// epsilon guard in `cosine_early_exit` absorbs this; ON/OFF/oracle must
4272    /// agree on all edges.
4273    #[test]
4274    fn vector_early_exit_razor_dim1536() {
4275        const MIN: f64 = 0.85;
4276        const DIM: usize = 1536;
4277        // target cosine = min + 5e-13: inside the dim-scale float-error zone.
4278        let target = MIN + 5e-13;
4279        let inv_sqrt = 1.0 / (DIM as f64).sqrt();
4280
4281        // a: unit-norm uniform vector — energy spread equally across all chunks.
4282        let a: Vec<f64> = vec![inv_sqrt; DIM];
4283
4284        // b = target * a + sqrt(1 - target^2) * e_perp
4285        // e_perp = [1, -1, 0, ..., 0] / sqrt(2) is perpendicular to uniform a:
4286        //   dot(a, e_perp) = inv_sqrt * (1 - 1) / sqrt(2) = 0  ✓
4287        // norm(b) = sqrt(target^2 + (1-target^2)) = 1            ✓
4288        // cos(a, b) = dot(a, b) = target * dot(a, a) = target    ✓
4289        let perp_scale = (1.0 - target * target).sqrt() / (2.0f64).sqrt();
4290        let mut b: Vec<f64> = vec![target * inv_sqrt; DIM];
4291        b[0] += perp_scale;
4292        b[1] -= perp_scale;
4293
4294        let def = RuleDef {
4295            name: "razor".into(),
4296            src_label: "Doc".into(),
4297            dst_label: "Doc".into(),
4298            predicate: Predicate::VectorSimilar {
4299                field: "emb".into(),
4300                min: MIN,
4301            },
4302            edge_type: "SIM".into(),
4303            weight_prop: None,
4304            max_edges: None,
4305            approximate: false,
4306            via_label: None,
4307            via_edge: None,
4308            via_dir: None,
4309        };
4310
4311        // Three independent fixtures with the same razor pair.
4312        let build_fx = || {
4313            let mut fx = Fx::new();
4314            let na = fx.add("Doc", "razor_a", vec![("emb", emb_val2(&a))]);
4315            let nb = fx.add("Doc", "razor_b", vec![("emb", emb_val2(&b))]);
4316            (fx, na, nb)
4317        };
4318
4319        let (mut fx_on, na, nb) = build_fx();
4320        let (mut fx_off, _, _) = build_fx();
4321        let (fx_oracle, _, _) = build_fx();
4322
4323        // ON
4324        let mut eng_on = RuleEngine::new();
4325        {
4326            let mut g = fx_on.g();
4327            eng_on.create_rule(def.clone(), &mut g).unwrap();
4328        }
4329        let edges_on = prov_pairs(&eng_on, "razor");
4330        assert!(
4331            edges_on.contains(&(na, nb)),
4332            "razor pair razor_a→razor_b must be present with early-exit ON (cos={target:.15}, min={MIN})"
4333        );
4334        assert!(
4335            edges_on.contains(&(nb, na)),
4336            "razor pair razor_b→razor_a must be present with early-exit ON"
4337        );
4338
4339        // OFF
4340        let mut eng_off = RuleEngine::new();
4341        {
4342            let mut g = fx_off.g();
4343            with_vector_early_exit(false, || {
4344                eng_off.create_rule(def.clone(), &mut g).unwrap();
4345            });
4346        }
4347        let edges_off = prov_pairs(&eng_off, "razor");
4348        assert_eq!(
4349            edges_on, edges_off,
4350            "razor dim=1536: early-exit ON vs OFF must produce identical edges"
4351        );
4352
4353        // Brute-force oracle.
4354        let ids = [na, nb];
4355        let mut oracle = BTreeSet::new();
4356        for &s in &ids {
4357            for &d in &ids {
4358                if s == d {
4359                    continue;
4360                }
4361                let skey = fx_oracle.ids.key_of(s).unwrap();
4362                let dkey = fx_oracle.ids.key_of(d).unwrap();
4363                let sg = |f: &str| fx_oracle.props.get(s, f).cloned();
4364                let dg = |f: &str| fx_oracle.props.get(d, f).cloned();
4365                if evaluate(
4366                    &def.predicate,
4367                    &NodeView {
4368                        key: skey,
4369                        props: &sg,
4370                    },
4371                    &NodeView {
4372                        key: dkey,
4373                        props: &dg,
4374                    },
4375                )
4376                .is_some()
4377                {
4378                    oracle.insert((s, d));
4379                }
4380            }
4381        }
4382        assert_eq!(
4383            edges_on, oracle,
4384            "razor dim=1536: early-exit ON vs brute-force oracle must be identical"
4385        );
4386    }
4387}