1use crate::def::{evaluate, is_keymatch_rooted, NodeView, Predicate, RuleDef};
2use crate::index::{
3 candidate_spec, candidate_spec_approx_with_k, ivf_drift_rebuild_threshold, CandidateSpec,
4 RuleIndex,
5};
6use core_storage::{ColumnStore, EdgeProps, IdMap, Interner, Topology, Value};
7use std::collections::{BTreeMap, BTreeSet};
8
9#[derive(Debug, Clone)]
16pub struct EngineEdgeDelta {
17 pub rule: String,
18 pub src_key: String,
20 pub dst_key: String,
22 pub edge_type: String,
24 pub etype_sym: u32,
26 pub src_id: u32,
28 pub dst_id: u32,
30 pub fired: bool,
32}
33
34#[cfg(test)]
35pub use crate::index::{with_ivf_drift_rebuild, with_vector_dim_reject, with_vector_early_exit};
36
37pub struct GraphMut<'a> {
39 pub ids: &'a IdMap,
40 pub syms: &'a mut Interner,
41 pub labels: &'a [u32],
42 pub props: &'a ColumnStore,
43 pub topo: &'a mut Topology,
44 pub edge_props: &'a mut EdgeProps,
45}
46
47pub const DEFAULT_MAX_EDGES: u64 = 1_000_000;
49
50type Triple = (u32, u32, u32);
52type Touch = (u32, u32, u32, u32);
54
55pub type SideIvfExport = (Vec<Vec<f64>>, BTreeMap<u32, usize>, u64);
58pub type RuleIvfExport = (SideIvfExport, SideIvfExport);
60
61#[derive(Debug, Default)]
62pub struct RuleEngine {
63 rules: BTreeMap<String, RuleDef>,
64 indexes: BTreeMap<String, RuleIndex>,
65 provenance: BTreeMap<String, BTreeSet<Triple>>,
66 owned: BTreeSet<Triple>,
67 by_node: BTreeMap<u32, BTreeSet<Touch>>,
70 rule_intern: BTreeMap<String, u32>,
75 intern_rule: Vec<String>,
76 tripped: BTreeMap<String, bool>,
77 fires: BTreeMap<String, u64>,
78 pending_deltas: Vec<EngineEdgeDelta>,
85 emit_deltas: bool,
98 rebuild_needed: BTreeSet<String>,
102}
103
104fn candidate_spec_for(def: &RuleDef) -> CandidateSpec<'_> {
111 if def.approximate {
112 let k = def.max_edges.map(|me| me.max(64)).unwrap_or(64) as usize;
115 candidate_spec_approx_with_k(&def.predicate, k)
116 } else {
117 candidate_spec(&def.predicate)
118 }
119}
120
121fn src_lookup_spec_for(def: &RuleDef) -> CandidateSpec<'_> {
127 if is_keymatch_rooted(&def.predicate) {
128 let field =
129 keymatch_field(&def.predicate).expect("keymatch-rooted predicate has a KeyMatch field");
130 CandidateSpec::Scalar { field }
131 } else {
132 candidate_spec_for(def)
133 }
134}
135
136fn predicate_covers_field(p: &Predicate, field: &str) -> bool {
138 match p {
139 Predicate::VectorSimilar { field: f, .. } => f == field,
140 Predicate::All(parts) | Predicate::Any(parts) => {
141 parts.iter().any(|q| predicate_covers_field(q, field))
142 }
143 _ => false,
144 }
145}
146
147fn keymatch_field(p: &Predicate) -> Option<&str> {
149 match p {
150 Predicate::KeyMatch { field } => Some(field),
151 Predicate::All(parts) => parts.first().and_then(keymatch_field),
152 Predicate::Any(_) => None,
153 _ => None,
154 }
155}
156
157fn compute_desired(
160 def: &RuleDef,
161 index: &RuleIndex,
162 n: u32,
163 on_src_side: bool,
164 g: &GraphMut<'_>,
165) -> BTreeMap<(u32, u32), f64> {
166 let (my_label, other_label) = if on_src_side {
167 (&def.src_label, &def.dst_label)
168 } else {
169 (&def.dst_label, &def.src_label)
170 };
171
172 let Some(my_sym) = g.syms.get(my_label) else {
173 return BTreeMap::new();
174 };
175 if g.labels.get(n as usize).copied() != Some(my_sym) {
176 return BTreeMap::new();
177 }
178 let other_sym = g.syms.get(other_label);
179
180 let n_key = match g.ids.key_of(n) {
181 Some(k) => k,
182 None => return BTreeMap::new(),
183 };
184 let n_get = |f: &str| g.props.get(n, f).cloned();
185
186 let spec = candidate_spec_for(def);
187 let candidates: BTreeSet<u32> = if on_src_side {
188 if is_keymatch_rooted(&def.predicate) {
189 let field = keymatch_field(&def.predicate).expect("ByKey always comes from KeyMatch");
193 match n_get(field) {
194 Some(Value::Str(ref target_key)) => match g.ids.get(target_key) {
195 Some(dst_id) => std::iter::once(dst_id).collect(),
196 None => BTreeSet::new(),
197 },
198 _ => BTreeSet::new(),
199 }
200 } else {
201 index.dst_side.candidates(&spec, &n_get)
202 }
203 } else {
204 let src_spec = src_lookup_spec_for(def);
206 if is_keymatch_rooted(&def.predicate) {
207 let key_getter = |_: &str| Some(Value::Str(n_key.to_string()));
210 index.src_side.candidates(&src_spec, &key_getter)
211 } else {
212 index.src_side.candidates(&src_spec, &n_get)
213 }
214 };
215
216 let n_early_exit_hint: Option<(Vec<f64>, f64, [f64; 8])> = if !def.approximate {
231 if let Predicate::VectorSimilar { field, .. } = &def.predicate {
232 if crate::index::vector_early_exit_enabled() {
233 let n_side = if on_src_side {
234 &index.src_side
235 } else {
236 &index.dst_side
237 };
238 if let Some(vn_v) = n_get(field) {
239 if let Some(vn) = crate::index::as_numeric_list(&vn_v) {
240 if let Some((norm_n, ckpts_n)) = n_side.fresh_ckpts_for(n, &vn) {
241 Some((vn, norm_n, *ckpts_n))
242 } else {
243 None
244 }
245 } else {
246 None
247 }
248 } else {
249 None
250 }
251 } else {
252 None
253 }
254 } else {
255 None
256 }
257 } else {
258 None
259 };
260
261 let mut out = BTreeMap::new();
262 for m in candidates {
263 if m == n {
264 continue; }
266 if g.labels.get(m as usize).copied() != other_sym {
267 continue; }
269 let m_key = match g.ids.key_of(m) {
270 Some(k) => k,
271 None => continue,
272 };
273 let m_get = |f: &str| g.props.get(m, f).cloned();
274 let (s_view, d_view, s_id, d_id) = if on_src_side {
275 (
276 NodeView {
277 key: n_key,
278 props: &n_get,
279 },
280 NodeView {
281 key: m_key,
282 props: &m_get,
283 },
284 n,
285 m,
286 )
287 } else {
288 (
289 NodeView {
290 key: m_key,
291 props: &m_get,
292 },
293 NodeView {
294 key: n_key,
295 props: &n_get,
296 },
297 m,
298 n,
299 )
300 };
301
302 if let (Some((ref vn, norm_n, ckpts_n)), Predicate::VectorSimilar { field, min }) =
304 (&n_early_exit_hint, &def.predicate)
305 {
306 let m_side = if on_src_side {
307 &index.dst_side
308 } else {
309 &index.src_side
310 };
311 if let Some(vm_v) = m_get(field) {
312 if let Some(vm) = crate::index::as_numeric_list(&vm_v) {
313 if let Some((norm_m, ckpts_m)) = m_side.fresh_ckpts_for(m, &vm) {
314 let (va, ckpts_a, na, vb, ckpts_b, nb) = if on_src_side {
315 (
316 vn.as_slice(),
317 ckpts_n,
318 *norm_n,
319 vm.as_slice(),
320 ckpts_m,
321 norm_m,
322 )
323 } else {
324 (
325 vm.as_slice(),
326 ckpts_m,
327 norm_m,
328 vn.as_slice(),
329 ckpts_n,
330 *norm_n,
331 )
332 };
333 match crate::def::cosine_early_exit(va, vb, ckpts_a, ckpts_b, na, nb, *min)
334 {
335 None => continue, Some(score) => {
337 out.insert((s_id, d_id), score);
338 continue; }
340 }
341 }
342 }
343 }
344 }
345
346 if let Some(score) = evaluate(&def.predicate, &s_view, &d_view) {
347 out.insert((s_id, d_id), score);
348 }
349 }
350 out
351}
352
353fn compute_desired_via(
366 def: &RuleDef,
367 anchor: ViaAnchor,
368 g: &GraphMut<'_>,
369) -> BTreeMap<(u32, u32), f64> {
370 let via_label = def.via_label.as_deref().unwrap();
371 let via_edge_str = def.via_edge.as_deref().unwrap();
372 let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
373
374 let src_sym = match g.syms.get(&def.src_label) {
375 Some(s) => s,
376 None => return BTreeMap::new(),
377 };
378 let via_sym = match g.syms.get(via_label) {
379 Some(s) => s,
380 None => return BTreeMap::new(),
381 };
382 let dst_sym = match g.syms.get(&def.dst_label) {
383 Some(s) => s,
384 None => return BTreeMap::new(),
385 };
386 let via_etype = match g.syms.get(via_edge_str) {
387 Some(e) => e,
388 None => return BTreeMap::new(),
389 };
390
391 let srcs: Vec<u32> = match anchor {
393 ViaAnchor::Src(src_id) => {
394 if g.labels.get(src_id as usize).copied() == Some(src_sym) {
395 vec![src_id]
396 } else {
397 return BTreeMap::new();
398 }
399 }
400 ViaAnchor::Dst(_) => {
401 (0..g.ids.len() as u32)
403 .filter(|&id| {
404 matches!(
405 g.labels.get(id as usize).copied(),
406 Some(s) if s != u32::MAX && s == src_sym
407 )
408 })
409 .collect()
410 }
411 };
412
413 let anchored_dst: Option<u32> = match anchor {
415 ViaAnchor::Dst(dst_id) => {
416 if g.labels.get(dst_id as usize).copied() == Some(dst_sym) {
417 Some(dst_id)
418 } else {
419 return BTreeMap::new();
420 }
421 }
422 _ => None,
423 };
424
425 let mut out = BTreeMap::new();
426
427 for src in srcs {
428 let _src_key = match g.ids.key_of(src) {
429 Some(k) => k,
430 None => continue,
431 };
432 let via_neighbors: Vec<u32> = g
434 .topo
435 .neighbors(via_etype, via_dir, src)
436 .iter()
437 .copied()
438 .filter(|&v| g.labels.get(v as usize).copied() == Some(via_sym))
439 .collect();
440
441 if via_neighbors.is_empty() {
442 continue;
443 }
444
445 let dsts: Vec<u32> = if let Some(dst_id) = anchored_dst {
447 vec![dst_id]
448 } else {
449 (0..g.ids.len() as u32)
450 .filter(|&id| {
451 id != src
452 && matches!(
453 g.labels.get(id as usize).copied(),
454 Some(s) if s != u32::MAX && s == dst_sym
455 )
456 })
457 .collect()
458 };
459
460 for dst in dsts {
461 if dst == src {
462 continue; }
464 let dst_key = match g.ids.key_of(dst) {
465 Some(k) => k,
466 None => continue,
467 };
468 let dst_get = |f: &str| g.props.get(dst, f).cloned();
469 let dst_view = NodeView {
470 key: dst_key,
471 props: &dst_get,
472 };
473
474 let mut best: Option<f64> = None;
476 for &via_id in &via_neighbors {
477 let via_key = match g.ids.key_of(via_id) {
478 Some(k) => k,
479 None => continue,
480 };
481 let via_get = |f: &str| g.props.get(via_id, f).cloned();
482 let via_view = NodeView {
483 key: via_key,
484 props: &via_get,
485 };
486 if let Some(score) = evaluate(&def.predicate, &via_view, &dst_view) {
487 best = Some(match best {
488 None => score,
489 Some(prev) => prev.max(score),
490 });
491 }
492 }
493
494 if let Some(score) = best {
495 out.insert((src, dst), score);
496 }
497 }
498 }
499
500 out
501}
502
503enum ViaAnchor {
505 Src(u32),
507 Dst(u32),
510}
511
512fn edge_budget(def: &RuleDef) -> u64 {
513 def.max_edges.unwrap_or(DEFAULT_MAX_EDGES)
516}
517
518pub(crate) fn filter_src_top_k(
540 per_src: BTreeMap<(u32, u32), f64>,
541 k: u64,
542 ids: &core_storage::IdMap,
543) -> BTreeMap<(u32, u32), f64> {
544 if per_src.len() as u64 <= k {
545 return per_src;
546 }
547 let mut candidates: Vec<((u32, u32), f64)> = per_src.into_iter().collect();
548 candidates.sort_by(|&((_, da), sa), &((_, db), sb)| {
550 sb.total_cmp(&sa).then_with(|| {
551 let ka = ids.key_of(da).unwrap_or("");
552 let kb = ids.key_of(db).unwrap_or("");
553 ka.cmp(kb)
554 })
555 });
556 candidates.truncate(k as usize);
557 candidates.into_iter().collect()
558}
559
560fn apply_per_src_top_k(
567 def: &RuleDef,
568 src: u32,
569 desired_from_src: BTreeMap<(u32, u32), f64>,
570 prov: &mut ProvSets<'_>,
571 g: &mut GraphMut<'_>,
572) {
573 let et = g.syms.intern(&def.edge_type);
574
575 let current: Vec<Triple> = {
579 let rid = prov.rule_intern.get(&def.name).copied();
580 prov.by_node
581 .get(&src)
582 .into_iter()
583 .flatten()
584 .filter(|(r, t, s, _d)| Some(*r) == rid && *t == et && *s == src)
585 .map(|(_, t, s, d)| (*t, *s, *d))
586 .collect()
587 };
588
589 for (t, s, d) in current {
591 if !desired_from_src.contains_key(&(s, d)) {
592 g.topo.remove_edge(t, s, d);
593 g.edge_props.remove_edge(t, s, d);
594 prov.remove(&def.name, (t, s, d), g.ids, g.syms);
595 }
596 }
597
598 for ((s, d), score) in &desired_from_src {
600 let triple = (et, *s, *d);
601 let already = prov.contains(&triple);
602 if !already {
603 let newly = g.topo.add_edge(et, *s, *d);
604 if newly {
605 prov.insert(&def.name, triple, g.ids, g.syms);
606 }
607 }
608 let is_owned = already || prov.contains(&triple);
609 if is_owned {
610 if let Some(p) = &def.weight_prop {
611 g.edge_props.set(et, *s, *d, p, Value::Float(*score));
612 }
613 }
614 }
615}
616
617fn intern_rule(intern: &mut BTreeMap<String, u32>, names: &mut Vec<String>, rule: &str) -> u32 {
619 if let Some(&id) = intern.get(rule) {
620 return id;
621 }
622 let id = names.len() as u32;
623 intern.insert(rule.to_string(), id);
624 names.push(rule.to_string());
625 id
626}
627
628type ByNodeRebuild = (
629 BTreeMap<u32, BTreeSet<Touch>>,
630 BTreeMap<String, u32>,
631 Vec<String>,
632);
633
634fn rebuild_by_node(provenance: &BTreeMap<String, BTreeSet<Triple>>) -> ByNodeRebuild {
635 let mut by_node = BTreeMap::new();
636 let mut intern = BTreeMap::new();
637 let mut names = Vec::new();
638 for (rule, set) in provenance {
639 let rid = intern_rule(&mut intern, &mut names, rule);
640 for &triple in set {
641 touch_insert(&mut by_node, rid, triple);
642 }
643 }
644 (by_node, intern, names)
645}
646
647fn touch_insert(by_node: &mut BTreeMap<u32, BTreeSet<Touch>>, rid: u32, triple: Triple) {
648 let (t, s, d) = triple;
649 let entry = (rid, t, s, d);
650 by_node.entry(s).or_default().insert(entry);
651 if s != d {
652 by_node.entry(d).or_default().insert(entry);
653 }
654}
655
656fn touch_remove(by_node: &mut BTreeMap<u32, BTreeSet<Touch>>, rid: u32, triple: Triple) {
657 let (t, s, d) = triple;
658 let entry = (rid, t, s, d);
659 if let Some(set) = by_node.get_mut(&s) {
660 set.remove(&entry);
661 if set.is_empty() {
662 by_node.remove(&s);
663 }
664 }
665 if s != d {
666 if let Some(set) = by_node.get_mut(&d) {
667 set.remove(&entry);
668 if set.is_empty() {
669 by_node.remove(&d);
670 }
671 }
672 }
673}
674
675#[cfg(test)]
676fn resolve_by_node(
677 by_node: &BTreeMap<u32, BTreeSet<Touch>>,
678 names: &[String],
679) -> BTreeMap<u32, BTreeSet<(String, Triple)>> {
680 by_node
681 .iter()
682 .map(|(&n, set)| {
683 let resolved = set
684 .iter()
685 .map(|&(rid, t, s, d)| (names[rid as usize].clone(), (t, s, d)))
686 .collect();
687 (n, resolved)
688 })
689 .collect()
690}
691
692struct ProvSets<'a> {
695 set: &'a mut BTreeSet<Triple>,
696 owned: &'a mut BTreeSet<Triple>,
697 by_node: &'a mut BTreeMap<u32, BTreeSet<Touch>>,
698 rule_intern: &'a mut BTreeMap<String, u32>,
699 intern_rule: &'a mut Vec<String>,
700 deltas: &'a mut Vec<EngineEdgeDelta>,
704 emit: bool,
707}
708
709impl ProvSets<'_> {
710 fn insert(&mut self, rule: &str, triple: Triple, ids: &IdMap, syms: &Interner) -> bool {
714 if !self.set.insert(triple) {
715 return false;
716 }
717 self.owned.insert(triple);
718 let rid = intern_rule(self.rule_intern, self.intern_rule, rule);
719 touch_insert(self.by_node, rid, triple);
720 let (etype, src, dst) = triple;
721 if self.emit {
722 if let (Some(sk), Some(dk), Some(et)) =
723 (ids.key_of(src), ids.key_of(dst), syms.resolve(etype))
724 {
725 self.deltas.push(EngineEdgeDelta {
726 rule: rule.to_string(),
727 src_key: sk.to_string(),
728 dst_key: dk.to_string(),
729 edge_type: et.to_string(),
730 etype_sym: etype,
731 src_id: src,
732 dst_id: dst,
733 fired: true,
734 });
735 }
736 }
737 true
738 }
739
740 fn remove(&mut self, rule: &str, triple: Triple, ids: &IdMap, syms: &Interner) -> bool {
741 if !self.set.remove(&triple) {
742 return false;
743 }
744 self.owned.remove(&triple);
745 let rid = intern_rule(self.rule_intern, self.intern_rule, rule);
746 touch_remove(self.by_node, rid, triple);
747 let (etype, src, dst) = triple;
748 if self.emit {
749 if let (Some(sk), Some(dk), Some(et)) =
750 (ids.key_of(src), ids.key_of(dst), syms.resolve(etype))
751 {
752 self.deltas.push(EngineEdgeDelta {
753 rule: rule.to_string(),
754 src_key: sk.to_string(),
755 dst_key: dk.to_string(),
756 edge_type: et.to_string(),
757 etype_sym: etype,
758 src_id: src,
759 dst_id: dst,
760 fired: false,
761 });
762 }
763 }
764 true
765 }
766
767 fn contains(&self, triple: &Triple) -> bool {
768 self.set.contains(triple)
769 }
770
771 fn len(&self) -> usize {
772 self.set.len()
773 }
774}
775
776fn apply_desired(
787 def: &RuleDef,
788 desired: BTreeMap<(u32, u32), f64>,
789 retract_touching: Option<u32>,
790 prov: &mut ProvSets<'_>,
791 tripped: &mut bool,
792 g: &mut GraphMut<'_>,
793) {
794 let budget = edge_budget(def);
795 let et = g.syms.intern(&def.edge_type);
796
797 let current: Vec<Triple> = match retract_touching {
798 None => prov
799 .set
800 .iter()
801 .filter(|(t, _, _)| *t == et)
802 .copied()
803 .collect(),
804 Some(n) => {
805 let rid = prov.rule_intern.get(&def.name).copied();
806 prov.by_node
807 .get(&n)
808 .into_iter()
809 .flatten()
810 .filter(|(r, t, _, _)| Some(*r) == rid && *t == et)
811 .map(|(_, t, s, d)| (*t, *s, *d))
812 .collect()
813 }
814 };
815
816 for (t, s, d) in current {
817 if !desired.contains_key(&(s, d)) {
818 g.topo.remove_edge(t, s, d);
819 g.edge_props.remove_edge(t, s, d);
820 prov.remove(&def.name, (t, s, d), g.ids, g.syms);
821 }
822 }
823
824 for ((s, d), score) in desired {
825 let triple = (et, s, d);
826 let already = prov.contains(&triple);
827 if !already {
828 if *tripped || prov.len() as u64 >= budget {
829 *tripped = true;
830 continue;
831 }
832 let newly = g.topo.add_edge(et, s, d);
833 if newly {
834 prov.insert(&def.name, triple, g.ids, g.syms);
835 }
836 }
837 let is_owned_here = already || prov.contains(&triple);
841 if is_owned_here {
842 if let Some(p) = &def.weight_prop {
843 g.edge_props.set(et, s, d, p, Value::Float(score));
844 }
845 }
846 }
847}
848
849#[cfg(test)]
857#[allow(dead_code)]
858fn compute_full_desired(
859 def: &RuleDef,
860 index: &RuleIndex,
861 g: &GraphMut<'_>,
862) -> BTreeMap<(u32, u32), f64> {
863 let mut desired = BTreeMap::new();
864 let src_sym = g.syms.get(&def.src_label);
865 for id in 0..g.ids.len() as u32 {
866 let label_sym = match g.labels.get(id as usize).copied() {
867 Some(s) if s != u32::MAX => s,
868 _ => continue,
869 };
870 if src_sym == Some(label_sym) {
871 desired.extend(compute_desired(def, index, id, true, g));
872 }
873 }
874 desired
875}
876
877fn pair_still_desired(def: &RuleDef, s: u32, d: u32, g: &GraphMut<'_>) -> bool {
884 let src_sym = match g.syms.get(&def.src_label) {
885 Some(sym) => sym,
886 None => return false,
887 };
888 let dst_sym = match g.syms.get(&def.dst_label) {
889 Some(sym) => sym,
890 None => return false,
891 };
892 if g.labels.get(s as usize).copied() != Some(src_sym) {
893 return false;
894 }
895 if g.labels.get(d as usize).copied() != Some(dst_sym) {
896 return false;
897 }
898 let s_key = match g.ids.key_of(s) {
899 Some(k) => k,
900 None => return false,
901 };
902 let d_key = match g.ids.key_of(d) {
903 Some(k) => k,
904 None => return false,
905 };
906 let s_get = |f: &str| g.props.get(s, f).cloned();
907 let d_get = |f: &str| g.props.get(d, f).cloned();
908 evaluate(
909 &def.predicate,
910 &NodeView {
911 key: s_key,
912 props: &s_get,
913 },
914 &NodeView {
915 key: d_key,
916 props: &d_get,
917 },
918 )
919 .is_some()
920}
921
922fn count_desired_up_to(def: &RuleDef, index: &RuleIndex, limit: u64, g: &GraphMut<'_>) -> u64 {
927 let mut count = 0u64;
928 let src_sym = g.syms.get(&def.src_label);
929 for id in 0..g.ids.len() as u32 {
930 let label_sym = match g.labels.get(id as usize).copied() {
931 Some(s) if s != u32::MAX => s,
932 _ => continue,
933 };
934 if src_sym != Some(label_sym) {
935 continue;
936 }
937 count += compute_desired(def, index, id, true, g).len() as u64;
938 if count > limit {
939 return count;
940 }
941 }
942 count
943}
944
945fn apply_streaming_create(
967 def: &RuleDef,
968 index: &RuleIndex,
969 prov: &mut ProvSets<'_>,
970 tripped: &mut bool,
971 g: &mut GraphMut<'_>,
972) {
973 let budget = edge_budget(def);
974 let et = g.syms.intern(&def.edge_type);
975 let src_sym = g.syms.get(&def.src_label);
976
977 'outer: for id in 0..g.ids.len() as u32 {
978 let label_sym = match g.labels.get(id as usize).copied() {
979 Some(s) if s != u32::MAX => s,
980 _ => continue,
981 };
982 if src_sym != Some(label_sym) {
983 continue;
984 }
985 let per_src = compute_desired(def, index, id, true, g);
986 for ((s, d), score) in per_src {
987 let triple = (et, s, d);
988 let already = prov.contains(&triple);
993 if !already {
994 if *tripped || prov.len() as u64 >= budget {
995 *tripped = true;
996 break 'outer;
997 }
998 let newly = g.topo.add_edge(et, s, d);
999 if newly {
1000 prov.insert(&def.name, triple, g.ids, g.syms);
1001 }
1002 }
1003 let is_owned_here = already || prov.contains(&triple);
1004 if is_owned_here {
1005 if let Some(p) = &def.weight_prop {
1006 g.edge_props.set(et, s, d, p, Value::Float(score));
1007 }
1008 }
1009 }
1010 }
1011}
1012
1013fn apply_streaming_create_top_k(
1020 def: &RuleDef,
1021 k: u64,
1022 index: &RuleIndex,
1023 prov: &mut ProvSets<'_>,
1024 g: &mut GraphMut<'_>,
1025) {
1026 let src_sym = g.syms.get(&def.src_label);
1027 for id in 0..g.ids.len() as u32 {
1028 let label_sym = match g.labels.get(id as usize).copied() {
1029 Some(s) if s != u32::MAX => s,
1030 _ => continue,
1031 };
1032 if src_sym != Some(label_sym) {
1033 continue;
1034 }
1035 let per_src = compute_desired(def, index, id, true, g);
1036 let top_k = filter_src_top_k(per_src, k, g.ids);
1037 apply_per_src_top_k(def, id, top_k, prov, g);
1038 }
1039}
1040
1041fn apply_streaming_rebuild_top_k(
1048 def: &RuleDef,
1049 k: u64,
1050 index: &RuleIndex,
1051 prov: &mut ProvSets<'_>,
1052 g: &mut GraphMut<'_>,
1053) {
1054 let et = g.syms.intern(&def.edge_type);
1055
1056 let existing_srcs: BTreeSet<u32> = prov
1059 .set
1060 .iter()
1061 .filter(|(t, _, _)| *t == et)
1062 .map(|(_, s, _)| *s)
1063 .collect();
1064
1065 let src_sym = g.syms.get(&def.src_label);
1066 let mut all_srcs: BTreeSet<u32> = existing_srcs;
1067 for id in 0..g.ids.len() as u32 {
1068 let label_sym = match g.labels.get(id as usize).copied() {
1069 Some(s) if s != u32::MAX => s,
1070 _ => continue,
1071 };
1072 if src_sym == Some(label_sym) {
1073 all_srcs.insert(id);
1074 }
1075 }
1076
1077 for src in all_srcs {
1078 let desired_src = compute_desired(def, index, src, true, g);
1079 let top_k = filter_src_top_k(desired_src, k, g.ids);
1080 apply_per_src_top_k(def, src, top_k, prov, g);
1081 }
1082}
1083
1084fn apply_streaming_rebuild(
1097 def: &RuleDef,
1098 index: &RuleIndex,
1099 prov: &mut ProvSets<'_>,
1100 tripped: &mut bool,
1101 g: &mut GraphMut<'_>,
1102) {
1103 let budget = edge_budget(def);
1104 let et = g.syms.intern(&def.edge_type);
1105
1106 let total = count_desired_up_to(def, index, budget, g);
1108 if total > budget {
1109 *tripped = true;
1110 return; }
1112
1113 *tripped = false;
1115
1116 let current: Vec<Triple> = prov
1119 .set
1120 .iter()
1121 .filter(|(t, _, _)| *t == et)
1122 .copied()
1123 .collect();
1124 for (t, s, d) in current {
1125 if !pair_still_desired(def, s, d, g) {
1126 g.topo.remove_edge(t, s, d);
1127 g.edge_props.remove_edge(t, s, d);
1128 prov.remove(&def.name, (t, s, d), g.ids, g.syms);
1129 }
1130 }
1131
1132 let src_sym = g.syms.get(&def.src_label);
1135 for id in 0..g.ids.len() as u32 {
1136 let label_sym = match g.labels.get(id as usize).copied() {
1137 Some(s) if s != u32::MAX => s,
1138 _ => continue,
1139 };
1140 if src_sym != Some(label_sym) {
1141 continue;
1142 }
1143 let per_src = compute_desired(def, index, id, true, g);
1144 for ((s, d), score) in per_src {
1145 let triple = (et, s, d);
1146 let already = prov.contains(&triple);
1147 if !already {
1148 let newly = g.topo.add_edge(et, s, d);
1149 if newly {
1150 prov.insert(&def.name, triple, g.ids, g.syms);
1151 }
1152 }
1153 let is_owned_here = already || prov.contains(&triple);
1154 if is_owned_here {
1155 if let Some(p) = &def.weight_prop {
1156 g.edge_props.set(et, s, d, p, Value::Float(score));
1157 }
1158 }
1159 }
1160 }
1161}
1162
1163fn bump_fires_for_participants(def: &RuleDef, g: &GraphMut<'_>, fires: &mut u64) {
1166 let src_sym = g.syms.get(&def.src_label);
1167 let dst_sym = g.syms.get(&def.dst_label);
1168 for id in 0..g.ids.len() as u32 {
1169 let label_sym = match g.labels.get(id as usize).copied() {
1170 Some(s) if s != u32::MAX => s,
1171 _ => continue,
1172 };
1173 if src_sym == Some(label_sym) || dst_sym == Some(label_sym) {
1174 *fires += 1;
1175 }
1176 }
1177}
1178
1179fn index_node_for_rule(
1181 id: u32,
1182 label_sym: u32,
1183 def: &RuleDef,
1184 index: &mut RuleIndex,
1185 syms: &Interner,
1186 props: &ColumnStore,
1187) {
1188 let get = |f: &str| props.get(id, f).cloned();
1189 if syms.get(&def.src_label) == Some(label_sym) {
1190 let spec = src_lookup_spec_for(def);
1191 index.src_side.insert(&spec, id, &get);
1192 }
1193 if syms.get(&def.dst_label) == Some(label_sym) {
1194 let spec = candidate_spec_for(def);
1195 index.dst_side.insert(&spec, id, &get);
1196 }
1197}
1198
1199impl RuleEngine {
1204 pub fn new() -> Self {
1205 Self::default()
1206 }
1207
1208 pub fn rules(&self) -> impl Iterator<Item = &RuleDef> {
1209 self.rules.values()
1210 }
1211
1212 pub fn is_owned(&self, etype: u32, src: u32, dst: u32) -> bool {
1213 self.owned.contains(&(etype, src, dst))
1214 }
1215
1216 pub fn provenance(&self) -> &BTreeMap<String, BTreeSet<(u32, u32, u32)>> {
1218 &self.provenance
1219 }
1220
1221 pub fn provenance_touching(
1223 &self,
1224 node: u32,
1225 ) -> impl Iterator<Item = (&str, u32, u32, u32)> + '_ {
1226 self.by_node
1227 .get(&node)
1228 .into_iter()
1229 .flatten()
1230 .map(|&(rid, t, s, d)| (self.intern_rule[rid as usize].as_str(), t, s, d))
1231 }
1232
1233 pub fn provenance_touching_len(&self, node: u32) -> usize {
1235 self.by_node.get(&node).map_or(0, BTreeSet::len)
1236 }
1237
1238 pub fn is_tripped(&self, name: &str) -> bool {
1241 self.tripped.get(name).copied().unwrap_or(false)
1242 }
1243
1244 pub fn fire_count(&self, name: &str) -> u64 {
1248 self.fires.get(name).copied().unwrap_or(0)
1249 }
1250
1251 pub fn drain_deltas(&mut self) -> Vec<EngineEdgeDelta> {
1265 std::mem::take(&mut self.pending_deltas)
1266 }
1267
1268 pub fn pending_delta_count(&self) -> usize {
1271 self.pending_deltas.len()
1272 }
1273
1274 pub fn pending_deltas_since(&self, cursor: usize) -> &[EngineEdgeDelta] {
1283 &self.pending_deltas[cursor..]
1284 }
1285
1286 #[allow(clippy::type_complexity)]
1290 pub fn to_persist(
1291 &self,
1292 ) -> (
1293 Vec<RuleDef>,
1294 BTreeMap<String, BTreeSet<(u32, u32, u32)>>,
1295 BTreeMap<String, bool>,
1296 BTreeMap<String, u64>,
1297 ) {
1298 (
1299 self.rules.values().cloned().collect(),
1300 self.provenance.clone(),
1301 self.tripped.clone(),
1302 self.fires.clone(),
1303 )
1304 }
1305
1306 pub fn from_persist(
1308 rules: Vec<RuleDef>,
1309 prov: BTreeMap<String, BTreeSet<(u32, u32, u32)>>,
1310 tripped: BTreeMap<String, bool>,
1311 fires: BTreeMap<String, u64>,
1312 ) -> Self {
1313 let mut owned = BTreeSet::new();
1314 for set in prov.values() {
1315 owned.extend(set.iter().copied());
1316 }
1317 let indexes = rules
1318 .iter()
1319 .map(|r| (r.name.clone(), RuleIndex::default()))
1320 .collect();
1321 let rules: BTreeMap<String, RuleDef> =
1322 rules.into_iter().map(|r| (r.name.clone(), r)).collect();
1323 let mut tripped = tripped;
1325 let mut fires = fires;
1326 for name in rules.keys() {
1327 tripped.entry(name.clone()).or_insert(false);
1328 fires.entry(name.clone()).or_insert(0);
1329 }
1330 let (by_node, rule_intern, intern_rule) = rebuild_by_node(&prov);
1331 Self {
1332 rules,
1333 indexes,
1334 provenance: prov,
1335 owned,
1336 by_node,
1337 rule_intern,
1338 intern_rule,
1339 tripped,
1340 fires,
1341 pending_deltas: Vec::new(),
1342 emit_deltas: false,
1343 rebuild_needed: BTreeSet::new(),
1344 }
1345 }
1346
1347 pub fn set_emit_deltas(&mut self, emit: bool) {
1353 self.emit_deltas = emit;
1354 }
1355
1356 pub fn emit_deltas(&self) -> bool {
1358 self.emit_deltas
1359 }
1360
1361 pub fn take_rebuild_needed(&mut self) -> Vec<String> {
1364 std::mem::take(&mut self.rebuild_needed)
1365 .into_iter()
1366 .collect()
1367 }
1368
1369 pub fn queue_rebuild_needed(&mut self, name: String) {
1373 self.rebuild_needed.insert(name);
1374 }
1375
1376 fn maybe_queue_ivf_rebuild(&mut self, rule_name: &str, def: &RuleDef) {
1377 if !def.approximate {
1378 return;
1379 }
1380 let Some(idx) = self.indexes.get(rule_name) else {
1381 return;
1382 };
1383 if idx.dst_side.ivf_drift > ivf_drift_rebuild_threshold() {
1384 self.rebuild_needed.insert(rule_name.to_string());
1385 }
1386 }
1387
1388 pub fn export_ivf_state(&self) -> BTreeMap<String, RuleIvfExport> {
1392 let mut out = BTreeMap::new();
1393 for (name, def) in &self.rules {
1394 if def.approximate {
1395 if let Some(idx) = self.indexes.get(name) {
1396 out.insert(
1397 name.clone(),
1398 (
1399 idx.src_side.export_ivf_state(),
1400 idx.dst_side.export_ivf_state(),
1401 ),
1402 );
1403 }
1404 }
1405 }
1406 out
1407 }
1408
1409 pub fn reindex_all(
1411 &mut self,
1412 ids: &IdMap,
1413 syms: &Interner,
1414 labels: &[u32],
1415 props: &ColumnStore,
1416 ) {
1417 for idx in self.indexes.values_mut() {
1418 *idx = RuleIndex::default();
1419 }
1420 let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1423
1424 for name in &rule_names {
1426 if self.rules[name].approximate {
1427 let idx = self.indexes.get_mut(name).unwrap();
1428 idx.src_side.init_hnsw(name);
1429 idx.dst_side.init_hnsw(name);
1430 }
1431 }
1432
1433 for id in 0..ids.len() as u32 {
1434 let label_sym = match labels.get(id as usize).copied() {
1435 Some(s) if s != u32::MAX => s,
1436 _ => continue,
1437 };
1438 for name in &rule_names {
1439 let def = self.rules[name].clone();
1440 let idx = self.indexes.get_mut(name).unwrap();
1441 index_node_for_rule(id, label_sym, &def, idx, syms, props);
1442 }
1443 }
1444 for name in &rule_names {
1447 if self.rules[name].approximate {
1448 let idx = self.indexes.get_mut(name).unwrap();
1449 idx.src_side.fit_ivf_clusters(name);
1450 idx.dst_side.fit_ivf_clusters(name);
1451 }
1452 }
1453 }
1454
1455 pub fn reindex_all_load_ivf(
1465 &mut self,
1466 ids: &IdMap,
1467 syms: &Interner,
1468 labels: &[u32],
1469 props: &ColumnStore,
1470 ivf_state: BTreeMap<String, RuleIvfExport>,
1471 ) {
1472 for idx in self.indexes.values_mut() {
1473 *idx = RuleIndex::default();
1474 }
1475 let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1476
1477 for name in &rule_names {
1480 if self.rules[name].approximate {
1481 let idx = self.indexes.get_mut(name).unwrap();
1482 idx.src_side.init_hnsw(name);
1483 idx.dst_side.init_hnsw(name);
1484 }
1485 }
1486
1487 for id in 0..ids.len() as u32 {
1488 let label_sym = match labels.get(id as usize).copied() {
1489 Some(s) if s != u32::MAX => s,
1490 _ => continue,
1491 };
1492 for name in &rule_names {
1493 let def = self.rules[name].clone();
1494 let idx = self.indexes.get_mut(name).unwrap();
1495 index_node_for_rule(id, label_sym, &def, idx, syms, props);
1496 }
1497 }
1498 for name in &rule_names {
1502 if !self.rules[name].approximate {
1503 continue;
1504 }
1505 let idx = self.indexes.get_mut(name).unwrap();
1506 if let Some(((sc, sa, sd), (dc, da, dd))) = ivf_state.get(name) {
1507 idx.src_side.load_ivf_state(sc.clone(), sa.clone(), *sd);
1508 idx.dst_side.load_ivf_state(dc.clone(), da.clone(), *dd);
1509 } else {
1510 idx.src_side.fit_ivf_clusters(name);
1512 idx.dst_side.fit_ivf_clusters(name);
1513 }
1514 }
1515 }
1516
1517 pub fn export_hnsw_state(&self) -> BTreeMap<String, (Vec<u8>, Vec<u8>)> {
1522 let mut out = BTreeMap::new();
1523 for (name, def) in &self.rules {
1524 if def.approximate {
1525 if let Some(idx) = self.indexes.get(name) {
1526 out.insert(
1527 name.clone(),
1528 (
1529 idx.src_side.export_hnsw_blob(),
1530 idx.dst_side.export_hnsw_blob(),
1531 ),
1532 );
1533 }
1534 }
1535 }
1536 out
1537 }
1538
1539 pub fn load_hnsw_state(&mut self, blobs: BTreeMap<String, (Vec<u8>, Vec<u8>)>) {
1544 for (name, (src_blob, dst_blob)) in blobs {
1545 if let Some(idx) = self.indexes.get_mut(&name) {
1546 if !src_blob.is_empty() {
1547 idx.src_side.load_hnsw_blob(&src_blob);
1548 }
1549 if !dst_blob.is_empty() {
1550 idx.dst_side.load_hnsw_blob(&dst_blob);
1551 }
1552 }
1553 }
1554 }
1555
1556 pub fn hnsw_search_dst(
1561 &self,
1562 field: &str,
1563 dst_label: &str,
1564 q: &[f64],
1565 k: usize,
1566 ) -> Option<Vec<(u32, f64)>> {
1567 for (name, def) in &self.rules {
1568 if !def.approximate || def.dst_label != dst_label {
1569 continue;
1570 }
1571 if !predicate_covers_field(&def.predicate, field) {
1573 continue;
1574 }
1575 if let Some(idx) = self.indexes.get(name) {
1576 if let Some(h) = idx.dst_side.hnsw_ref() {
1577 if !h.is_empty() {
1578 return Some(h.search(q, k));
1579 }
1580 }
1581 }
1582 }
1583 None
1584 }
1585
1586 pub fn create_rule(&mut self, def: RuleDef, g: &mut GraphMut<'_>) -> Result<(), String> {
1589 def.validate()?;
1590 if self.rules.contains_key(&def.name) {
1591 return Err(format!("rule {:?} already exists", def.name));
1592 }
1593 let name = def.name.clone();
1594 self.rules.insert(name.clone(), def);
1595 self.indexes.insert(name.clone(), RuleIndex::default());
1596 self.provenance.entry(name.clone()).or_default();
1597 self.tripped.insert(name.clone(), false);
1598 self.fires.insert(name.clone(), 0);
1599
1600 let n_total = g.ids.len() as u32;
1602 let def = self.rules[&name].clone();
1603
1604 if def.approximate {
1607 let idx = self.indexes.get_mut(&name).unwrap();
1608 idx.src_side.init_hnsw(&name);
1609 idx.dst_side.init_hnsw(&name);
1610 }
1611
1612 for id in 0..n_total {
1613 let label_sym = match g.labels.get(id as usize).copied() {
1614 Some(s) if s != u32::MAX => s,
1615 _ => continue,
1616 };
1617 let idx = self.indexes.get_mut(&name).unwrap();
1618 index_node_for_rule(id, label_sym, &def, idx, g.syms, g.props);
1619 }
1620
1621 if def.approximate {
1624 let idx = self.indexes.get_mut(&name).unwrap();
1625 idx.src_side.fit_ivf_clusters(&name);
1626 idx.dst_side.fit_ivf_clusters(&name);
1627 }
1628
1629 let mut prov = ProvSets {
1634 set: self.provenance.get_mut(&name).unwrap(),
1635 owned: &mut self.owned,
1636 by_node: &mut self.by_node,
1637 rule_intern: &mut self.rule_intern,
1638 intern_rule: &mut self.intern_rule,
1639 deltas: &mut self.pending_deltas,
1640 emit: self.emit_deltas,
1641 };
1642 if def.via_label.is_some() {
1643 let budget = edge_budget(&def);
1645 let et = g.syms.intern(&def.edge_type);
1646 let src_sym = g.syms.get(&def.src_label);
1647 let tripped = self.tripped.get_mut(&name).unwrap();
1648 'via_outer: for id in 0..g.ids.len() as u32 {
1649 let label_sym = match g.labels.get(id as usize).copied() {
1650 Some(s) if s != u32::MAX => s,
1651 _ => continue,
1652 };
1653 if src_sym != Some(label_sym) {
1654 continue;
1655 }
1656 let per_src = compute_desired_via(&def, ViaAnchor::Src(id), g);
1657 if let Some(k) = def.max_edges {
1658 let top_k = filter_src_top_k(per_src, k, g.ids);
1659 apply_per_src_top_k(&def, id, top_k, &mut prov, g);
1660 } else {
1661 for ((s, d), score) in per_src {
1662 let triple = (et, s, d);
1663 let already = prov.contains(&triple);
1664 if !already {
1665 if *tripped || prov.len() as u64 >= budget {
1666 *tripped = true;
1667 break 'via_outer;
1668 }
1669 let newly = g.topo.add_edge(et, s, d);
1670 if newly {
1671 prov.insert(&name, triple, g.ids, g.syms);
1672 }
1673 }
1674 let is_owned_here = already || prov.contains(&triple);
1675 if is_owned_here {
1676 if let Some(p) = &def.weight_prop {
1677 g.edge_props.set(et, s, d, p, Value::Float(score));
1678 }
1679 }
1680 }
1681 }
1682 }
1683 } else if let Some(k) = def.max_edges {
1684 apply_streaming_create_top_k(&def, k, &self.indexes[&name], &mut prov, g);
1685 } else {
1686 let tripped = self.tripped.get_mut(&name).unwrap();
1687 apply_streaming_create(&def, &self.indexes[&name], &mut prov, tripped, g);
1688 }
1689 let fires = self.fires.get_mut(&name).unwrap();
1692 bump_fires_for_participants(&def, g, fires);
1693
1694 Ok(())
1695 }
1696
1697 pub fn delete_rule(&mut self, name: &str, g: &mut GraphMut<'_>) -> Result<(), String> {
1699 if !self.rules.contains_key(name) {
1700 return Err(format!("rule {:?} not found", name));
1701 }
1702 let def = self.rules.remove(name).unwrap();
1703 self.indexes.remove(name);
1704 self.tripped.remove(name);
1705 self.fires.remove(name);
1706 let mut leftover = self.provenance.remove(name).unwrap_or_default();
1707 let _et = g.syms.intern(&def.edge_type);
1709 let triples: Vec<Triple> = leftover.iter().copied().collect();
1710 let mut sets = ProvSets {
1711 set: &mut leftover,
1712 owned: &mut self.owned,
1713 by_node: &mut self.by_node,
1714 rule_intern: &mut self.rule_intern,
1715 intern_rule: &mut self.intern_rule,
1716 deltas: &mut self.pending_deltas,
1717 emit: self.emit_deltas,
1718 };
1719 for triple in triples {
1720 let (t, s, d) = triple;
1721 g.topo.remove_edge(t, s, d);
1722 g.edge_props.remove_edge(t, s, d);
1723 sets.remove(name, triple, g.ids, g.syms);
1724 }
1725 let same_etype_survivors: Vec<String> = self
1731 .rules
1732 .values()
1733 .filter(|r| r.edge_type == def.edge_type)
1734 .map(|r| r.name.clone())
1735 .collect();
1736 for survivor in same_etype_survivors {
1737 let _ = self.rebuild(&survivor, g);
1739 }
1740 Ok(())
1741 }
1742
1743 pub fn on_node_changed(
1753 &mut self,
1754 n: u32,
1755 changed: Option<(&str, Option<Value>)>,
1756 g: &mut GraphMut<'_>,
1757 ) {
1758 let n_label = g.labels.get(n as usize).copied();
1759 let rule_names: Vec<String> = self.rules.keys().cloned().collect();
1760
1761 for rule_name in rule_names {
1762 let def = self.rules[&rule_name].clone();
1763
1764 if def.via_label.is_some() {
1765 self.on_node_changed_via(&rule_name, &def, n, n_label, changed.clone(), g);
1767 } else {
1768 let src_sym = g.syms.get(&def.src_label);
1770 let dst_sym = g.syms.get(&def.dst_label);
1771 let as_src = src_sym.is_some() && n_label == src_sym;
1772 let as_dst = dst_sym.is_some() && n_label == dst_sym;
1773
1774 let fires = match changed {
1775 None => as_src || as_dst,
1776 Some((field, _)) => def.watched_fields().contains(field) && (as_src || as_dst),
1777 };
1778 if !fires {
1779 continue;
1780 }
1781 *self.fires.entry(rule_name.clone()).or_default() += 1;
1782
1783 if let Some((field, ref old_val)) = changed {
1785 let old_val_cloned = old_val.clone();
1786 let old_getter = |f: &str| {
1787 if f == field {
1788 old_val_cloned.clone()
1789 } else {
1790 g.props.get(n, f).cloned()
1791 }
1792 };
1793 let idx = self.indexes.get_mut(&rule_name).unwrap();
1794 if as_src {
1795 let spec = src_lookup_spec_for(&def);
1796 idx.src_side.remove(&spec, n, &old_getter);
1797 }
1798 if as_dst {
1799 let spec = candidate_spec_for(&def);
1800 idx.dst_side.remove(&spec, n, &old_getter);
1801 }
1802 }
1803
1804 {
1805 let cur_getter = |f: &str| g.props.get(n, f).cloned();
1806 let idx = self.indexes.get_mut(&rule_name).unwrap();
1807 if as_src {
1808 let spec = src_lookup_spec_for(&def);
1809 idx.src_side.insert(&spec, n, &cur_getter);
1810 }
1811 if as_dst {
1812 let spec = candidate_spec_for(&def);
1813 idx.dst_side.insert(&spec, n, &cur_getter);
1814 }
1815 }
1816
1817 self.maybe_queue_ivf_rebuild(&rule_name, &def);
1818
1819 if let Some(k) = def.max_edges {
1821 let et = g.syms.intern(&def.edge_type);
1822 let affected_srcs_for_n_dst: BTreeSet<u32> = if as_dst {
1823 let rid = self.rule_intern.get(&def.name).copied();
1824 self.by_node
1825 .get(&n)
1826 .into_iter()
1827 .flatten()
1828 .filter(|(r, t, _s, d)| Some(*r) == rid && *t == et && *d == n)
1829 .map(|(_, _, s, _)| *s)
1830 .collect()
1831 } else {
1832 BTreeSet::new()
1833 };
1834
1835 let mut prov = ProvSets {
1836 set: self.provenance.entry(rule_name.clone()).or_default(),
1837 owned: &mut self.owned,
1838 by_node: &mut self.by_node,
1839 rule_intern: &mut self.rule_intern,
1840 intern_rule: &mut self.intern_rule,
1841 deltas: &mut self.pending_deltas,
1842 emit: self.emit_deltas,
1843 };
1844
1845 if as_src {
1846 let desired_n_src =
1847 compute_desired(&def, &self.indexes[&rule_name], n, true, g);
1848 let top_k = filter_src_top_k(desired_n_src, k, g.ids);
1849 apply_per_src_top_k(&def, n, top_k, &mut prov, g);
1850 }
1851
1852 if as_dst {
1853 let new_desired =
1854 compute_desired(&def, &self.indexes[&rule_name], n, false, g);
1855 let new_srcs: BTreeSet<u32> = new_desired.keys().map(|(s, _)| *s).collect();
1856 let affected_srcs: BTreeSet<u32> =
1857 affected_srcs_for_n_dst.union(&new_srcs).copied().collect();
1858 for src in affected_srcs {
1859 if src == n {
1860 continue;
1861 }
1862 let desired_src =
1863 compute_desired(&def, &self.indexes[&rule_name], src, true, g);
1864 let top_k = filter_src_top_k(desired_src, k, g.ids);
1865 apply_per_src_top_k(&def, src, top_k, &mut prov, g);
1866 }
1867 }
1868 } else {
1869 let mut desired = BTreeMap::new();
1870 if as_src {
1871 desired.extend(compute_desired(
1872 &def,
1873 &self.indexes[&rule_name],
1874 n,
1875 true,
1876 g,
1877 ));
1878 }
1879 if as_dst {
1880 desired.extend(compute_desired(
1881 &def,
1882 &self.indexes[&rule_name],
1883 n,
1884 false,
1885 g,
1886 ));
1887 }
1888 let tripped = self.tripped.entry(rule_name.clone()).or_default();
1889 apply_desired(
1890 &def,
1891 desired,
1892 Some(n),
1893 &mut ProvSets {
1894 set: self.provenance.entry(rule_name).or_default(),
1895 owned: &mut self.owned,
1896 by_node: &mut self.by_node,
1897 rule_intern: &mut self.rule_intern,
1898 intern_rule: &mut self.intern_rule,
1899 deltas: &mut self.pending_deltas,
1900 emit: self.emit_deltas,
1901 },
1902 tripped,
1903 g,
1904 );
1905 }
1906 }
1907 }
1908 }
1909
1910 fn on_node_changed_via(
1925 &mut self,
1926 rule_name: &str,
1927 def: &RuleDef,
1928 n: u32,
1929 n_label: Option<u32>,
1930 changed: Option<(&str, Option<Value>)>,
1931 g: &mut GraphMut<'_>,
1932 ) {
1933 let src_sym = g.syms.get(&def.src_label);
1934 let dst_sym = g.syms.get(&def.dst_label);
1935 let via_sym = def.via_label.as_deref().and_then(|l| g.syms.get(l));
1936
1937 let as_src = src_sym.is_some() && n_label == src_sym;
1938 let as_dst = dst_sym.is_some() && n_label == dst_sym;
1939 let as_via = via_sym.is_some() && n_label == via_sym;
1940
1941 let fires = match changed {
1945 None => as_src || as_via || as_dst,
1946 Some((field, _)) => {
1947 let wf = def.watched_fields();
1948 (wf.contains(field)) && (as_src || as_via || as_dst)
1949 }
1950 };
1951 if !fires {
1952 return;
1953 }
1954 *self.fires.entry(rule_name.to_string()).or_default() += 1;
1955
1956 let mut affected_srcs: BTreeSet<u32> = BTreeSet::new();
1958 if as_src {
1959 affected_srcs.insert(n);
1960 }
1961 if as_via {
1962 let via_edge_str = def.via_edge.as_deref().unwrap();
1964 let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
1965 let rev_dir = match via_dir {
1966 core_storage::Direction::Out => core_storage::Direction::In,
1967 core_storage::Direction::In => core_storage::Direction::Out,
1968 };
1969 if let (Some(via_etype), Some(s_sym)) = (g.syms.get(via_edge_str), src_sym) {
1970 for &src in g.topo.neighbors(via_etype, rev_dir, n).as_ref() {
1971 if g.labels.get(src as usize).copied() == Some(s_sym) {
1972 affected_srcs.insert(src);
1973 }
1974 }
1975 }
1976 }
1977 if as_dst {
1978 let desired_touching_n = compute_desired_via(def, ViaAnchor::Dst(n), g);
1980 for (src, _dst) in desired_touching_n.keys() {
1981 affected_srcs.insert(*src);
1982 }
1983 let et = g.syms.intern(&def.edge_type);
1985 let rid = self.rule_intern.get(rule_name).copied();
1986 let old_srcs: Vec<u32> = self
1987 .by_node
1988 .get(&n)
1989 .into_iter()
1990 .flatten()
1991 .filter(|(r, t, _s, d)| Some(*r) == rid && *t == et && *d == n)
1992 .map(|(_, _, s, _)| *s)
1993 .collect();
1994 affected_srcs.extend(old_srcs);
1995 }
1996
1997 let affected_srcs: Vec<u32> = affected_srcs.into_iter().collect();
1999
2000 if let Some(k) = def.max_edges {
2001 let mut prov = ProvSets {
2002 set: self.provenance.entry(rule_name.to_string()).or_default(),
2003 owned: &mut self.owned,
2004 by_node: &mut self.by_node,
2005 rule_intern: &mut self.rule_intern,
2006 intern_rule: &mut self.intern_rule,
2007 deltas: &mut self.pending_deltas,
2008 emit: self.emit_deltas,
2009 };
2010 for src in affected_srcs {
2011 let desired_src = compute_desired_via(def, ViaAnchor::Src(src), g);
2012 let top_k = filter_src_top_k(desired_src, k, g.ids);
2013 apply_per_src_top_k(def, src, top_k, &mut prov, g);
2014 }
2015 } else {
2016 let tripped = self.tripped.entry(rule_name.to_string()).or_default();
2017 let budget = edge_budget(def);
2018 for src in affected_srcs {
2021 let desired_src = compute_desired_via(def, ViaAnchor::Src(src), g);
2022 if !*tripped {
2023 let mut prov = ProvSets {
2024 set: self.provenance.entry(rule_name.to_string()).or_default(),
2025 owned: &mut self.owned,
2026 by_node: &mut self.by_node,
2027 rule_intern: &mut self.rule_intern,
2028 intern_rule: &mut self.intern_rule,
2029 deltas: &mut self.pending_deltas,
2030 emit: self.emit_deltas,
2031 };
2032 apply_desired(def, desired_src, Some(src), &mut prov, tripped, g);
2033 }
2034 let _ = budget;
2038 }
2039 }
2040 }
2041
2042 pub fn on_edge_changed(
2055 &mut self,
2056 etype_str: &str,
2057 src_id: u32,
2058 dst_id: u32,
2059 g: &mut GraphMut<'_>,
2060 ) {
2061 let rule_names: Vec<String> = self.rules.keys().cloned().collect();
2062 for rule_name in rule_names {
2063 let def = self.rules[&rule_name].clone();
2064 let Some(ref via_edge) = def.via_edge else {
2065 continue; };
2067 if via_edge != etype_str {
2068 continue; }
2070
2071 let src_sym = match g.syms.get(&def.src_label) {
2073 Some(s) => s,
2074 None => continue,
2075 };
2076 let via_sym = match def.via_label.as_deref().and_then(|l| g.syms.get(l)) {
2077 Some(s) => s,
2078 None => continue,
2079 };
2080 let via_dir = def.via_dir.unwrap_or(core_storage::Direction::Out);
2084 let (rule_src, rule_via) = match via_dir {
2085 core_storage::Direction::Out => (src_id, dst_id),
2086 core_storage::Direction::In => (dst_id, src_id),
2087 };
2088
2089 if g.labels.get(rule_src as usize).copied() != Some(src_sym) {
2090 continue;
2091 }
2092 if g.labels.get(rule_via as usize).copied() != Some(via_sym) {
2093 continue;
2094 }
2095
2096 *self.fires.entry(rule_name.clone()).or_default() += 1;
2098 let desired_src = compute_desired_via(&def, ViaAnchor::Src(rule_src), g);
2099
2100 if let Some(k) = def.max_edges {
2101 let mut prov = ProvSets {
2102 set: self.provenance.entry(rule_name).or_default(),
2103 owned: &mut self.owned,
2104 by_node: &mut self.by_node,
2105 rule_intern: &mut self.rule_intern,
2106 intern_rule: &mut self.intern_rule,
2107 deltas: &mut self.pending_deltas,
2108 emit: self.emit_deltas,
2109 };
2110 let top_k = filter_src_top_k(desired_src, k, g.ids);
2111 apply_per_src_top_k(&def, rule_src, top_k, &mut prov, g);
2112 } else {
2113 let tripped = self.tripped.entry(rule_name.clone()).or_default();
2114 let mut prov = ProvSets {
2115 set: self.provenance.entry(rule_name).or_default(),
2116 owned: &mut self.owned,
2117 by_node: &mut self.by_node,
2118 rule_intern: &mut self.rule_intern,
2119 intern_rule: &mut self.intern_rule,
2120 deltas: &mut self.pending_deltas,
2121 emit: self.emit_deltas,
2122 };
2123 apply_desired(&def, desired_src, Some(rule_src), &mut prov, tripped, g);
2124 }
2125 }
2126 }
2127
2128 pub fn on_node_removed(&mut self, n: u32, g: &mut GraphMut<'_>) {
2136 let n_label = g.labels.get(n as usize).copied();
2137 let rule_names: Vec<String> = self.rules.keys().cloned().collect();
2138
2139 for rule_name in rule_names {
2140 let def = self.rules[&rule_name].clone();
2141 let src_sym = g.syms.get(&def.src_label);
2142 let dst_sym = g.syms.get(&def.dst_label);
2143 let as_src = src_sym.is_some() && n_label == src_sym;
2144 let as_dst = dst_sym.is_some() && n_label == dst_sym;
2145
2146 {
2147 let cur_getter = |f: &str| g.props.get(n, f).cloned();
2148 let idx = self.indexes.get_mut(&rule_name).unwrap();
2149 if as_src {
2150 let spec = src_lookup_spec_for(&def);
2151 idx.src_side.remove(&spec, n, &cur_getter);
2152 }
2153 if as_dst {
2154 let spec = candidate_spec_for(&def);
2155 idx.dst_side.remove(&spec, n, &cur_getter);
2156 }
2157 }
2158
2159 self.maybe_queue_ivf_rebuild(&rule_name, &def);
2160 }
2161
2162 let touching: Vec<(String, Triple)> = self
2163 .by_node
2164 .get(&n)
2165 .into_iter()
2166 .flatten()
2167 .map(|&(rid, t, s, d)| (self.intern_rule[rid as usize].clone(), (t, s, d)))
2168 .collect();
2169
2170 let topk_backfill: Vec<(String, u32)> = touching
2174 .iter()
2175 .filter_map(|(rule_name, triple)| {
2176 let &(_, s, d) = triple;
2177 let def = self.rules.get(rule_name)?;
2178 def.max_edges?; if d == n && s != n {
2180 Some((rule_name.clone(), s))
2181 } else {
2182 None
2183 }
2184 })
2185 .collect();
2186
2187 for (rule_name, triple) in touching {
2188 let (t, s, d) = triple;
2189 g.topo.remove_edge(t, s, d);
2190 g.edge_props.remove_edge(t, s, d);
2191 if let Some(set) = self.provenance.get_mut(&rule_name) {
2192 ProvSets {
2193 set,
2194 owned: &mut self.owned,
2195 by_node: &mut self.by_node,
2196 rule_intern: &mut self.rule_intern,
2197 intern_rule: &mut self.intern_rule,
2198 deltas: &mut self.pending_deltas,
2199 emit: self.emit_deltas,
2200 }
2201 .remove(&rule_name, triple, g.ids, g.syms);
2202 }
2203 }
2204
2205 for (rule_name, src) in topk_backfill {
2210 let def = self.rules[&rule_name].clone();
2211 let k = def.max_edges.unwrap(); let desired_src = compute_desired(&def, &self.indexes[&rule_name], src, true, g);
2213 let top_k = filter_src_top_k(desired_src, k, g.ids);
2214 let mut prov = ProvSets {
2215 set: self.provenance.entry(rule_name.clone()).or_default(),
2216 owned: &mut self.owned,
2217 by_node: &mut self.by_node,
2218 rule_intern: &mut self.rule_intern,
2219 intern_rule: &mut self.intern_rule,
2220 deltas: &mut self.pending_deltas,
2221 emit: self.emit_deltas,
2222 };
2223 apply_per_src_top_k(&def, src, top_k, &mut prov, g);
2224 }
2225 }
2226
2227 pub fn rebuild(&mut self, name: &str, g: &mut GraphMut<'_>) -> Result<(), String> {
2235 if !self.rules.contains_key(name) {
2236 return Err(format!("rule {:?} not found", name));
2237 }
2238 self.rebuild_needed.remove(name);
2239 let def = self.rules[name].clone();
2240
2241 *self.indexes.get_mut(name).unwrap() = RuleIndex::default();
2243
2244 if def.approximate {
2246 let idx = self.indexes.get_mut(name).unwrap();
2247 idx.src_side.init_hnsw(name);
2248 idx.dst_side.init_hnsw(name);
2249 }
2250
2251 let n_total = g.ids.len() as u32;
2252 for id in 0..n_total {
2253 let label_sym = match g.labels.get(id as usize).copied() {
2254 Some(s) if s != u32::MAX => s,
2255 _ => continue,
2256 };
2257 let idx = self.indexes.get_mut(name).unwrap();
2258 index_node_for_rule(id, label_sym, &def, idx, g.syms, g.props);
2259 }
2260
2261 if def.approximate {
2264 let idx = self.indexes.get_mut(name).unwrap();
2265 idx.src_side.fit_ivf_clusters(name);
2266 idx.dst_side.fit_ivf_clusters(name);
2267 }
2268
2269 let mut prov = ProvSets {
2273 set: self.provenance.get_mut(name).unwrap(),
2274 owned: &mut self.owned,
2275 by_node: &mut self.by_node,
2276 rule_intern: &mut self.rule_intern,
2277 intern_rule: &mut self.intern_rule,
2278 deltas: &mut self.pending_deltas,
2279 emit: self.emit_deltas,
2280 };
2281 if let Some(k) = def.max_edges {
2282 apply_streaming_rebuild_top_k(&def, k, &self.indexes[name], &mut prov, g);
2283 } else {
2284 let tripped = self.tripped.get_mut(name).unwrap();
2285 apply_streaming_rebuild(&def, &self.indexes[name], &mut prov, tripped, g);
2286 }
2287 let fires = self.fires.entry(name.to_string()).or_default();
2288 bump_fires_for_participants(&def, g, fires);
2289
2290 Ok(())
2291 }
2292
2293 #[cfg(test)]
2294 fn by_node_consistent(&self) -> bool {
2295 let (rebuilt, intern, names) = rebuild_by_node(&self.provenance);
2296 resolve_by_node(&self.by_node, &self.intern_rule) == resolve_by_node(&rebuilt, &names)
2297 && intern.len() == names.len()
2298 }
2299}
2300
2301#[cfg(test)]
2306mod tests {
2307 use super::*;
2308 use crate::def::{evaluate, NodeView, Predicate, RuleDef};
2309 use core_storage::{ColumnStore, Direction, EdgeProps, IdMap, Interner, Topology, Value};
2310
2311 struct Fx {
2312 ids: IdMap,
2313 syms: Interner,
2314 labels: Vec<u32>,
2315 props: ColumnStore,
2316 topo: Topology,
2317 eprops: EdgeProps,
2318 }
2319 impl Fx {
2320 fn new() -> Self {
2321 Fx {
2322 ids: IdMap::new(),
2323 syms: Interner::new(),
2324 labels: vec![],
2325 props: ColumnStore::new(),
2326 topo: Topology::new(),
2327 eprops: EdgeProps::new(),
2328 }
2329 }
2330 fn add(&mut self, label: &str, key: &str, props: Vec<(&str, Value)>) -> u32 {
2331 let id = self.ids.get_or_insert(key);
2332 let sym = self.syms.intern(label);
2333 self.labels.resize(id as usize + 1, u32::MAX);
2334 self.labels[id as usize] = sym;
2335 for (f, v) in props {
2336 self.props.set(id, f, v);
2337 }
2338 id
2339 }
2340 fn g(&mut self) -> GraphMut<'_> {
2341 GraphMut {
2342 ids: &self.ids,
2343 syms: &mut self.syms,
2344 labels: &self.labels,
2345 props: &self.props,
2346 topo: &mut self.topo,
2347 edge_props: &mut self.eprops,
2348 }
2349 }
2350 }
2351
2352 fn tags(items: &[&str]) -> Value {
2353 Value::List(items.iter().map(|s| Value::Str((*s).into())).collect())
2354 }
2355
2356 fn overlap_rule() -> RuleDef {
2357 RuleDef {
2358 name: "rel".into(),
2359 src_label: "A".into(),
2360 dst_label: "A".into(),
2361 predicate: Predicate::Overlap {
2362 field: "tags".into(),
2363 min: 0.4,
2364 },
2365 edge_type: "REL".into(),
2366 weight_prop: Some("score".into()),
2367 max_edges: None,
2368 approximate: false,
2369 via_label: None,
2370 via_edge: None,
2371 via_dir: None,
2372 }
2373 }
2374
2375 fn emb(xs: &[f64]) -> Value {
2376 Value::List(xs.iter().copied().map(Value::Float).collect())
2377 }
2378
2379 fn approx_vec_rule() -> RuleDef {
2380 RuleDef {
2381 name: "sim".into(),
2382 src_label: "V".into(),
2383 dst_label: "V".into(),
2384 predicate: Predicate::VectorSimilar {
2385 field: "emb".into(),
2386 min: 0.5,
2387 },
2388 edge_type: "SIM".into(),
2389 weight_prop: None,
2390 max_edges: None,
2391 approximate: true,
2392 via_label: None,
2393 via_edge: None,
2394 via_dir: None,
2395 }
2396 }
2397
2398 #[test]
2399 fn approximate_rule_rebuilds_after_drift_threshold() {
2400 with_ivf_drift_rebuild(1, || {
2401 let mut fx = Fx::new();
2402 let mut ids = Vec::new();
2403 for i in 0..6 {
2404 let x = i as f64 * 0.2;
2405 ids.push(fx.add("V", &format!("v{i}"), vec![("emb", emb(&[x, 1.0 - x]))]));
2406 }
2407 let mut eng = RuleEngine::new();
2408 {
2409 let mut g = fx.g();
2410 eng.create_rule(approx_vec_rule(), &mut g).unwrap();
2411 }
2412 assert!(eng.take_rebuild_needed().is_empty());
2413 {
2414 let mut g = fx.g();
2415 eng.on_node_removed(ids[0], &mut g);
2416 }
2417 assert!(
2418 eng.take_rebuild_needed().is_empty(),
2419 "drift=1 is not > threshold 1"
2420 );
2421 {
2422 let mut g = fx.g();
2423 eng.on_node_removed(ids[1], &mut g);
2424 }
2425 assert_eq!(eng.take_rebuild_needed(), vec!["sim".to_string()]);
2426 {
2427 let mut g = fx.g();
2428 eng.rebuild("sim", &mut g).unwrap();
2429 }
2430 assert!(
2431 eng.take_rebuild_needed().is_empty(),
2432 "rebuild must reset drift and not re-queue itself"
2433 );
2434 let drift = eng
2435 .export_ivf_state()
2436 .get("sim")
2437 .map(|(_, dst)| dst.2)
2438 .unwrap();
2439 assert_eq!(drift, 0, "rebuild resets dst-side IVF drift");
2440 });
2441 }
2442
2443 #[test]
2444 fn backfill_creates_edges_with_scores_and_delete_removes_exactly_them() {
2445 let mut fx = Fx::new();
2446 let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2447 let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2448 let _c = fx.add("A", "c", vec![("tags", tags(&["q"]))]);
2449 let et = fx.syms.intern("REL");
2451 fx.topo.add_edge(et, a, b);
2452 let mut eng = RuleEngine::new();
2453 let mut g = fx.g();
2454 eng.create_rule(overlap_rule(), &mut g).unwrap();
2455 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2457 assert_eq!(
2458 g.edge_props.get(et, b, a, "score"),
2459 Some(&Value::Float(1.0))
2460 );
2461 assert!(!eng.is_owned(et, a, b));
2462 assert!(eng.is_owned(et, b, a));
2463 eng.delete_rule("rel", &mut g).unwrap();
2464 assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b)); assert!(!g.topo.neighbors(et, Direction::Out, b).contains(&a)); assert_eq!(g.edge_props.get(et, b, a, "score"), None);
2467 }
2468
2469 #[test]
2470 fn incremental_update_adds_and_removes_edges() {
2471 let mut fx = Fx::new();
2472 let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2473 let b = fx.add("A", "b", vec![("tags", tags(&["y", "z"]))]);
2474 let et = fx.syms.intern("REL");
2475 let mut eng = RuleEngine::new();
2476 {
2477 let mut g = fx.g();
2478 eng.create_rule(overlap_rule(), &mut g).unwrap(); assert_eq!(g.topo.edge_count(), 0);
2480 }
2481 let old = fx.props.get(b, "tags").cloned();
2483 fx.props.set(b, "tags", tags(&["x", "y"]));
2484 {
2485 let mut g = fx.g();
2486 eng.on_node_changed(b, Some(("tags", old)), &mut g);
2487 assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2488 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2489 }
2490 let old = fx.props.get(b, "tags").cloned();
2492 fx.props.set(b, "tags", tags(&["qqq"]));
2493 let mut g = fx.g();
2494 eng.on_node_changed(b, Some(("tags", old)), &mut g);
2495 assert_eq!(g.topo.edge_count(), 0);
2496 assert_eq!(g.edge_props.get(et, a, b, "score"), None);
2497 }
2498
2499 #[test]
2500 fn key_match_new_node_links_and_rebuild_is_noop() {
2501 let mut fx = Fx::new();
2502 fx.add("C", "c1", vec![]);
2503 let mut eng = RuleEngine::new();
2504 {
2505 let mut g = fx.g();
2506 eng.create_rule(
2507 RuleDef {
2508 name: "fk".into(),
2509 src_label: "T".into(),
2510 dst_label: "C".into(),
2511 predicate: Predicate::KeyMatch {
2512 field: "cid".into(),
2513 },
2514 edge_type: "AT".into(),
2515 weight_prop: None,
2516 max_edges: None,
2517 approximate: false,
2518 via_label: None,
2519 via_edge: None,
2520 via_dir: None,
2521 },
2522 &mut g,
2523 )
2524 .unwrap();
2525 }
2526 let t = fx.add("T", "t1", vec![("cid", Value::Str("c1".into()))]);
2527 let (at, c1, count_before) = {
2528 let mut g = fx.g();
2529 eng.on_node_changed(t, None, &mut g);
2530 let at = g.syms.get("AT").unwrap();
2531 let c1 = g.ids.get("c1").unwrap();
2532 assert!(g.topo.neighbors(at, Direction::Out, t).contains(&c1));
2533 (at, c1, g.topo.edge_count())
2534 };
2535 let mut g = fx.g();
2536 eng.rebuild("fk", &mut g).unwrap();
2537 assert_eq!(g.topo.edge_count(), count_before); assert!(g.topo.neighbors(at, Direction::Out, t).contains(&c1));
2539 }
2540
2541 #[test]
2542 fn score_refresh_on_persisting_owned_edge() {
2543 let mut fx = Fx::new();
2546 let a = fx.add("A", "a", vec![("tags", tags(&["x", "y", "z"]))]);
2547 let b = fx.add("A", "b", vec![("tags", tags(&["x", "y", "q"]))]);
2548 let et = fx.syms.intern("SIM");
2549 let mut eng = RuleEngine::new();
2550 {
2551 let mut g = fx.g();
2552 eng.create_rule(
2553 RuleDef {
2554 name: "sim".into(),
2555 src_label: "A".into(),
2556 dst_label: "A".into(),
2557 predicate: Predicate::Overlap {
2558 field: "tags".into(),
2559 min: 0.2,
2560 },
2561 edge_type: "SIM".into(),
2562 weight_prop: Some("score".into()),
2563 max_edges: None,
2564 approximate: false,
2565 via_label: None,
2566 via_edge: None,
2567 via_dir: None,
2568 },
2569 &mut g,
2570 )
2571 .unwrap();
2572 assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2574 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2575 assert!(eng.is_owned(et, a, b) || eng.is_owned(et, b, a));
2576 let check = |v: Option<&Value>| {
2577 if let Some(Value::Float(f)) = v {
2578 assert!(
2579 (f - 0.5).abs() < 1e-9,
2580 "initial score should be 0.5, got {f}"
2581 );
2582 }
2583 };
2584 check(g.edge_props.get(et, a, b, "score"));
2585 check(g.edge_props.get(et, b, a, "score"));
2586 }
2587 let old = fx.props.get(b, "tags").cloned();
2589 fx.props.set(b, "tags", tags(&["x", "y", "z"]));
2590 {
2591 let mut g = fx.g();
2592 eng.on_node_changed(b, Some(("tags", old)), &mut g);
2593 assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2595 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2596 assert_eq!(
2598 g.edge_props.get(et, a, b, "score"),
2599 Some(&Value::Float(1.0)),
2600 "score on a→b must refresh to 1.0"
2601 );
2602 assert_eq!(
2603 g.edge_props.get(et, b, a, "score"),
2604 Some(&Value::Float(1.0)),
2605 "score on b→a must refresh to 1.0"
2606 );
2607 }
2608 }
2609
2610 #[test]
2611 fn dst_side_keymatch_links_when_c_node_inserted_after_t() {
2612 let mut fx = Fx::new();
2614 let t = fx.add("T", "t1", vec![("cid", Value::Str("c9".into()))]);
2616 let mut eng = RuleEngine::new();
2617 {
2618 let mut g = fx.g();
2619 eng.create_rule(
2620 RuleDef {
2621 name: "fk".into(),
2622 src_label: "T".into(),
2623 dst_label: "C".into(),
2624 predicate: Predicate::KeyMatch {
2625 field: "cid".into(),
2626 },
2627 edge_type: "AT".into(),
2628 weight_prop: None,
2629 max_edges: None,
2630 approximate: false,
2631 via_label: None,
2632 via_edge: None,
2633 via_dir: None,
2634 },
2635 &mut g,
2636 )
2637 .unwrap();
2638 let at = g.syms.intern("AT");
2640 assert_eq!(g.topo.edge_count(), 0, "no C node yet → no edge");
2641 let _ = at;
2643 }
2644 let c9 = fx.add("C", "c9", vec![]);
2646 {
2647 let mut g = fx.g();
2648 eng.on_node_changed(c9, None, &mut g);
2649 let at = g.syms.get("AT").unwrap();
2650 assert!(
2652 g.topo.neighbors(at, Direction::Out, t).contains(&c9),
2653 "T→C edge must appear when C node is inserted"
2654 );
2655 assert!(eng.is_owned(at, t, c9));
2656 }
2657 }
2658
2659 #[test]
2660 fn on_node_removed_retracts_both_sides_and_deindexes() {
2661 let mut fx = Fx::new();
2662 let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2663 let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2664 let et = fx.syms.intern("REL");
2665 let mut eng = RuleEngine::new();
2666 {
2667 let mut g = fx.g();
2668 eng.create_rule(overlap_rule(), &mut g).unwrap();
2669 assert!(g.topo.neighbors(et, Direction::Out, a).contains(&b));
2670 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&a));
2671 }
2672 {
2673 let mut g = fx.g();
2674 eng.on_node_removed(a, &mut g);
2675 assert!(!g.topo.neighbors(et, Direction::Out, a).contains(&b));
2676 assert!(!g.topo.neighbors(et, Direction::Out, b).contains(&a));
2677 assert_eq!(g.edge_props.get(et, a, b, "score"), None);
2678 assert_eq!(g.edge_props.get(et, b, a, "score"), None);
2679 assert!(!eng.is_owned(et, a, b));
2680 assert!(!eng.is_owned(et, b, a));
2681 }
2682 let c = fx.add("A", "c", vec![("tags", tags(&["x", "y"]))]);
2684 {
2685 let mut g = fx.g();
2686 eng.on_node_changed(c, None, &mut g);
2687 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&c));
2688 assert!(g.topo.neighbors(et, Direction::Out, c).contains(&b));
2689 assert!(!g.topo.neighbors(et, Direction::Out, c).contains(&a));
2690 assert!(!g.topo.neighbors(et, Direction::Out, a).contains(&c));
2691 }
2692 {
2694 let mut g = fx.g();
2695 eng.on_node_removed(a, &mut g);
2696 assert!(g.topo.neighbors(et, Direction::Out, b).contains(&c));
2697 }
2698 }
2699
2700 #[test]
2701 fn duplicate_name_and_unknown_delete_error() {
2702 let mut fx = Fx::new();
2703 let mut eng = RuleEngine::new();
2704 let mut g = fx.g();
2705 eng.create_rule(overlap_rule(), &mut g).unwrap();
2706 assert!(eng.create_rule(overlap_rule(), &mut g).is_err());
2707 assert!(eng.delete_rule("nope", &mut g).is_err());
2708 }
2709
2710 #[test]
2716 fn coowned_edge_type_survives_first_delete_gone_after_second() {
2717 let mut fx = Fx::new();
2718 let a = fx.add("A", "a", vec![("tags", tags(&["x", "y"]))]);
2719 let b = fx.add("A", "b", vec![("tags", tags(&["x", "y"]))]);
2720 let mut eng = RuleEngine::new();
2721 {
2722 let mut g = fx.g();
2723 eng.create_rule(
2725 RuleDef {
2726 name: "r1".into(),
2727 src_label: "A".into(),
2728 dst_label: "A".into(),
2729 predicate: Predicate::Overlap {
2730 field: "tags".into(),
2731 min: 0.1,
2732 },
2733 edge_type: "REL2".into(),
2734 weight_prop: None,
2735 max_edges: None,
2736 approximate: false,
2737 via_label: None,
2738 via_edge: None,
2739 via_dir: None,
2740 },
2741 &mut g,
2742 )
2743 .unwrap();
2744 eng.create_rule(
2746 RuleDef {
2747 name: "r2".into(),
2748 src_label: "A".into(),
2749 dst_label: "A".into(),
2750 predicate: Predicate::Overlap {
2751 field: "tags".into(),
2752 min: 0.2,
2753 },
2754 edge_type: "REL2".into(),
2755 weight_prop: None,
2756 max_edges: None,
2757 approximate: false,
2758 via_label: None,
2759 via_edge: None,
2760 via_dir: None,
2761 },
2762 &mut g,
2763 )
2764 .unwrap();
2765
2766 let et = g.syms.intern("REL2");
2767 assert!(
2769 g.topo.neighbors(et, Direction::Out, a).contains(&b),
2770 "a→b must exist after both rules created"
2771 );
2772 assert!(
2773 g.topo.neighbors(et, Direction::Out, b).contains(&a),
2774 "b→a must exist after both rules created"
2775 );
2776
2777 eng.delete_rule("r1", &mut g).unwrap();
2779 assert!(
2780 g.topo.neighbors(et, Direction::Out, a).contains(&b),
2781 "a→b must survive R1 deletion (R2 rebuilds and claims it)"
2782 );
2783 assert!(
2784 g.topo.neighbors(et, Direction::Out, b).contains(&a),
2785 "b→a must survive R1 deletion (R2 rebuilds and claims it)"
2786 );
2787 assert!(
2789 eng.is_owned(et, a, b),
2790 "a→b must be owned by R2 after rebuild"
2791 );
2792 assert!(
2793 eng.is_owned(et, b, a),
2794 "b→a must be owned by R2 after rebuild"
2795 );
2796
2797 eng.delete_rule("r2", &mut g).unwrap();
2799 assert!(
2800 !g.topo.neighbors(et, Direction::Out, a).contains(&b),
2801 "a→b must be gone after both rules deleted"
2802 );
2803 assert!(
2804 !g.topo.neighbors(et, Direction::Out, b).contains(&a),
2805 "b→a must be gone after both rules deleted"
2806 );
2807 }
2808 }
2809
2810 fn topk_eq_rule(k: u64) -> RuleDef {
2812 RuleDef {
2813 name: "eq".into(),
2814 src_label: "N".into(),
2815 dst_label: "N".into(),
2816 predicate: Predicate::FieldEqual { field: "k".into() },
2817 edge_type: "EQ".into(),
2818 weight_prop: None,
2819 max_edges: Some(k),
2820 approximate: false,
2821 via_label: None,
2822 via_edge: None,
2823 via_dir: None,
2824 }
2825 }
2826
2827 fn prov_pairs(eng: &RuleEngine, name: &str) -> BTreeSet<(u32, u32)> {
2828 eng.provenance()
2829 .get(name)
2830 .map(|s| s.iter().map(|&(_, a, b)| (a, b)).collect())
2831 .unwrap_or_default()
2832 }
2833
2834 #[test]
2838 fn topk_k1_keeps_best_scored_dst() {
2839 let mut fx = Fx::new();
2840 let mut eng = RuleEngine::new();
2841 {
2842 let mut g = fx.g();
2843 eng.create_rule(topk_eq_rule(1), &mut g).unwrap();
2844 }
2845 let mut ids = Vec::new();
2847 for i in 0..4usize {
2848 let id = fx.add(
2849 "N",
2850 &format!("n{i}"),
2851 vec![("k", Value::Str("const".into()))],
2852 );
2853 ids.push(id);
2854 let mut g = fx.g();
2855 eng.on_node_changed(id, None, &mut g);
2856 }
2857 let et = fx.syms.get("EQ").unwrap();
2858 let expected_dsts = [ids[1], ids[0], ids[0], ids[0]];
2864 for (i, (&src, &expected_dst)) in ids.iter().zip(expected_dsts.iter()).enumerate() {
2865 let out: Vec<u32> = fx.topo.neighbors(et, Direction::Out, src).to_vec();
2866 assert_eq!(
2867 out,
2868 vec![expected_dst],
2869 "src n{i} should point only to the best dst"
2870 );
2871 }
2872 assert_eq!(eng.provenance()["eq"].len(), 4);
2873 assert!(!eng.is_tripped("eq"), "top-k rules never trip");
2874 }
2875
2876 #[test]
2879 fn topk_insert_evict() {
2880 let mut fx = Fx::new();
2884 let rule = RuleDef {
2885 name: "nw".into(),
2886 src_label: "S".into(),
2887 dst_label: "D".into(),
2888 predicate: Predicate::NumericWithin {
2889 field: "v".into(),
2890 tolerance: 10.0,
2891 },
2892 edge_type: "NEAR".into(),
2893 weight_prop: Some("score".into()),
2894 max_edges: Some(1),
2895 approximate: false,
2896 via_label: None,
2897 via_edge: None,
2898 via_dir: None,
2899 };
2900 let mut eng = RuleEngine::new();
2901 {
2902 let mut g = fx.g();
2903 eng.create_rule(rule, &mut g).unwrap();
2904 }
2905
2906 let s0 = fx.add("S", "s0", vec![("v", Value::Float(0.0))]);
2908 let d_far = fx.add("D", "d_far", vec![("v", Value::Float(9.0))]);
2910 {
2911 let mut g = fx.g();
2912 eng.on_node_changed(s0, None, &mut g);
2913 eng.on_node_changed(d_far, None, &mut g);
2914 }
2915 let et = fx.syms.get("NEAR").unwrap();
2916 assert!(fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2918 assert_eq!(eng.provenance()["nw"].len(), 1);
2919
2920 let d_close = fx.add("D", "d_close", vec![("v", Value::Float(1.0))]);
2922 {
2923 let mut g = fx.g();
2924 eng.on_node_changed(d_close, None, &mut g);
2925 }
2926 let out: Vec<u32> = fx.topo.neighbors(et, Direction::Out, s0).to_vec();
2928 assert_eq!(out, vec![d_close], "d_close should evict d_far");
2929 assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2930 assert_eq!(eng.provenance()["nw"].len(), 1);
2931 assert!(eng.by_node_consistent());
2932 }
2933
2934 #[test]
2936 fn topk_retract_backfill() {
2937 let mut fx = Fx::new();
2938 let rule = RuleDef {
2939 name: "nw".into(),
2940 src_label: "S".into(),
2941 dst_label: "D".into(),
2942 predicate: Predicate::NumericWithin {
2943 field: "v".into(),
2944 tolerance: 10.0,
2945 },
2946 edge_type: "NEAR".into(),
2947 weight_prop: Some("score".into()),
2948 max_edges: Some(1),
2949 approximate: false,
2950 via_label: None,
2951 via_edge: None,
2952 via_dir: None,
2953 };
2954 let mut eng = RuleEngine::new();
2955
2956 let s0 = fx.add("S", "s0", vec![("v", Value::Float(0.0))]);
2957 let d_close = fx.add("D", "d_close", vec![("v", Value::Float(1.0))]); let d_far = fx.add("D", "d_far", vec![("v", Value::Float(8.0))]); {
2960 let mut g = fx.g();
2961 eng.create_rule(rule, &mut g).unwrap();
2962 }
2963 let et = fx.syms.get("NEAR").unwrap();
2964 assert!(fx.topo.neighbors(et, Direction::Out, s0).contains(&d_close));
2966 assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far));
2967 assert_eq!(eng.provenance()["nw"].len(), 1);
2968
2969 let old = fx.props.get(d_close, "v").cloned();
2971 fx.props.set(d_close, "v", Value::Float(50.0));
2972 {
2973 let mut g = fx.g();
2974 eng.on_node_changed(d_close, Some(("v", old)), &mut g);
2975 }
2976 assert!(!fx.topo.neighbors(et, Direction::Out, s0).contains(&d_close));
2978 assert!(
2979 fx.topo.neighbors(et, Direction::Out, s0).contains(&d_far),
2980 "d_far should backfill after d_close retracted"
2981 );
2982 assert_eq!(eng.provenance()["nw"].len(), 1);
2983 assert!(eng.by_node_consistent());
2984 }
2985
2986 #[test]
2988 fn topk_tie_broken_by_dst_key() {
2989 let mut fx = Fx::new();
2991 let mut eng = RuleEngine::new();
2992 {
2993 let mut g = fx.g();
2994 eng.create_rule(topk_eq_rule(2), &mut g).unwrap();
2995 }
2996 for name in ["a", "b", "c", "d", "e"] {
2999 let id = fx.add("N", name, vec![("k", Value::Str("x".into()))]);
3000 let mut g = fx.g();
3001 eng.on_node_changed(id, None, &mut g);
3002 }
3003 let et = fx.syms.get("EQ").unwrap();
3004 let get_id = |key: &str| fx.ids.get(key).unwrap();
3005 let a = get_id("a");
3007 let b = get_id("b");
3008 let c = get_id("c");
3009 let out_a: BTreeSet<u32> = fx
3010 .topo
3011 .neighbors(et, Direction::Out, a)
3012 .iter()
3013 .copied()
3014 .collect();
3015 assert!(out_a.contains(&b), "a→b (b is best key after a)");
3016 assert!(out_a.contains(&c), "a→c (c is 2nd best key)");
3017 assert_eq!(out_a.len(), 2);
3018 let e = get_id("e");
3020 let out_e: BTreeSet<u32> = fx
3021 .topo
3022 .neighbors(et, Direction::Out, e)
3023 .iter()
3024 .copied()
3025 .collect();
3026 assert!(out_e.contains(&a), "e→a");
3027 assert!(out_e.contains(&b), "e→b");
3028 assert_eq!(out_e.len(), 2);
3029 assert!(eng.by_node_consistent());
3030 }
3031
3032 #[test]
3034 fn topk_k_larger_than_candidate_count() {
3035 let mut fx = Fx::new();
3036 let mut eng = RuleEngine::new();
3037 {
3038 let mut g = fx.g();
3039 eng.create_rule(topk_eq_rule(100), &mut g).unwrap();
3041 }
3042 for i in 0..4usize {
3043 let id = fx.add("N", &format!("n{i}"), vec![("k", Value::Str("c".into()))]);
3044 let mut g = fx.g();
3045 eng.on_node_changed(id, None, &mut g);
3046 }
3047 assert_eq!(eng.provenance()["eq"].len(), 12);
3049 assert!(!eng.is_tripped("eq"));
3050 }
3051
3052 #[test]
3055 fn topk_rebuild_exact() {
3056 let mut fx = Fx::new();
3057 let mut eng = RuleEngine::new();
3058 {
3059 let mut g = fx.g();
3060 eng.create_rule(topk_eq_rule(1), &mut g).unwrap();
3061 }
3062 let _a = fx.add("N", "a", vec![("k", Value::Str("x".into()))]);
3064 let _b = fx.add("N", "b", vec![("k", Value::Str("x".into()))]);
3065 let _c = fx.add("N", "c", vec![("k", Value::Str("x".into()))]);
3066 {
3067 let mut g = fx.g();
3068 eng.on_node_changed(_a, None, &mut g);
3069 eng.on_node_changed(_b, None, &mut g);
3070 eng.on_node_changed(_c, None, &mut g);
3071 }
3072 assert_eq!(eng.provenance()["eq"].len(), 3);
3073
3074 {
3076 let mut g = fx.g();
3077 eng.rebuild("eq", &mut g).unwrap();
3078 }
3079 assert_eq!(eng.provenance()["eq"].len(), 3);
3080 assert!(!eng.is_tripped("eq"));
3081 assert!(eng.by_node_consistent());
3082 }
3083
3084 #[test]
3086 fn topk_by_node_consistent() {
3087 let mut fx = Fx::new();
3088 let mut eng = RuleEngine::new();
3089 {
3090 let mut g = fx.g();
3091 eng.create_rule(topk_eq_rule(2), &mut g).unwrap();
3092 }
3093 for i in 0..5usize {
3094 let id = fx.add(
3095 "N",
3096 &format!("n{i}"),
3097 vec![("k", Value::Str("const".into()))],
3098 );
3099 let mut g = fx.g();
3100 eng.on_node_changed(id, None, &mut g);
3101 }
3102 assert!(eng.by_node_consistent(), "consistent after insertions");
3103
3104 let id2 = fx.ids.get("n2").unwrap();
3106 let old = fx.props.get(id2, "k").cloned();
3107 fx.props.set(id2, "k", Value::Str("other".into()));
3108 {
3109 let mut g = fx.g();
3110 eng.on_node_changed(id2, Some(("k", old)), &mut g);
3111 }
3112 assert!(eng.by_node_consistent(), "consistent after eviction");
3113
3114 {
3115 let mut g = fx.g();
3116 eng.rebuild("eq", &mut g).unwrap();
3117 }
3118 assert!(eng.by_node_consistent(), "consistent after rebuild");
3119 }
3120
3121 fn numeric_rule() -> RuleDef {
3122 RuleDef {
3123 name: "nw".into(),
3124 src_label: "C".into(),
3125 dst_label: "C".into(),
3126 predicate: Predicate::NumericWithin {
3127 field: "year".into(),
3128 tolerance: 2.0,
3129 },
3130 edge_type: "NEAR".into(),
3131 weight_prop: Some("score".into()),
3132 max_edges: None,
3133 approximate: false,
3134 via_label: None,
3135 via_edge: None,
3136 via_dir: None,
3137 }
3138 }
3139
3140 fn geo_rule() -> RuleDef {
3141 RuleDef {
3142 name: "geo".into(),
3143 src_label: "City".into(),
3144 dst_label: "City".into(),
3145 predicate: Predicate::GeoRadius {
3146 field: "loc".into(),
3147 km: 400.0,
3148 },
3149 edge_type: "NEAR_GEO".into(),
3150 weight_prop: Some("score".into()),
3151 max_edges: None,
3152 approximate: false,
3153 via_label: None,
3154 via_edge: None,
3155 via_dir: None,
3156 }
3157 }
3158
3159 fn vec_rule() -> RuleDef {
3160 RuleDef {
3161 name: "vec".into(),
3162 src_label: "Doc".into(),
3163 dst_label: "Doc".into(),
3164 predicate: Predicate::VectorSimilar {
3165 field: "emb".into(),
3166 min: 0.9,
3167 },
3168 edge_type: "SIM".into(),
3169 weight_prop: Some("score".into()),
3170 max_edges: None,
3171 approximate: false,
3172 via_label: None,
3173 via_edge: None,
3174 via_dir: None,
3175 }
3176 }
3177
3178 fn pair_edges(topo: &Topology, et: u32, a: u32, b: u32) -> bool {
3179 topo.neighbors(et, Direction::Out, a).contains(&b)
3180 && topo.neighbors(et, Direction::Out, b).contains(&a)
3181 }
3182
3183 #[test]
3184 fn numeric_within_incremental_crosses_bucket_and_clears_old_index() {
3185 let mut fx = Fx::new();
3186 let a = fx.add("C", "a", vec![("year", Value::Float(10.0))]);
3187 let b = fx.add("C", "b", vec![("year", Value::Float(12.0))]);
3188 let et = fx.syms.intern("NEAR");
3189 let mut eng = RuleEngine::new();
3190 {
3191 let mut g = fx.g();
3192 eng.create_rule(numeric_rule(), &mut g).unwrap();
3193 assert!(pair_edges(g.topo, et, a, b));
3195 }
3196
3197 let old = fx.props.get(b, "year").cloned();
3200 fx.props.set(b, "year", Value::Float(16.1));
3201 {
3202 let mut g = fx.g();
3203 eng.on_node_changed(b, Some(("year", old)), &mut g);
3204 assert!(!pair_edges(g.topo, et, a, b));
3205 assert_eq!(g.topo.edge_count(), 0);
3206 }
3207 let def = numeric_rule();
3208 let spec = candidate_spec_for(&def);
3209 let old_map: std::collections::HashMap<_, _> =
3210 [("year".to_string(), Value::Float(12.0))].into();
3211 let old_get = |f: &str| old_map.get(f).cloned();
3212 let src_hits = eng.indexes["nw"].src_side.candidates(&spec, &old_get);
3213 let dst_hits = eng.indexes["nw"].dst_side.candidates(&spec, &old_get);
3214 assert!(!src_hits.contains(&b), "old src bucket must drop b");
3215 assert!(!dst_hits.contains(&b), "old dst bucket must drop b");
3216 assert!(src_hits.contains(&a));
3217
3218 let old = fx.props.get(b, "year").cloned();
3220 fx.props.set(b, "year", Value::Float(11.9));
3221 let mut g = fx.g();
3222 eng.on_node_changed(b, Some(("year", old)), &mut g);
3223 assert!(pair_edges(g.topo, et, a, b));
3224 }
3225
3226 fn loc_val(lat: f64, lon: f64) -> Value {
3227 Value::List(vec![Value::Float(lat), Value::Float(lon)])
3228 }
3229
3230 fn emb_val(vals: &[f64]) -> Value {
3231 Value::List(vals.iter().copied().map(Value::Float).collect())
3232 }
3233
3234 #[test]
3235 fn rebuild_is_noop_for_numeric_geo_and_vector() {
3236 let mut fx = Fx::new();
3237 let ca = fx.add("C", "ca", vec![("year", Value::Int(1998))]);
3238 let cb = fx.add("C", "cb", vec![("year", Value::Float(2000.0))]);
3239 let pa = fx.add("City", "paris", vec![("loc", loc_val(48.8566, 2.3522))]);
3240 let lo = fx.add("City", "london", vec![("loc", loc_val(51.5074, -0.1278))]);
3241 let da = fx.add("Doc", "d1", vec![("emb", emb_val(&[1.0, 0.0]))]);
3242 let db = fx.add("Doc", "d2", vec![("emb", emb_val(&[1.0, 0.0]))]);
3243
3244 let mut eng = RuleEngine::new();
3245 {
3246 let mut g = fx.g();
3247 eng.create_rule(numeric_rule(), &mut g).unwrap();
3248 eng.create_rule(geo_rule(), &mut g).unwrap();
3249 eng.create_rule(vec_rule(), &mut g).unwrap();
3250 }
3251
3252 let (near, ngeo, sim) = (
3253 fx.syms.get("NEAR").unwrap(),
3254 fx.syms.get("NEAR_GEO").unwrap(),
3255 fx.syms.get("SIM").unwrap(),
3256 );
3257 assert!(pair_edges(&fx.topo, near, ca, cb));
3258 assert!(pair_edges(&fx.topo, ngeo, pa, lo));
3259 assert!(pair_edges(&fx.topo, sim, da, db));
3260 let before = fx.topo.edge_count();
3261
3262 {
3263 let mut g = fx.g();
3264 eng.rebuild("nw", &mut g).unwrap();
3265 eng.rebuild("geo", &mut g).unwrap();
3266 eng.rebuild("vec", &mut g).unwrap();
3267 }
3268 assert_eq!(fx.topo.edge_count(), before);
3269 assert!(pair_edges(&fx.topo, near, ca, cb));
3270 assert!(pair_edges(&fx.topo, ngeo, pa, lo));
3271 assert!(pair_edges(&fx.topo, sim, da, db));
3272 }
3273
3274 fn fk_rule() -> RuleDef {
3275 RuleDef {
3276 name: "works_at".into(),
3277 src_label: "T".into(),
3278 dst_label: "C".into(),
3279 predicate: Predicate::KeyMatch {
3280 field: "cid".into(),
3281 },
3282 edge_type: "AT".into(),
3283 weight_prop: None,
3284 max_edges: None,
3285 approximate: false,
3286 via_label: None,
3287 via_edge: None,
3288 via_dir: None,
3289 }
3290 }
3291
3292 #[test]
3293 fn by_node_matches_rebuild_after_mutation_storm() {
3294 let mut fx = Fx::new();
3295 let hub = fx.add("C", "hub", vec![]);
3296 let other = fx.add("C", "other", vec![]);
3297 let mut people = Vec::new();
3298 for i in 0..40 {
3299 let cid = if i < 30 { "hub" } else { "other" };
3300 people.push(fx.add(
3301 "T",
3302 &format!("t{i}"),
3303 vec![("cid", Value::Str(cid.into())), ("tags", tags(&["x", "y"]))],
3304 ));
3305 }
3306 let mut overlap = overlap_rule();
3307 overlap.src_label = "T".into();
3308 overlap.dst_label = "T".into();
3309 let mut eng = RuleEngine::new();
3310 {
3311 let mut g = fx.g();
3312 eng.create_rule(fk_rule(), &mut g).unwrap();
3313 eng.create_rule(overlap, &mut g).unwrap();
3314 }
3315 assert!(eng.by_node_consistent());
3316 assert_eq!(eng.provenance_touching_len(hub), 30);
3317
3318 for (i, &id) in people.iter().enumerate().take(15) {
3320 let old = fx.props.get(id, "cid").cloned();
3321 fx.props.set(id, "cid", Value::Str("other".into()));
3322 let mut g = fx.g();
3323 eng.on_node_changed(id, Some(("cid", old)), &mut g);
3324 assert!(
3325 eng.by_node_consistent(),
3326 "inconsistent after cid update {i}"
3327 );
3328 }
3329 for &id in people.iter().take(8) {
3330 let old = fx.props.get(id, "tags").cloned();
3331 fx.props.set(id, "tags", tags(&["q"]));
3332 let mut g = fx.g();
3333 eng.on_node_changed(id, Some(("tags", old)), &mut g);
3334 }
3335 assert!(eng.by_node_consistent());
3336
3337 {
3339 let mut g = fx.g();
3340 eng.on_node_removed(people[0], &mut g);
3341 }
3342 fx.labels[people[0] as usize] = u32::MAX;
3343 assert!(eng.by_node_consistent());
3344 assert_eq!(eng.provenance_touching_len(people[0]), 0);
3345
3346 {
3347 let mut g = fx.g();
3348 eng.rebuild("works_at", &mut g).unwrap();
3349 eng.rebuild("rel", &mut g).unwrap();
3350 }
3351 assert!(eng.by_node_consistent());
3352
3353 {
3354 let mut g = fx.g();
3355 eng.delete_rule("rel", &mut g).unwrap();
3356 }
3357 assert!(eng.by_node_consistent());
3358 assert_eq!(eng.provenance_touching(people[1]).count(), 1);
3359
3360 let (defs, prov, tripped, fires) = eng.to_persist();
3362 let restored = RuleEngine::from_persist(defs, prov, tripped, fires);
3363 assert!(restored.by_node_consistent());
3364 assert_eq!(
3365 restored.provenance_touching_len(hub),
3366 eng.provenance_touching_len(hub)
3367 );
3368 assert_eq!(
3369 restored.provenance_touching_len(other),
3370 eng.provenance_touching_len(other)
3371 );
3372 }
3373
3374 #[test]
3375 fn provenance_touching_high_degree_hub() {
3376 let mut fx = Fx::new();
3377 let hub = fx.add("C", "hub", vec![]);
3378 let mut first = None;
3379 for i in 0..256 {
3380 let id = fx.add(
3381 "T",
3382 &format!("t{i}"),
3383 vec![("cid", Value::Str("hub".into()))],
3384 );
3385 if first.is_none() {
3386 first = Some(id);
3387 }
3388 }
3389 let first = first.unwrap();
3390 let mut eng = RuleEngine::new();
3391 {
3392 let mut g = fx.g();
3393 eng.create_rule(fk_rule(), &mut g).unwrap();
3394 }
3395 assert!(eng.by_node_consistent());
3396 assert_eq!(eng.provenance_touching_len(hub), 256);
3397 assert_eq!(eng.provenance_touching_len(first), 1);
3398 let hits: Vec<_> = eng.provenance_touching(first).collect();
3399 assert_eq!(hits.len(), 1);
3400 assert_eq!(hits[0].0, "works_at");
3401 assert_eq!(hits[0].2, first);
3402 assert_eq!(hits[0].3, hub);
3403 }
3404
3405 #[test]
3413 fn by_node_consistent_across_inserts_and_rebuild() {
3414 let mut fx = Fx::new();
3415 let mut eng = RuleEngine::new();
3416 let rule = RuleDef {
3417 name: "eq".into(),
3418 src_label: "N".into(),
3419 dst_label: "N".into(),
3420 predicate: Predicate::FieldEqual { field: "k".into() },
3421 edge_type: "EQ".into(),
3422 weight_prop: None,
3423 max_edges: None, approximate: false,
3425 via_label: None,
3426 via_edge: None,
3427 via_dir: None,
3428 };
3429 {
3430 let mut g = fx.g();
3431 eng.create_rule(rule, &mut g).unwrap();
3432 }
3433 let mut ids = Vec::new();
3434 for i in 0..6 {
3435 let id = fx.add(
3436 "N",
3437 &format!("n{i}"),
3438 vec![("k", Value::Str("const".into()))],
3439 );
3440 ids.push(id);
3441 let mut g = fx.g();
3442 eng.on_node_changed(id, None, &mut g);
3443 }
3444 assert_eq!(eng.provenance()["eq"].len(), 30);
3446 assert!(!eng.is_tripped("eq"));
3447 assert!(eng.by_node_consistent(), "consistent after insertions");
3448
3449 let old = fx.props.get(ids[3], "k").cloned();
3451 fx.props.set(ids[3], "k", Value::Str("other".into()));
3452 {
3453 let mut g = fx.g();
3454 eng.on_node_changed(ids[3], Some(("k", old)), &mut g);
3455 }
3456 assert!(eng.by_node_consistent(), "consistent after property change");
3457
3458 {
3459 let mut g = fx.g();
3460 eng.rebuild("eq", &mut g).unwrap();
3461 }
3462 assert!(!eng.is_tripped("eq"));
3463 assert!(eng.by_node_consistent(), "consistent after rebuild");
3464 }
3465
3466 fn mix64(mut x: u64) -> u64 {
3467 x = x.wrapping_add(0x9E3779B97F4A7C15);
3468 x = (x ^ (x >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
3469 x = (x ^ (x >> 27)).wrapping_mul(0x94D049BB133111EB);
3470 x ^ (x >> 31)
3471 }
3472
3473 fn rand_emb(seed: u64, i: u32, dim: usize) -> Value {
3474 let vals: Vec<f64> = (0..dim)
3475 .map(|d| {
3476 let bits = mix64(seed ^ ((i as u64 + 1).wrapping_mul(0x100000001)) ^ (d as u64));
3477 let mut f = (bits as f64) / (u64::MAX as f64) * 2.0 - 1.0;
3478 if f == 0.0 {
3479 f = 1.0;
3480 }
3481 f
3482 })
3483 .collect();
3484 emb_val(&vals)
3485 }
3486
3487 fn seed_docs(n: u32, seed: u64) -> (Fx, Vec<u32>) {
3488 let dims = [2usize, 3, 4, 8];
3489 let mut fx = Fx::new();
3490 let mut ids = Vec::new();
3491 for i in 0..n {
3492 let dim = dims[(i as usize) % dims.len()];
3493 ids.push(fx.add(
3494 "Doc",
3495 &format!("d{i}"),
3496 vec![("emb", rand_emb(seed, i, dim))],
3497 ));
3498 }
3499 (fx, ids)
3500 }
3501
3502 #[test]
3505 fn vector_dim_reject_matches_unfiltered_and_oracle() {
3506 const N: u32 = 500;
3507 const SEED: u64 = 0xC0FF_EE00_D15C;
3508 let def = vec_rule();
3509
3510 let (mut fx_on, ids) = seed_docs(N, SEED);
3511 let mut eng_on = RuleEngine::new();
3512 {
3513 let mut g = fx_on.g();
3514 eng_on.create_rule(def.clone(), &mut g).unwrap();
3515 }
3516 let on = prov_pairs(&eng_on, "vec");
3517 assert!(!on.is_empty(), "seeded set must produce some edges");
3518
3519 let (mut fx_off, _) = seed_docs(N, SEED);
3520 let mut eng_off = RuleEngine::new();
3521 {
3522 let mut g = fx_off.g();
3523 with_vector_dim_reject(false, || {
3524 eng_off.create_rule(def.clone(), &mut g).unwrap();
3525 });
3526 }
3527 assert_eq!(on, prov_pairs(&eng_off, "vec"), "filter vs no-filter");
3528
3529 let mut brute = BTreeSet::new();
3530 for &s in &ids {
3531 for &d in &ids {
3532 if s == d {
3533 continue;
3534 }
3535 let skey = fx_on.ids.key_of(s).unwrap();
3536 let dkey = fx_on.ids.key_of(d).unwrap();
3537 let sget = |f: &str| fx_on.props.get(s, f).cloned();
3538 let dget = |f: &str| fx_on.props.get(d, f).cloned();
3539 if evaluate(
3540 &def.predicate,
3541 &NodeView {
3542 key: skey,
3543 props: &sget,
3544 },
3545 &NodeView {
3546 key: dkey,
3547 props: &dget,
3548 },
3549 )
3550 .is_some()
3551 {
3552 brute.insert((s, d));
3553 }
3554 }
3555 }
3556 assert_eq!(on, brute, "filter vs brute-force evaluate");
3557 }
3558
3559 #[test]
3562 fn vector_dim_change_updates_cache_and_matches_fresh_build() {
3563 let mut fx = Fx::new();
3564 let a = fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0]))]);
3565 let b = fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0]))]);
3566 let c = fx.add("Doc", "c", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3567 let mut eng = RuleEngine::new();
3568 {
3569 let mut g = fx.g();
3570 eng.create_rule(vec_rule(), &mut g).unwrap();
3571 }
3572 assert_eq!(eng.indexes["vec"].src_side.vec_dim(a), Some(2));
3573 assert_eq!(eng.indexes["vec"].src_side.vec_dim(c), Some(3));
3574 assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(a, b), (b, a)]));
3575
3576 let old = fx.props.get(b, "emb").cloned();
3577 fx.props.set(b, "emb", emb_val(&[1.0, 0.0, 0.0]));
3578 {
3579 let mut g = fx.g();
3580 eng.on_node_changed(b, Some(("emb", old)), &mut g);
3581 }
3582 assert_eq!(eng.indexes["vec"].src_side.vec_dim(b), Some(3));
3583 assert_eq!(eng.indexes["vec"].dst_side.vec_dim(b), Some(3));
3584 let after = prov_pairs(&eng, "vec");
3585 assert_eq!(after, BTreeSet::from([(b, c), (c, b)]));
3586
3587 let mut fresh_fx = Fx::new();
3589 let fa = fresh_fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0]))]);
3590 let fb = fresh_fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3591 let fc = fresh_fx.add("Doc", "c", vec![("emb", emb_val(&[1.0, 0.0, 0.0]))]);
3592 let mut fresh = RuleEngine::new();
3593 {
3594 let mut g = fresh_fx.g();
3595 fresh.create_rule(vec_rule(), &mut g).unwrap();
3596 }
3597 assert_eq!(
3598 prov_pairs(&fresh, "vec"),
3599 BTreeSet::from([(fb, fc), (fc, fb)])
3600 );
3601 assert_eq!(fresh.indexes["vec"].src_side.vec_dim(fb), Some(3));
3602 assert_eq!(fresh.indexes["vec"].src_side.vec_dim(fa), Some(2));
3603 }
3604
3605 #[test]
3625 fn streaming_topk_order_identity_property_test() {
3626 fn reference_topk(rule: &RuleDef, k: u64, fx: &mut Fx) -> BTreeSet<(u32, u32)> {
3629 let mut idx = RuleIndex::default();
3630 for id in 0..fx.ids.len() as u32 {
3631 let label_sym = match fx.labels.get(id as usize).copied() {
3632 Some(s) if s != u32::MAX => s,
3633 _ => continue,
3634 };
3635 index_node_for_rule(id, label_sym, rule, &mut idx, &fx.syms, &fx.props);
3636 }
3637 let src_sym = fx.syms.get(&rule.src_label);
3638 let mut out = BTreeSet::new();
3639 let ids_snap: Vec<u32> = (0..fx.ids.len() as u32).collect();
3640 for id in ids_snap {
3641 let label_sym = match fx.labels.get(id as usize).copied() {
3642 Some(s) if s != u32::MAX => s,
3643 _ => continue,
3644 };
3645 if src_sym != Some(label_sym) {
3646 continue;
3647 }
3648 let g = GraphMut {
3649 ids: &fx.ids,
3650 syms: &mut fx.syms,
3651 labels: &fx.labels,
3652 props: &fx.props,
3653 topo: &mut fx.topo,
3654 edge_props: &mut fx.eprops,
3655 };
3656 let per_src = compute_desired(rule, &idx, id, true, &g);
3657 let mut candidates: Vec<((u32, u32), f64)> = per_src.into_iter().collect();
3659 candidates.sort_by(|&((_, da), sa), &((_, db), sb)| {
3660 sb.total_cmp(&sa).then_with(|| {
3661 let ka = fx.ids.key_of(da).unwrap_or("");
3662 let kb = fx.ids.key_of(db).unwrap_or("");
3663 ka.cmp(kb)
3664 })
3665 });
3666 candidates.truncate(k as usize);
3667 out.extend(candidates.into_iter().map(|(k, _)| k));
3668 }
3669 out
3670 }
3671
3672 fn streaming_pairs(rule: RuleDef, fx: &mut Fx) -> BTreeSet<(u32, u32)> {
3674 let name = rule.name.clone();
3675 let mut eng = RuleEngine::new();
3676 eng.create_rule(rule, &mut fx.g()).unwrap();
3677 eng.provenance()
3678 .get(&name)
3679 .map(|s| s.iter().map(|&(_, a, b)| (a, b)).collect())
3680 .unwrap_or_default()
3681 }
3682
3683 for seed in [0u64, 1, 42, 0xDEAD_BEEF, 0x1234_5678, 99, 12_648_430, 7] {
3688 for k in [1u64, 2, 3, 5] {
3689 let rule = RuleDef {
3690 name: "eq".into(),
3691 src_label: "N".into(),
3692 dst_label: "N".into(),
3693 predicate: Predicate::FieldEqual { field: "k".into() },
3694 edge_type: "EQ".into(),
3695 weight_prop: None,
3696 max_edges: Some(k),
3697 approximate: false,
3698 via_label: None,
3699 via_edge: None,
3700 via_dir: None,
3701 };
3702
3703 let build = || {
3704 let mut fx = Fx::new();
3705 for i in 0..12u32 {
3706 let h = mix64(seed ^ (i as u64 + 1));
3707 let val = match h % 3 {
3708 0 => "a",
3709 1 => "b",
3710 _ => "c",
3711 };
3712 fx.add(
3713 "N",
3714 &format!("n{i:02}"),
3715 vec![("k", Value::Str(val.into()))],
3716 );
3717 }
3718 fx
3719 };
3720
3721 let expected = reference_topk(&rule, k, &mut build());
3722 let actual = streaming_pairs(rule, &mut build());
3723
3724 assert_eq!(
3725 expected, actual,
3726 "FieldEqual seed={seed} k={k}: streaming top-k must match brute-force top-k"
3727 );
3728 }
3729 }
3730
3731 for seed in [0u64, 1, 42, 7] {
3736 for k in [1u64, 2, 4] {
3737 let rule = RuleDef {
3738 name: "nw".into(),
3739 src_label: "S".into(),
3740 dst_label: "D".into(),
3741 predicate: Predicate::NumericWithin {
3742 field: "v".into(),
3743 tolerance: 10.0,
3744 },
3745 edge_type: "NEAR".into(),
3746 weight_prop: Some("score".into()),
3747 max_edges: Some(k),
3748 approximate: false,
3749 via_label: None,
3750 via_edge: None,
3751 via_dir: None,
3752 };
3753
3754 let build = || {
3755 let mut fx = Fx::new();
3756 for i in 0..6u32 {
3757 let h = mix64(seed ^ (i as u64 + 1));
3758 let v = (h % 20) as f64;
3759 fx.add("S", &format!("s{i}"), vec![("v", Value::Float(v))]);
3760 }
3761 for i in 0..8u32 {
3762 let h = mix64(seed ^ (i as u64 + 101));
3763 let v = (h % 20) as f64;
3764 fx.add("D", &format!("d{i}"), vec![("v", Value::Float(v))]);
3765 }
3766 fx
3767 };
3768
3769 let expected = reference_topk(&rule, k, &mut build());
3770 let actual = streaming_pairs(rule, &mut build());
3771
3772 assert_eq!(
3773 expected, actual,
3774 "NumericWithin seed={seed} k={k}: streaming top-k must match brute-force top-k"
3775 );
3776 }
3777 }
3778
3779 for seed in [0u64, 1, 42, 7] {
3786 for k in [1u64, 2] {
3787 let rule = RuleDef {
3788 name: "fk".into(),
3789 src_label: "T".into(),
3790 dst_label: "C".into(),
3791 predicate: Predicate::KeyMatch {
3792 field: "cid".into(),
3793 },
3794 edge_type: "AT".into(),
3795 weight_prop: None,
3796 max_edges: Some(k),
3797 approximate: false,
3798 via_label: None,
3799 via_edge: None,
3800 via_dir: None,
3801 };
3802
3803 let build = || {
3804 let mut fx = Fx::new();
3805 for i in 0..4u32 {
3807 fx.add("C", &format!("c{i}"), vec![]);
3808 }
3809 for i in 0..8u32 {
3811 let h = mix64(seed ^ (i as u64 + 1));
3812 let cid = format!("c{}", h % 4);
3813 fx.add("T", &format!("t{i}"), vec![("cid", Value::Str(cid))]);
3814 }
3815 fx
3816 };
3817
3818 let expected = reference_topk(&rule, k, &mut build());
3819 let actual = streaming_pairs(rule, &mut build());
3820
3821 assert_eq!(
3822 expected, actual,
3823 "KeyMatch seed={seed} k={k}: streaming top-k must match brute-force top-k"
3824 );
3825 }
3826 }
3827
3828 {
3834 let cluster_a: &[(&str, f64, f64)] = &[
3836 ("va0", 1.0_f64, 0.0_f64),
3837 ("va1", 0.98_f64, 0.199_f64), ("va2", 0.97_f64, 0.243_f64), ];
3840 let cluster_b: &[(&str, f64, f64)] = &[
3841 ("vb0", 0.0_f64, 1.0_f64),
3842 ("vb1", 0.1_f64, 0.995_f64),
3843 ("vb2", 0.05_f64, 0.999_f64),
3844 ];
3845 for k in [1u64, 2] {
3846 let rule = RuleDef {
3847 name: "vsim".into(),
3848 src_label: "V".into(),
3849 dst_label: "V".into(),
3850 predicate: Predicate::VectorSimilar {
3851 field: "emb".into(),
3852 min: 0.9,
3853 },
3854 edge_type: "VSIM".into(),
3855 weight_prop: Some("score".into()),
3856 max_edges: Some(k),
3857 approximate: false,
3858 via_label: None,
3859 via_edge: None,
3860 via_dir: None,
3861 };
3862
3863 let build = || {
3864 let mut fx = Fx::new();
3865 let mut add_v = |key: &str, x: f64, y: f64| {
3866 let norm = (x * x + y * y).sqrt();
3867 let v = Value::List(vec![Value::Float(x / norm), Value::Float(y / norm)]);
3868 fx.add("V", key, vec![("emb", v)]);
3869 };
3870 for &(k, x, y) in cluster_a.iter().chain(cluster_b.iter()) {
3871 add_v(k, x, y);
3872 }
3873 fx
3874 };
3875
3876 let expected = reference_topk(&rule, k, &mut build());
3877 let actual = streaming_pairs(rule, &mut build());
3878
3879 assert_eq!(
3880 expected, actual,
3881 "VectorSimilar/ScanAll k={k}: streaming top-k must match brute-force top-k"
3882 );
3883 }
3884 }
3885 }
3886
3887 #[test]
3913 #[ignore]
3914 fn streaming_peak_transient_bound() {
3915 use std::sync::{
3916 atomic::{AtomicBool, AtomicU64, Ordering},
3917 Arc,
3918 };
3919
3920 fn peak_rss_during<F: FnOnce()>(f: F) -> u64 {
3923 let done = Arc::new(AtomicBool::new(false));
3924 let peak = Arc::new(AtomicU64::new(0));
3925 let done2 = done.clone();
3926 let peak2 = peak.clone();
3927 let pid = std::process::id().to_string();
3928
3929 let handle = std::thread::spawn(move || {
3930 while !done2.load(Ordering::Relaxed) {
3931 let rss = std::process::Command::new("ps")
3932 .args(["-o", "rss=", "-p", &pid])
3933 .output()
3934 .ok()
3935 .and_then(|o| String::from_utf8(o.stdout).ok())
3936 .and_then(|s| s.trim().parse::<u64>().ok())
3937 .unwrap_or(0)
3938 * 1024;
3939 peak2.fetch_max(rss, Ordering::Relaxed);
3940 std::thread::sleep(std::time::Duration::from_millis(1));
3941 }
3942 });
3943
3944 f();
3945
3946 done.store(true, Ordering::Relaxed);
3947 let _ = handle.join();
3948 peak.load(Ordering::Relaxed)
3949 }
3950
3951 let mut fx = Fx::new();
3955 for i in 0..500u32 {
3956 fx.add(
3957 "Talent",
3958 &format!("t{i}"),
3959 vec![("k", Value::Str("same".into()))],
3960 );
3961 }
3962 for i in 0..500u32 {
3963 fx.add(
3964 "Company",
3965 &format!("c{i}"),
3966 vec![("k", Value::Str("same".into()))],
3967 );
3968 }
3969 let rule = RuleDef {
3970 name: "eq_tc".into(),
3971 src_label: "Talent".into(),
3972 dst_label: "Company".into(),
3973 predicate: Predicate::FieldEqual { field: "k".into() },
3974 edge_type: "EQ".into(),
3975 weight_prop: None,
3976 max_edges: Some(2), approximate: false,
3978 via_label: None,
3979 via_edge: None,
3980 via_dir: None,
3981 };
3982
3983 let pid = std::process::id().to_string();
3985 let baseline = std::process::Command::new("ps")
3986 .args(["-o", "rss=", "-p", &pid])
3987 .output()
3988 .ok()
3989 .and_then(|o| String::from_utf8(o.stdout).ok())
3990 .and_then(|s| s.trim().parse::<u64>().ok())
3991 .unwrap_or(0)
3992 * 1024;
3993
3994 let mut eng = RuleEngine::new();
3995 let peak = peak_rss_during(|| {
3996 eng.create_rule(rule, &mut fx.g()).unwrap();
3997 });
3998
3999 let peak_delta = peak.saturating_sub(baseline);
4000
4001 assert!(
4005 peak_delta < 3 * 1024 * 1024,
4006 "peak transient delta {} bytes ({} KiB) exceeded 3 MiB; \
4007 streaming path may be building the full pairs map",
4008 peak_delta,
4009 peak_delta / 1024
4010 );
4011 assert_eq!(eng.provenance()["eq_tc"].len(), 1_000); assert!(!eng.is_tripped("eq_tc")); eprintln!(
4014 "streaming_peak_transient_bound: baseline={baseline} peak={peak} \
4015 delta={peak_delta} bytes ({} KiB)",
4016 peak_delta / 1024
4017 );
4018 }
4019
4020 fn near_threshold_pair(dim: usize, min: f64) -> (Vec<f64>, Vec<f64>) {
4027 let cos_target = min + 1e-6; let sin_small = (1.0 - cos_target * cos_target).sqrt();
4031 let mut a = vec![0.0f64; dim];
4032 a[0] = 1.0;
4033 let mut b = vec![0.0f64; dim];
4034 b[0] = cos_target;
4035 if dim > 1 {
4036 b[1] = sin_small;
4037 }
4038 (a, b)
4039 }
4040
4041 fn emb_val2(xs: &[f64]) -> Value {
4042 Value::List(xs.iter().copied().map(Value::Float).collect())
4043 }
4044
4045 fn make_early_exit_fixture(seed: u64, min: f64) -> (Fx, Vec<u32>, usize, usize) {
4049 let dims = [2usize, 4, 8, 16];
4050 let n = 100u32;
4051 let mut fx = Fx::new();
4052 let mut ids = Vec::new();
4053 for i in 0..n {
4054 let dim = dims[(i as usize) % dims.len()];
4055 let emb = rand_emb(seed, i, dim);
4056 ids.push(fx.add("Doc", &format!("d{i}"), vec![("emb", emb)]));
4057 }
4058 let (va, vb) = near_threshold_pair(8, min);
4060 let nt_a = fx.add("Doc", "nt_a", vec![("emb", emb_val2(&va))]);
4061 let nt_b = fx.add("Doc", "nt_b", vec![("emb", emb_val2(&vb))]);
4062 ids.push(nt_a);
4063 ids.push(nt_b);
4064 (fx, ids, nt_a as usize, nt_b as usize)
4065 }
4066
4067 #[test]
4071 fn vector_early_exit_identity_proof() {
4072 const SEED: u64 = 0xEA_4E_5A;
4073 const MIN: f64 = 0.85;
4074
4075 let def = RuleDef {
4076 name: "vec".into(),
4077 src_label: "Doc".into(),
4078 dst_label: "Doc".into(),
4079 predicate: Predicate::VectorSimilar {
4080 field: "emb".into(),
4081 min: MIN,
4082 },
4083 edge_type: "SIM".into(),
4084 weight_prop: Some("score".into()),
4085 max_edges: None,
4086 approximate: false,
4087 via_label: None,
4088 via_edge: None,
4089 via_dir: None,
4090 };
4091
4092 let (mut fx_on, ids, nt_a, nt_b) = make_early_exit_fixture(SEED, MIN);
4094 let (mut fx_off, _, _, _) = make_early_exit_fixture(SEED, MIN);
4095 let (fx_oracle, _, _, _) = make_early_exit_fixture(SEED, MIN);
4096
4097 let nt_a = nt_a as u32;
4098 let nt_b = nt_b as u32;
4099
4100 let mut eng_on = RuleEngine::new();
4102 {
4103 let mut g = fx_on.g();
4104 eng_on.create_rule(def.clone(), &mut g).unwrap();
4105 }
4106 let edges_on = prov_pairs(&eng_on, "vec");
4107 assert!(!edges_on.is_empty(), "should produce some edges");
4108
4109 assert!(
4111 edges_on.contains(&(nt_a, nt_b)),
4112 "near-threshold pair nt_a→nt_b must match with early-exit ON"
4113 );
4114 assert!(
4115 edges_on.contains(&(nt_b, nt_a)),
4116 "near-threshold pair nt_b→nt_a must match with early-exit ON"
4117 );
4118
4119 let mut eng_off = RuleEngine::new();
4121 {
4122 let mut g = fx_off.g();
4123 with_vector_early_exit(false, || {
4124 eng_off.create_rule(def.clone(), &mut g).unwrap();
4125 });
4126 }
4127 let edges_off = prov_pairs(&eng_off, "vec");
4128 assert_eq!(
4129 edges_on, edges_off,
4130 "early-exit ON vs OFF must produce identical edges"
4131 );
4132
4133 let mut oracle = BTreeSet::new();
4135 for &s in &ids {
4136 for &d in &ids {
4137 if s == d {
4138 continue;
4139 }
4140 let skey = fx_oracle.ids.key_of(s).unwrap();
4141 let dkey = fx_oracle.ids.key_of(d).unwrap();
4142 let sg = |f: &str| fx_oracle.props.get(s, f).cloned();
4143 let dg = |f: &str| fx_oracle.props.get(d, f).cloned();
4144 if evaluate(
4145 &def.predicate,
4146 &NodeView {
4147 key: skey,
4148 props: &sg,
4149 },
4150 &NodeView {
4151 key: dkey,
4152 props: &dg,
4153 },
4154 )
4155 .is_some()
4156 {
4157 oracle.insert((s, d));
4158 }
4159 }
4160 }
4161 assert_eq!(
4162 edges_on, oracle,
4163 "early-exit ON vs brute-force oracle must be identical"
4164 );
4165 }
4166
4167 #[test]
4170 fn vector_early_exit_checkpoint_coherence() {
4171 let mut fx = Fx::new();
4172 let a = fx.add("Doc", "a", vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0]))]);
4174 let b = fx.add("Doc", "b", vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0]))]);
4175 let c = fx.add(
4177 "Doc",
4178 "c",
4179 vec![("emb", emb_val(&[1.0, 0.0, 0.0, 0.0, 0.0, 0.0]))],
4180 );
4181 let def = RuleDef {
4182 name: "vec".into(),
4183 src_label: "Doc".into(),
4184 dst_label: "Doc".into(),
4185 predicate: Predicate::VectorSimilar {
4186 field: "emb".into(),
4187 min: 0.9,
4188 },
4189 edge_type: "SIM".into(),
4190 weight_prop: None,
4191 max_edges: None,
4192 approximate: false,
4193 via_label: None,
4194 via_edge: None,
4195 via_dir: None,
4196 };
4197
4198 let mut eng = RuleEngine::new();
4199 {
4200 let mut g = fx.g();
4201 eng.create_rule(def.clone(), &mut g).unwrap();
4202 }
4203
4204 assert!(
4206 eng.indexes["vec"].src_side.vec_ckpts(a).is_some(),
4207 "a must have src checkpoints"
4208 );
4209 assert!(
4210 eng.indexes["vec"].dst_side.vec_ckpts(b).is_some(),
4211 "b must have dst checkpoints"
4212 );
4213 assert!(
4214 eng.indexes["vec"].src_side.vec_ckpts(c).is_some(),
4215 "c must have src checkpoints (dim=6)"
4216 );
4217
4218 let ckpts_a = *eng.indexes["vec"].src_side.vec_ckpts(a).unwrap();
4220 let norm_a = eng.indexes["vec"].src_side.vec_meta(a).unwrap().1;
4221 assert!(
4222 (ckpts_a[0] - norm_a).abs() < 1e-12,
4223 "ckpts[0] must equal the full L2 norm"
4224 );
4225
4226 assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(a, b), (b, a)]));
4228
4229 let old_b = fx.props.get(b, "emb").cloned();
4231 fx.props
4232 .set(b, "emb", emb_val(&[1.0, 0.0, 0.0, 0.0, 0.0, 0.0]));
4233 {
4234 let mut g = fx.g();
4235 eng.on_node_changed(b, Some(("emb", old_b)), &mut g);
4236 }
4237 assert_eq!(eng.indexes["vec"].src_side.vec_dim(b), Some(6));
4239 assert_eq!(eng.indexes["vec"].dst_side.vec_dim(b), Some(6));
4240 assert!(eng.indexes["vec"].src_side.vec_ckpts(b).is_some());
4242 assert_eq!(prov_pairs(&eng, "vec"), BTreeSet::from([(b, c), (c, b)]));
4244
4245 let wrong_live = vec![2.0f64, 0.0, 0.0, 0.0, 0.0, 0.0]; let gate_result = eng.indexes["vec"].src_side.fresh_ckpts_for(b, &wrong_live);
4249 assert!(
4250 gate_result.is_none(),
4251 "freshness gate must reject a mismatched-norm live vector"
4252 );
4253
4254 let correct_live = vec![1.0f64, 0.0, 0.0, 0.0, 0.0, 0.0];
4256 let gate_result = eng.indexes["vec"]
4257 .src_side
4258 .fresh_ckpts_for(b, &correct_live);
4259 assert!(
4260 gate_result.is_some(),
4261 "freshness gate must accept the matching live vector"
4262 );
4263 }
4264
4265 #[test]
4274 fn vector_early_exit_razor_dim1536() {
4275 const MIN: f64 = 0.85;
4276 const DIM: usize = 1536;
4277 let target = MIN + 5e-13;
4279 let inv_sqrt = 1.0 / (DIM as f64).sqrt();
4280
4281 let a: Vec<f64> = vec![inv_sqrt; DIM];
4283
4284 let perp_scale = (1.0 - target * target).sqrt() / (2.0f64).sqrt();
4290 let mut b: Vec<f64> = vec![target * inv_sqrt; DIM];
4291 b[0] += perp_scale;
4292 b[1] -= perp_scale;
4293
4294 let def = RuleDef {
4295 name: "razor".into(),
4296 src_label: "Doc".into(),
4297 dst_label: "Doc".into(),
4298 predicate: Predicate::VectorSimilar {
4299 field: "emb".into(),
4300 min: MIN,
4301 },
4302 edge_type: "SIM".into(),
4303 weight_prop: None,
4304 max_edges: None,
4305 approximate: false,
4306 via_label: None,
4307 via_edge: None,
4308 via_dir: None,
4309 };
4310
4311 let build_fx = || {
4313 let mut fx = Fx::new();
4314 let na = fx.add("Doc", "razor_a", vec![("emb", emb_val2(&a))]);
4315 let nb = fx.add("Doc", "razor_b", vec![("emb", emb_val2(&b))]);
4316 (fx, na, nb)
4317 };
4318
4319 let (mut fx_on, na, nb) = build_fx();
4320 let (mut fx_off, _, _) = build_fx();
4321 let (fx_oracle, _, _) = build_fx();
4322
4323 let mut eng_on = RuleEngine::new();
4325 {
4326 let mut g = fx_on.g();
4327 eng_on.create_rule(def.clone(), &mut g).unwrap();
4328 }
4329 let edges_on = prov_pairs(&eng_on, "razor");
4330 assert!(
4331 edges_on.contains(&(na, nb)),
4332 "razor pair razor_a→razor_b must be present with early-exit ON (cos={target:.15}, min={MIN})"
4333 );
4334 assert!(
4335 edges_on.contains(&(nb, na)),
4336 "razor pair razor_b→razor_a must be present with early-exit ON"
4337 );
4338
4339 let mut eng_off = RuleEngine::new();
4341 {
4342 let mut g = fx_off.g();
4343 with_vector_early_exit(false, || {
4344 eng_off.create_rule(def.clone(), &mut g).unwrap();
4345 });
4346 }
4347 let edges_off = prov_pairs(&eng_off, "razor");
4348 assert_eq!(
4349 edges_on, edges_off,
4350 "razor dim=1536: early-exit ON vs OFF must produce identical edges"
4351 );
4352
4353 let ids = [na, nb];
4355 let mut oracle = BTreeSet::new();
4356 for &s in &ids {
4357 for &d in &ids {
4358 if s == d {
4359 continue;
4360 }
4361 let skey = fx_oracle.ids.key_of(s).unwrap();
4362 let dkey = fx_oracle.ids.key_of(d).unwrap();
4363 let sg = |f: &str| fx_oracle.props.get(s, f).cloned();
4364 let dg = |f: &str| fx_oracle.props.get(d, f).cloned();
4365 if evaluate(
4366 &def.predicate,
4367 &NodeView {
4368 key: skey,
4369 props: &sg,
4370 },
4371 &NodeView {
4372 key: dkey,
4373 props: &dg,
4374 },
4375 )
4376 .is_some()
4377 {
4378 oracle.insert((s, d));
4379 }
4380 }
4381 }
4382 assert_eq!(
4383 edges_on, oracle,
4384 "razor dim=1536: early-exit ON vs brute-force oracle must be identical"
4385 );
4386 }
4387}