1use std::collections::{BTreeMap, HashSet};
17use std::sync::{Arc, OnceLock};
18
19use core_query::cypher::{execute, is_write_tokens, lex, parse, plan, Params};
20use core_query::{expand, neighborhood, Dir, GraphView, ResultSet};
21use core_storage::fulltext::FulltextIndex;
22use core_storage::v8::seam::{ColumnsView, TopologyView};
23use core_storage::v8::MappedBase;
24use core_storage::wal::WalRecord;
25use core_storage::{
26 ColumnStore, Direction, EdgeProps, EdgePropsView, GraphError, IdMap, Interner, Result,
27 Topology, Value,
28};
29
30use crate::db::{EdgeInfo, NodeInfo};
31use crate::mask::{NodeMask, RoleMaskCache};
32use crate::roles::RoleDef;
33
34pub const FOLD_EVERY_K: usize = 64;
37
38pub struct CommitDelta {
40 pub records: Vec<WalRecord>,
42 pub derived_inserts: Vec<(u32, u32, u32)>,
44 pub derived_deletes: Vec<(u32, u32, u32)>,
46}
47
48#[derive(Clone)]
52pub struct FrozenOverlay {
53 pub ids: IdMap,
54 pub syms: Interner,
55 pub topo: Topology,
56 pub props: ColumnStore,
57 pub labels: Vec<u32>,
58 pub edge_props: EdgeProps,
59 pub roles: Option<Vec<RoleDef>>,
60 pub fulltext: FulltextIndex,
61}
62
63pub struct ReaderSnapshot {
72 pub frozen: Arc<FrozenOverlay>,
74 pub base: Option<Arc<MappedBase>>,
76 pub deltas: Vec<Arc<CommitDelta>>,
78 pub version: u64,
82 role_masks: Arc<RoleMaskCache>,
84 cache: OnceLock<std::result::Result<FrozenOverlay, String>>,
92}
93
94fn build_tv<'a>(topo: &'a Topology, base: &'a Option<Arc<MappedBase>>) -> TopologyView<'a> {
97 match base {
98 None => TopologyView::owned(topo),
99 Some(b) => {
100 let csr = b.topology().expect("base CSR CRC already verified at open");
101 TopologyView::with_base(topo, csr)
102 }
103 }
104}
105
106fn build_cv<'a>(props: &'a ColumnStore, base: &'a Option<Arc<MappedBase>>) -> ColumnsView<'a> {
107 match base {
108 None => ColumnsView::owned(props),
109 Some(b) => {
110 let cols = b
111 .columns()
112 .expect("base columns CRC already verified at open");
113 let strings = b
114 .string_table()
115 .transpose()
116 .expect("base strings CRC already verified at open");
117 ColumnsView::with_base_cached(props, cols, b.mixed_cache()).with_shared_strings(strings)
118 }
119 }
120}
121
122fn build_epv<'a>(
123 edge_props: &'a EdgeProps,
124 base: &'a Option<Arc<MappedBase>>,
125) -> EdgePropsView<'a> {
126 match base {
127 None => EdgePropsView::owned(edge_props),
128 Some(b) => {
129 let archived = b
130 .edge_props_section()
131 .expect("base edge_props CRC already verified at open");
132 EdgePropsView::with_base(edge_props, archived)
133 }
134 }
135}
136
137fn make_view<'a>(
138 state: &'a FrozenOverlay,
139 base: &'a Option<Arc<MappedBase>>,
140 mask: Option<&'a core_query::visible::VisibleSet>,
141) -> GraphView<'a> {
142 GraphView {
143 ids: &state.ids,
144 syms: &state.syms,
145 labels: &state.labels,
146 props: build_cv(&state.props, base),
147 topo: build_tv(&state.topo, base),
148 edge_props: build_epv(&state.edge_props, base),
149 mask,
150 prop_index: None,
153 }
154}
155
156#[allow(clippy::too_many_arguments)]
165fn apply_one(
166 ids: &mut IdMap,
167 syms: &mut Interner,
168 topo: &mut Topology,
169 props: &mut ColumnStore,
170 edge_props: &mut EdgeProps,
171 labels: &mut Vec<u32>,
172 fulltext: &mut FulltextIndex,
173 rec: &WalRecord,
174) -> Result<()> {
175 match rec {
176 WalRecord::Intern { id, text } => {
177 let got = syms.intern(text);
178 if got != *id {
179 return Err(GraphError::Corrupt {
180 detail: format!(
181 "mvcc delta intern mismatch for {text:?}: expected {id}, got {got}"
182 ),
183 });
184 }
185 }
186
187 WalRecord::InsertNodeId {
188 label,
189 key,
190 props: node_props,
191 } => {
192 let node_id = ids.try_insert(key)?;
193 if labels.len() <= node_id as usize {
194 labels.resize(node_id as usize + 1, u32::MAX);
195 }
196 labels[node_id as usize] = *label;
197 let label_str = syms
198 .resolve(*label)
199 .ok_or_else(|| GraphError::Corrupt {
200 detail: format!("mvcc delta: unknown label sym {label}"),
201 })?
202 .to_string();
203 for (field_sym, value) in node_props {
204 let field = syms
205 .resolve(*field_sym)
206 .ok_or_else(|| GraphError::Corrupt {
207 detail: format!("mvcc delta: unknown field sym {field_sym}"),
208 })?
209 .to_string();
210 props.set(node_id, &field, value.clone());
211 if fulltext.is_enabled(&label_str, &field) {
212 fulltext.add_tokens(node_id, &field, value);
213 }
214 }
215 }
216
217 WalRecord::InsertNode {
218 label,
219 key,
220 props: node_props,
221 } => {
222 let label_sym = syms.intern(label);
223 let node_id = ids.try_insert(key)?;
224 if labels.len() <= node_id as usize {
225 labels.resize(node_id as usize + 1, u32::MAX);
226 }
227 labels[node_id as usize] = label_sym;
228 for (field, value) in node_props {
229 props.set(node_id, field, value.clone());
230 if fulltext.is_enabled(label, field) {
231 fulltext.add_tokens(node_id, field, value);
232 }
233 }
234 }
235
236 WalRecord::SetPropId { id, field, value } => {
237 if let Some(field_str) = syms.resolve(*field).map(str::to_string) {
238 props.set(*id, &field_str, value.clone());
239 if let Some(&label_sym) = labels.get(*id as usize) {
240 if let Some(label_str) = syms.resolve(label_sym) {
241 if fulltext.is_enabled(label_str, &field_str) {
242 fulltext.add_tokens(*id, &field_str, value);
243 }
244 }
245 }
246 }
247 }
248
249 WalRecord::SetProp { key, field, value } => {
250 if let Some(node_id) = ids.get(key) {
251 props.set(node_id, field, value.clone());
252 if let Some(&label_sym) = labels.get(node_id as usize) {
253 if let Some(label_str) = syms.resolve(label_sym) {
254 if fulltext.is_enabled(label_str, field) {
255 fulltext.add_tokens(node_id, field, value);
256 }
257 }
258 }
259 }
260 }
261
262 WalRecord::RemoveProp { key, field } => {
263 if let Some(node_id) = ids.get(key) {
264 props.remove(node_id, field);
265 props.record_prop_tombstone(node_id, field);
270 fulltext.remove_node_field(node_id, field);
271 }
272 }
273
274 WalRecord::DeleteNode { key } => {
275 if let Some(node_id) = ids.delete(key) {
276 props.remove_all(node_id);
277 fulltext.remove_node(node_id);
278 if let Some(slot) = labels.get_mut(node_id as usize) {
280 *slot = u32::MAX;
281 }
282 let etypes: Vec<u32> = topo.etypes().collect();
286 let mut doomed = Vec::new();
287 for et in &etypes {
288 for &dst in topo.neighbors(*et, Direction::Out, node_id).as_ref() {
289 doomed.push((*et, node_id, dst));
290 }
291 for &src in topo.neighbors(*et, Direction::In, node_id).as_ref() {
292 doomed.push((*et, src, node_id));
293 }
294 }
295 for (et, s, d) in doomed {
296 topo.remove_edge(et, s, d);
297 edge_props.remove_edge(et, s, d);
298 }
299 }
300 }
301
302 WalRecord::InsertEdgeId { etype, src, dst } => {
303 topo.add_edge(*etype, *src, *dst);
304 }
305
306 WalRecord::InsertEdge {
307 edge_type,
308 src_key,
309 dst_key,
310 } => {
311 let etype = syms.intern(edge_type);
312 if let (Some(src), Some(dst)) = (ids.get(src_key), ids.get(dst_key)) {
313 topo.add_edge(etype, src, dst);
314 }
315 }
316
317 WalRecord::DeleteEdge {
318 edge_type,
319 src_key,
320 dst_key,
321 } => {
322 if let Some(etype) = syms.get(edge_type) {
323 if let (Some(src), Some(dst)) = (ids.get(src_key), ids.get(dst_key)) {
324 topo.remove_edge(etype, src, dst);
325 }
326 }
327 }
328
329 WalRecord::EnableFulltext { label, field } => {
330 fulltext.enable(label, field);
331 }
332
333 WalRecord::DisableFulltext { label, field } => {
334 fulltext.disable(label, field);
335 }
336
337 WalRecord::Batch(inner) => {
338 for r in inner {
339 apply_one(ids, syms, topo, props, edge_props, labels, fulltext, r)?;
340 }
341 }
342
343 WalRecord::CreateRule { .. }
346 | WalRecord::DeleteRule { .. }
347 | WalRecord::RebuildRule { .. }
348 | WalRecord::CreateView { .. }
349 | WalRecord::DeleteView { .. }
350 | WalRecord::EnableIndex { .. }
354 | WalRecord::DisableIndex { .. }
355 | WalRecord::SetEdgeCount { count: 0, .. }
359 | WalRecord::DerivedEdgeAdded { .. }
363 | WalRecord::DerivedEdgeRetracted { .. } => {}
364
365 WalRecord::SetEdgeCount {
368 etype,
369 src,
370 dst,
371 count,
372 } => {
373 edge_props.set(
374 *etype,
375 *src,
376 *dst,
377 crate::db::EDGE_COUNT_PROP,
378 core_storage::Value::Int(*count as i64),
379 );
380 }
381
382 WalRecord::RenameNode { old_key, new_key } => {
383 if ids.get(old_key).is_some() {
386 ids.rename(old_key, new_key).map_err(|_| GraphError::Corrupt {
387 detail: format!("mvcc delta RenameNode {old_key}→{new_key} failed"),
388 })?;
389 }
390 }
391 }
392 Ok(())
393}
394
395fn mask_for_role_from(
404 state: &FrozenOverlay,
405 base: &Option<Arc<MappedBase>>,
406 role: &str,
407) -> Result<NodeMask> {
408 let roles = state.roles.as_ref().ok_or_else(|| GraphError::Corrupt {
409 detail: "roles.json was corrupt at open; fix the file and re-open".into(),
410 })?;
411 let def = roles
412 .iter()
413 .find(|r| r.name == role)
414 .ok_or_else(|| GraphError::KeyNotFound {
415 key: format!("role:{role}"),
416 })?;
417 let mut visible = HashSet::new();
418 for key in &def.keys {
420 if let Some(id) = state.ids.get(key) {
421 visible.insert(id);
422 }
423 }
424 let props = def
425 .visible_where
426 .as_ref()
427 .map(|_| build_cv(&state.props, base));
428 for label_name in &def.labels {
429 if let Some(sym) = state.syms.get(label_name) {
430 for (i, &s) in state.labels.iter().enumerate() {
431 if s != sym {
432 continue;
433 }
434 let id = i as u32;
435 match (&def.visible_where, &props) {
436 (Some(pred), Some(view)) => {
437 let value = view.get(id, &pred.field).map(|vr| vr.into_value());
438 if pred.holds(value.as_ref()) {
439 visible.insert(id);
440 }
441 }
442 _ => {
443 visible.insert(id);
444 }
445 }
446 }
447 }
448 }
449 if def.namespaces.is_some() {
455 let cv = build_cv(&state.props, base);
456 visible.retain(|&id| {
457 let value = cv.get(id, core_storage::NS_PROP).map(|vr| vr.into_value());
458 def.sees_namespace(core_storage::namespace_of_value(value.as_ref()))
459 });
460 }
461 Ok(NodeMask::from_ids(visible))
462}
463
464impl ReaderSnapshot {
465 fn materialize(&self) -> Result<FrozenOverlay> {
470 let mut w = (*self.frozen).clone();
471 for delta in &self.deltas {
472 for rec in &delta.records {
473 apply_one(
474 &mut w.ids,
475 &mut w.syms,
476 &mut w.topo,
477 &mut w.props,
478 &mut w.edge_props,
479 &mut w.labels,
480 &mut w.fulltext,
481 rec,
482 )?;
483 }
484 for &(etype, src, dst) in &delta.derived_inserts {
485 w.topo.add_edge(etype, src, dst);
486 }
487 for &(etype, src, dst) in &delta.derived_deletes {
488 w.topo.remove_edge(etype, src, dst);
489 }
490 }
491 if !self.deltas.is_empty() {
492 let cv = build_cv(&w.props, &self.base);
495 w.fulltext.rebuild_all(&w.ids, &w.labels, &w.syms, cv);
496 }
497 Ok(w)
498 }
499
500 pub(crate) fn new(
505 frozen: Arc<FrozenOverlay>,
506 base: Option<Arc<MappedBase>>,
507 deltas: Vec<Arc<CommitDelta>>,
508 version: u64,
509 role_masks: Arc<RoleMaskCache>,
510 ) -> Self {
511 Self {
512 frozen,
513 base,
514 deltas,
515 version,
516 role_masks,
517 cache: OnceLock::new(),
518 }
519 }
520
521 fn effective(&self) -> Result<&FrozenOverlay> {
531 if self.deltas.is_empty() {
532 return Ok(&self.frozen);
533 }
534 let cached = self
535 .cache
536 .get_or_init(|| self.materialize().map_err(|e| e.to_string()));
537 cached
538 .as_ref()
539 .map_err(|e| GraphError::Corrupt { detail: e.clone() })
540 }
541
542 pub fn mask_for_role(&self, role: &str) -> Result<NodeMask> {
554 self.role_masks
555 .get_or_build(role, self.version, || {
556 mask_for_role_from(self.effective()?, &self.base, role)
557 })
558 .map(|m| (*m).clone())
559 }
560
561 pub fn mask_for_namespace(&self, namespace: &str) -> Result<NodeMask> {
568 let state = self.effective()?;
569 let cv = build_cv(&state.props, &self.base);
570 let mut visible = HashSet::new();
571 for (i, &sym) in state.labels.iter().enumerate() {
572 if sym == u32::MAX {
573 continue; }
575 let id = i as u32;
576 let value = cv.get(id, core_storage::NS_PROP).map(|vr| vr.into_value());
577 if core_storage::namespace_of_value(value.as_ref()) == namespace {
578 visible.insert(id);
579 }
580 }
581 Ok(NodeMask::from_ids(visible))
582 }
583
584 pub fn resolve_key(&self, key: &str) -> Option<u32> {
589 self.effective().ok()?.ids.get(key)
590 }
591
592 pub fn query(&self, cypher: &str, params: &BTreeMap<String, Value>) -> Result<ResultSet> {
594 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
595 detail: format!("lex: {e}"),
596 })?;
597 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
598 detail: format!("parse: {e}"),
599 })?;
600 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
601 detail: format!("plan: {e}"),
602 })?;
603 let state = self.effective()?;
604 let view = make_view(state, &self.base, None);
605 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
606 detail: format!("execute: {e}"),
607 })
608 }
609
610 pub fn query_masked(
614 &self,
615 cypher: &str,
616 params: &BTreeMap<String, Value>,
617 mask: &NodeMask,
618 ) -> Result<ResultSet> {
619 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
620 detail: format!("lex: {e}"),
621 })?;
622 if is_write_tokens(&tokens) {
623 return Err(GraphError::QueryError {
624 detail: "masked queries are read-only".into(),
625 });
626 }
627 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
628 detail: format!("parse: {e}"),
629 })?;
630 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
631 detail: format!("plan: {e}"),
632 })?;
633 let state = self.effective()?;
634 let view = make_view(state, &self.base, Some(&mask.visible));
635 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
636 detail: format!("execute: {e}"),
637 })
638 }
639
640 pub fn node_info(&self, key: &str) -> Option<NodeInfo> {
642 node_info_from(key, self.effective().ok()?, &self.base)
643 }
644
645 pub fn node_edges(&self, key: &str) -> Result<Vec<EdgeInfo>> {
650 node_edges_from(key, self.effective()?, &self.base)
651 }
652
653 pub fn neighborhood_masked(
658 &self,
659 key: &str,
660 depth: u32,
661 edge_types: Option<&[&str]>,
662 dir: Dir,
663 mask: &NodeMask,
664 ) -> Option<ResultSet> {
665 neighborhood_masked_from(
666 key,
667 self.effective().ok()?,
668 &self.base,
669 depth,
670 edge_types,
671 dir,
672 mask,
673 )
674 }
675
676 pub fn node_edges_scoped(&self, key: &str, mask: &NodeMask) -> Result<Vec<EdgeInfo>> {
690 let state = self.effective()?;
691 if !state.ids.get(key).is_some_and(|id| mask.contains_id(id)) {
692 return Err(GraphError::KeyNotFound { key: key.into() });
693 }
694 let edges = node_edges_from(key, state, &self.base)?;
695 Ok(edges
696 .into_iter()
697 .filter(|e| {
698 let other = if e.src_key == key {
699 &e.dst_key
700 } else {
701 &e.src_key
702 };
703 state.ids.get(other).is_some_and(|id| mask.contains_id(id))
704 })
705 .collect())
706 }
707
708 pub fn neighborhood_scoped(
718 &self,
719 key: &str,
720 depth: u32,
721 edge_types: Option<&[&str]>,
722 dir: Dir,
723 mask: &NodeMask,
724 ) -> Result<ResultSet> {
725 let state = self.effective()?;
726 if !state.ids.get(key).is_some_and(|id| mask.contains_id(id)) {
727 return Err(GraphError::KeyNotFound { key: key.into() });
728 }
729 neighborhood_masked_from(key, state, &self.base, depth, edge_types, dir, mask)
730 .ok_or_else(|| GraphError::KeyNotFound { key: key.into() })
731 }
732}
733
734fn node_info_from(
737 key: &str,
738 state: &FrozenOverlay,
739 base: &Option<Arc<MappedBase>>,
740) -> Option<NodeInfo> {
741 let id = state.ids.get(key)?;
742 let label_sym = *state.labels.get(id as usize)?;
743 if label_sym == u32::MAX {
744 return None;
745 }
746 let label = state.syms.resolve(label_sym)?.to_string();
747 let cv = build_cv(&state.props, base);
748 let mut props = BTreeMap::new();
749 for field in cv.field_names() {
750 if let Some(vr) = cv.get(id, &field) {
751 props.insert(field, vr.into_value());
752 }
753 }
754 Some(NodeInfo {
755 key: key.to_string(),
756 label,
757 props,
758 })
759}
760
761fn node_edges_from(
762 key: &str,
763 state: &FrozenOverlay,
764 base: &Option<Arc<MappedBase>>,
765) -> Result<Vec<EdgeInfo>> {
766 let id = state
767 .ids
768 .get(key)
769 .ok_or_else(|| GraphError::KeyNotFound { key: key.into() })?;
770 let tv = build_tv(&state.topo, base);
771 let mut edges = Vec::new();
772 for etype in tv.etypes() {
773 let edge_type = state
774 .syms
775 .resolve(etype)
776 .ok_or_else(|| GraphError::Corrupt {
777 detail: format!("reader: topology etype {etype} not in interner"),
778 })?
779 .to_string();
780 for dir in [Direction::Out, Direction::In] {
781 for &nbr in tv.neighbors(etype, dir, id).as_ref() {
782 let (src_key, dst_key) = match dir {
783 Direction::Out => (
784 key.to_string(),
785 state
786 .ids
787 .key_of(nbr)
788 .ok_or_else(|| GraphError::Corrupt {
789 detail: format!("topology id {nbr} has no key"),
790 })?
791 .to_string(),
792 ),
793 Direction::In => (
794 state
795 .ids
796 .key_of(nbr)
797 .ok_or_else(|| GraphError::Corrupt {
798 detail: format!("topology id {nbr} has no key"),
799 })?
800 .to_string(),
801 key.to_string(),
802 ),
803 };
804 edges.push(EdgeInfo {
805 edge_type: edge_type.clone(),
806 src_key,
807 dst_key,
808 derived: false,
809 });
810 }
811 }
812 }
813 edges.sort_by(|a, b| {
814 a.edge_type
815 .cmp(&b.edge_type)
816 .then(a.src_key.cmp(&b.src_key))
817 .then(a.dst_key.cmp(&b.dst_key))
818 });
819 edges.dedup();
820 Ok(edges)
821}
822
823fn neighborhood_masked_from(
824 key: &str,
825 state: &FrozenOverlay,
826 base: &Option<Arc<MappedBase>>,
827 depth: u32,
828 edge_types: Option<&[&str]>,
829 dir: Dir,
830 mask: &NodeMask,
831) -> Option<ResultSet> {
832 let start_id = state.ids.get(key)?;
833 let view = make_view(state, base, Some(&mask.visible));
834 let resolved: Option<Vec<u32>> = edge_types.map(|names| {
835 names
836 .iter()
837 .filter_map(|name| view.syms.get(name))
838 .collect()
839 });
840 let nb = neighborhood(&view, start_id, depth, resolved.as_deref(), dir);
841 let mut rs = ResultSet::new(vec!["key".into(), "label".into(), "depth".into()]);
842 let mut visited: Vec<(u32, u32)> = Vec::with_capacity(nb.nodes.len() + 1);
844 visited.push((start_id, 0));
845 for (nid, d) in &nb.nodes {
846 let k = view.key_of(*nid);
847 let lbl = view
848 .label_of(*nid)
849 .expect("real nodes always have a label; u32::MAX sentinel cannot occur");
850 rs.push_row(vec![
851 Some(Value::Str(k.to_string())),
852 Some(Value::Str(lbl.to_string())),
853 Some(Value::Int(*d as i64)),
854 ]);
855 visited.push((*nid, *d));
856 }
857 if mask.mode() == crate::mask::MaskMode::Stub {
865 let raw_view = make_view(state, base, None);
866 let mut seen: HashSet<u32> = visited.iter().map(|(id, _)| *id).collect();
867 for (node_id, node_depth) in &visited {
868 if *node_depth >= depth {
869 continue;
870 }
871 for e in expand(&raw_view, *node_id, resolved.as_deref(), dir) {
872 let nbr = if e.src == *node_id { e.dst } else { e.src };
873 if !mask.contains_id(nbr) && seen.insert(nbr) {
874 if let Some(k) = state.ids.key_of(nbr) {
875 rs.push_row(vec![
876 Some(Value::Str(k.to_string())),
877 None,
878 Some(Value::Int((*node_depth + 1) as i64)),
879 ]);
880 }
881 }
882 }
883 }
884 }
885 Some(rs)
886}