Skip to main content

core_rules/
engine.rs

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