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: Arc<IdMap>,
54 pub syms: Arc<Interner>,
55 pub topo: Arc<Topology>,
56 pub props: Arc<ColumnStore>,
57 pub labels: Arc<Vec<u32>>,
58 pub edge_props: Arc<EdgeProps>,
59 pub roles: Option<Arc<Vec<RoleDef>>>,
60 pub fulltext: Arc<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();
480 for delta in &self.deltas {
481 for rec in &delta.records {
482 apply_one(
483 Arc::make_mut(&mut w.ids),
484 Arc::make_mut(&mut w.syms),
485 Arc::make_mut(&mut w.topo),
486 Arc::make_mut(&mut w.props),
487 Arc::make_mut(&mut w.edge_props),
488 Arc::make_mut(&mut w.labels),
489 Arc::make_mut(&mut w.fulltext),
490 rec,
491 )?;
492 }
493 for &(etype, src, dst) in &delta.derived_inserts {
494 Arc::make_mut(&mut w.topo).add_edge(etype, src, dst);
495 }
496 for &(etype, src, dst) in &delta.derived_deletes {
497 Arc::make_mut(&mut w.topo).remove_edge(etype, src, dst);
498 }
499 }
500 if !self.deltas.is_empty() {
501 let cv = build_cv(&w.props, &self.base);
504 let (ids, labels, syms) = (w.ids.clone(), w.labels.clone(), w.syms.clone());
505 Arc::make_mut(&mut w.fulltext).rebuild_all(&ids, &labels, &syms, cv);
506 }
507 Ok(w)
508 }
509
510 pub(crate) fn new(
515 frozen: Arc<FrozenOverlay>,
516 base: Option<Arc<MappedBase>>,
517 deltas: Vec<Arc<CommitDelta>>,
518 version: u64,
519 role_masks: Arc<RoleMaskCache>,
520 ) -> Self {
521 Self {
522 frozen,
523 base,
524 deltas,
525 version,
526 role_masks,
527 cache: OnceLock::new(),
528 }
529 }
530
531 fn effective(&self) -> Result<&FrozenOverlay> {
541 if self.deltas.is_empty() {
542 return Ok(&self.frozen);
543 }
544 let cached = self
545 .cache
546 .get_or_init(|| self.materialize().map_err(|e| e.to_string()));
547 cached
548 .as_ref()
549 .map_err(|e| GraphError::Corrupt { detail: e.clone() })
550 }
551
552 pub fn mask_for_role(&self, role: &str) -> Result<NodeMask> {
564 self.role_masks
565 .get_or_build(role, self.version, || {
566 mask_for_role_from(self.effective()?, &self.base, role)
567 })
568 .map(|m| (*m).clone())
569 }
570
571 pub fn mask_for_namespace(&self, namespace: &str) -> Result<NodeMask> {
578 let state = self.effective()?;
579 let cv = build_cv(&state.props, &self.base);
580 let mut visible = HashSet::new();
581 for (i, &sym) in state.labels.iter().enumerate() {
582 if sym == u32::MAX {
583 continue; }
585 let id = i as u32;
586 let value = cv.get(id, core_storage::NS_PROP).map(|vr| vr.into_value());
587 if core_storage::namespace_of_value(value.as_ref()) == namespace {
588 visible.insert(id);
589 }
590 }
591 Ok(NodeMask::from_ids(visible))
592 }
593
594 pub fn resolve_key(&self, key: &str) -> Option<u32> {
599 self.effective().ok()?.ids.get(key)
600 }
601
602 pub fn query(&self, cypher: &str, params: &BTreeMap<String, Value>) -> Result<ResultSet> {
604 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
605 detail: format!("lex: {e}"),
606 })?;
607 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
608 detail: format!("parse: {e}"),
609 })?;
610 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
611 detail: format!("plan: {e}"),
612 })?;
613 let state = self.effective()?;
614 let view = make_view(state, &self.base, None);
615 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
616 detail: format!("execute: {e}"),
617 })
618 }
619
620 pub fn query_masked(
624 &self,
625 cypher: &str,
626 params: &BTreeMap<String, Value>,
627 mask: &NodeMask,
628 ) -> Result<ResultSet> {
629 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
630 detail: format!("lex: {e}"),
631 })?;
632 if is_write_tokens(&tokens) {
633 return Err(GraphError::QueryError {
634 detail: "masked queries are read-only".into(),
635 });
636 }
637 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
638 detail: format!("parse: {e}"),
639 })?;
640 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
641 detail: format!("plan: {e}"),
642 })?;
643 let state = self.effective()?;
644 let view = make_view(state, &self.base, Some(&mask.visible));
645 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
646 detail: format!("execute: {e}"),
647 })
648 }
649
650 pub fn node_info(&self, key: &str) -> Option<NodeInfo> {
652 node_info_from(key, self.effective().ok()?, &self.base)
653 }
654
655 pub fn node_edges(&self, key: &str) -> Result<Vec<EdgeInfo>> {
660 node_edges_from(key, self.effective()?, &self.base)
661 }
662
663 pub fn neighborhood_masked(
668 &self,
669 key: &str,
670 depth: u32,
671 edge_types: Option<&[&str]>,
672 dir: Dir,
673 mask: &NodeMask,
674 ) -> Option<ResultSet> {
675 neighborhood_masked_from(
676 key,
677 self.effective().ok()?,
678 &self.base,
679 depth,
680 edge_types,
681 dir,
682 mask,
683 )
684 }
685
686 pub fn node_edges_scoped(&self, key: &str, mask: &NodeMask) -> Result<Vec<EdgeInfo>> {
700 let state = self.effective()?;
701 if !state.ids.get(key).is_some_and(|id| mask.contains_id(id)) {
702 return Err(GraphError::KeyNotFound { key: key.into() });
703 }
704 let edges = node_edges_from(key, state, &self.base)?;
705 Ok(edges
706 .into_iter()
707 .filter(|e| {
708 let other = if e.src_key == key {
709 &e.dst_key
710 } else {
711 &e.src_key
712 };
713 state.ids.get(other).is_some_and(|id| mask.contains_id(id))
714 })
715 .collect())
716 }
717
718 pub fn neighborhood_scoped(
728 &self,
729 key: &str,
730 depth: u32,
731 edge_types: Option<&[&str]>,
732 dir: Dir,
733 mask: &NodeMask,
734 ) -> Result<ResultSet> {
735 let state = self.effective()?;
736 if !state.ids.get(key).is_some_and(|id| mask.contains_id(id)) {
737 return Err(GraphError::KeyNotFound { key: key.into() });
738 }
739 neighborhood_masked_from(key, state, &self.base, depth, edge_types, dir, mask)
740 .ok_or_else(|| GraphError::KeyNotFound { key: key.into() })
741 }
742}
743
744fn node_info_from(
747 key: &str,
748 state: &FrozenOverlay,
749 base: &Option<Arc<MappedBase>>,
750) -> Option<NodeInfo> {
751 let id = state.ids.get(key)?;
752 let label_sym = *state.labels.get(id as usize)?;
753 if label_sym == u32::MAX {
754 return None;
755 }
756 let label = state.syms.resolve(label_sym)?.to_string();
757 let cv = build_cv(&state.props, base);
758 let mut props = BTreeMap::new();
759 for field in cv.field_names() {
760 if let Some(vr) = cv.get(id, &field) {
761 props.insert(field, vr.into_value());
762 }
763 }
764 Some(NodeInfo {
765 key: key.to_string(),
766 label,
767 props,
768 })
769}
770
771fn node_edges_from(
772 key: &str,
773 state: &FrozenOverlay,
774 base: &Option<Arc<MappedBase>>,
775) -> Result<Vec<EdgeInfo>> {
776 let id = state
777 .ids
778 .get(key)
779 .ok_or_else(|| GraphError::KeyNotFound { key: key.into() })?;
780 let tv = build_tv(&state.topo, base);
781 let mut edges = Vec::new();
782 for etype in tv.etypes() {
783 let edge_type = state
784 .syms
785 .resolve(etype)
786 .ok_or_else(|| GraphError::Corrupt {
787 detail: format!("reader: topology etype {etype} not in interner"),
788 })?
789 .to_string();
790 for dir in [Direction::Out, Direction::In] {
791 for &nbr in tv.neighbors(etype, dir, id).as_ref() {
792 let (src_key, dst_key) = match dir {
793 Direction::Out => (
794 key.to_string(),
795 state
796 .ids
797 .key_of(nbr)
798 .ok_or_else(|| GraphError::Corrupt {
799 detail: format!("topology id {nbr} has no key"),
800 })?
801 .to_string(),
802 ),
803 Direction::In => (
804 state
805 .ids
806 .key_of(nbr)
807 .ok_or_else(|| GraphError::Corrupt {
808 detail: format!("topology id {nbr} has no key"),
809 })?
810 .to_string(),
811 key.to_string(),
812 ),
813 };
814 edges.push(EdgeInfo {
815 edge_type: edge_type.clone(),
816 src_key,
817 dst_key,
818 derived: false,
819 });
820 }
821 }
822 }
823 edges.sort_by(|a, b| {
824 a.edge_type
825 .cmp(&b.edge_type)
826 .then(a.src_key.cmp(&b.src_key))
827 .then(a.dst_key.cmp(&b.dst_key))
828 });
829 edges.dedup();
830 Ok(edges)
831}
832
833fn neighborhood_masked_from(
834 key: &str,
835 state: &FrozenOverlay,
836 base: &Option<Arc<MappedBase>>,
837 depth: u32,
838 edge_types: Option<&[&str]>,
839 dir: Dir,
840 mask: &NodeMask,
841) -> Option<ResultSet> {
842 let start_id = state.ids.get(key)?;
843 let view = make_view(state, base, Some(&mask.visible));
844 let resolved: Option<Vec<u32>> = edge_types.map(|names| {
845 names
846 .iter()
847 .filter_map(|name| view.syms.get(name))
848 .collect()
849 });
850 let nb = neighborhood(&view, start_id, depth, resolved.as_deref(), dir);
851 let mut rs = ResultSet::new(vec!["key".into(), "label".into(), "depth".into()]);
852 let mut visited: Vec<(u32, u32)> = Vec::with_capacity(nb.nodes.len() + 1);
854 visited.push((start_id, 0));
855 for (nid, d) in &nb.nodes {
856 let k = view.key_of(*nid);
857 let lbl = view
858 .label_of(*nid)
859 .expect("real nodes always have a label; u32::MAX sentinel cannot occur");
860 rs.push_row(vec![
861 Some(Value::Str(k.to_string())),
862 Some(Value::Str(lbl.to_string())),
863 Some(Value::Int(*d as i64)),
864 ]);
865 visited.push((*nid, *d));
866 }
867 if mask.mode() == crate::mask::MaskMode::Stub {
875 let raw_view = make_view(state, base, None);
876 let mut seen: HashSet<u32> = visited.iter().map(|(id, _)| *id).collect();
877 for (node_id, node_depth) in &visited {
878 if *node_depth >= depth {
879 continue;
880 }
881 for e in expand(&raw_view, *node_id, resolved.as_deref(), dir) {
882 let nbr = if e.src == *node_id { e.dst } else { e.src };
883 if !mask.contains_id(nbr) && seen.insert(nbr) {
884 if let Some(k) = state.ids.key_of(nbr) {
885 rs.push_row(vec![
886 Some(Value::Str(k.to_string())),
887 None,
888 Some(Value::Int((*node_depth + 1) as i64)),
889 ]);
890 }
891 }
892 }
893 }
894 }
895 Some(rs)
896}