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 HashSet<u32>>,
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::DerivedEdgeAdded { .. }
359 | WalRecord::DerivedEdgeRetracted { .. } => {}
360
361 WalRecord::RenameNode { old_key, new_key } => {
362 if ids.get(old_key).is_some() {
365 ids.rename(old_key, new_key).map_err(|_| GraphError::Corrupt {
366 detail: format!("mvcc delta RenameNode {old_key}→{new_key} failed"),
367 })?;
368 }
369 }
370 }
371 Ok(())
372}
373
374fn mask_for_role_from(
383 state: &FrozenOverlay,
384 base: &Option<Arc<MappedBase>>,
385 role: &str,
386) -> Result<NodeMask> {
387 let roles = state.roles.as_ref().ok_or_else(|| GraphError::Corrupt {
388 detail: "roles.json was corrupt at open; fix the file and re-open".into(),
389 })?;
390 let def = roles
391 .iter()
392 .find(|r| r.name == role)
393 .ok_or_else(|| GraphError::KeyNotFound {
394 key: format!("role:{role}"),
395 })?;
396 let mut visible = HashSet::new();
397 for key in &def.keys {
399 if let Some(id) = state.ids.get(key) {
400 visible.insert(id);
401 }
402 }
403 let props = def
404 .visible_where
405 .as_ref()
406 .map(|_| build_cv(&state.props, base));
407 for label_name in &def.labels {
408 if let Some(sym) = state.syms.get(label_name) {
409 for (i, &s) in state.labels.iter().enumerate() {
410 if s != sym {
411 continue;
412 }
413 let id = i as u32;
414 match (&def.visible_where, &props) {
415 (Some(pred), Some(view)) => {
416 let value = view.get(id, &pred.field).map(|vr| vr.into_value());
417 if pred.holds(value.as_ref()) {
418 visible.insert(id);
419 }
420 }
421 _ => {
422 visible.insert(id);
423 }
424 }
425 }
426 }
427 }
428 Ok(NodeMask::from_ids(visible))
429}
430
431impl ReaderSnapshot {
432 fn materialize(&self) -> Result<FrozenOverlay> {
437 let mut w = (*self.frozen).clone();
438 for delta in &self.deltas {
439 for rec in &delta.records {
440 apply_one(
441 &mut w.ids,
442 &mut w.syms,
443 &mut w.topo,
444 &mut w.props,
445 &mut w.edge_props,
446 &mut w.labels,
447 &mut w.fulltext,
448 rec,
449 )?;
450 }
451 for &(etype, src, dst) in &delta.derived_inserts {
452 w.topo.add_edge(etype, src, dst);
453 }
454 for &(etype, src, dst) in &delta.derived_deletes {
455 w.topo.remove_edge(etype, src, dst);
456 }
457 }
458 if !self.deltas.is_empty() {
459 let cv = build_cv(&w.props, &self.base);
462 w.fulltext.rebuild_all(&w.ids, &w.labels, &w.syms, cv);
463 }
464 Ok(w)
465 }
466
467 pub(crate) fn new(
472 frozen: Arc<FrozenOverlay>,
473 base: Option<Arc<MappedBase>>,
474 deltas: Vec<Arc<CommitDelta>>,
475 version: u64,
476 role_masks: Arc<RoleMaskCache>,
477 ) -> Self {
478 Self {
479 frozen,
480 base,
481 deltas,
482 version,
483 role_masks,
484 cache: OnceLock::new(),
485 }
486 }
487
488 fn effective(&self) -> Result<&FrozenOverlay> {
498 if self.deltas.is_empty() {
499 return Ok(&self.frozen);
500 }
501 let cached = self
502 .cache
503 .get_or_init(|| self.materialize().map_err(|e| e.to_string()));
504 cached
505 .as_ref()
506 .map_err(|e| GraphError::Corrupt { detail: e.clone() })
507 }
508
509 pub fn mask_for_role(&self, role: &str) -> Result<NodeMask> {
521 self.role_masks
522 .get_or_build(role, self.version, || {
523 mask_for_role_from(self.effective()?, &self.base, role)
524 })
525 .map(|m| (*m).clone())
526 }
527
528 pub fn resolve_key(&self, key: &str) -> Option<u32> {
533 self.effective().ok()?.ids.get(key)
534 }
535
536 pub fn query(&self, cypher: &str, params: &BTreeMap<String, Value>) -> Result<ResultSet> {
538 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
539 detail: format!("lex: {e}"),
540 })?;
541 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
542 detail: format!("parse: {e}"),
543 })?;
544 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
545 detail: format!("plan: {e}"),
546 })?;
547 let state = self.effective()?;
548 let view = make_view(state, &self.base, None);
549 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
550 detail: format!("execute: {e}"),
551 })
552 }
553
554 pub fn query_masked(
558 &self,
559 cypher: &str,
560 params: &BTreeMap<String, Value>,
561 mask: &NodeMask,
562 ) -> Result<ResultSet> {
563 let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
564 detail: format!("lex: {e}"),
565 })?;
566 if is_write_tokens(&tokens) {
567 return Err(GraphError::QueryError {
568 detail: "masked queries are read-only".into(),
569 });
570 }
571 let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
572 detail: format!("parse: {e}"),
573 })?;
574 let ops = plan(&ast).map_err(|e| GraphError::QueryError {
575 detail: format!("plan: {e}"),
576 })?;
577 let state = self.effective()?;
578 let view = make_view(state, &self.base, Some(&mask.visible));
579 execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
580 detail: format!("execute: {e}"),
581 })
582 }
583
584 pub fn node_info(&self, key: &str) -> Option<NodeInfo> {
586 node_info_from(key, self.effective().ok()?, &self.base)
587 }
588
589 pub fn node_edges(&self, key: &str) -> Result<Vec<EdgeInfo>> {
594 node_edges_from(key, self.effective()?, &self.base)
595 }
596
597 pub fn neighborhood_masked(
602 &self,
603 key: &str,
604 depth: u32,
605 edge_types: Option<&[&str]>,
606 dir: Dir,
607 mask: &NodeMask,
608 ) -> Option<ResultSet> {
609 neighborhood_masked_from(
610 key,
611 self.effective().ok()?,
612 &self.base,
613 depth,
614 edge_types,
615 dir,
616 mask,
617 )
618 }
619}
620
621fn node_info_from(
624 key: &str,
625 state: &FrozenOverlay,
626 base: &Option<Arc<MappedBase>>,
627) -> Option<NodeInfo> {
628 let id = state.ids.get(key)?;
629 let label_sym = *state.labels.get(id as usize)?;
630 if label_sym == u32::MAX {
631 return None;
632 }
633 let label = state.syms.resolve(label_sym)?.to_string();
634 let cv = build_cv(&state.props, base);
635 let mut props = BTreeMap::new();
636 for field in cv.field_names() {
637 if let Some(vr) = cv.get(id, &field) {
638 props.insert(field, vr.into_value());
639 }
640 }
641 Some(NodeInfo {
642 key: key.to_string(),
643 label,
644 props,
645 })
646}
647
648fn node_edges_from(
649 key: &str,
650 state: &FrozenOverlay,
651 base: &Option<Arc<MappedBase>>,
652) -> Result<Vec<EdgeInfo>> {
653 let id = state
654 .ids
655 .get(key)
656 .ok_or_else(|| GraphError::KeyNotFound { key: key.into() })?;
657 let tv = build_tv(&state.topo, base);
658 let mut edges = Vec::new();
659 for etype in tv.etypes() {
660 let edge_type = state
661 .syms
662 .resolve(etype)
663 .ok_or_else(|| GraphError::Corrupt {
664 detail: format!("reader: topology etype {etype} not in interner"),
665 })?
666 .to_string();
667 for dir in [Direction::Out, Direction::In] {
668 for &nbr in tv.neighbors(etype, dir, id).as_ref() {
669 let (src_key, dst_key) = match dir {
670 Direction::Out => (
671 key.to_string(),
672 state
673 .ids
674 .key_of(nbr)
675 .ok_or_else(|| GraphError::Corrupt {
676 detail: format!("topology id {nbr} has no key"),
677 })?
678 .to_string(),
679 ),
680 Direction::In => (
681 state
682 .ids
683 .key_of(nbr)
684 .ok_or_else(|| GraphError::Corrupt {
685 detail: format!("topology id {nbr} has no key"),
686 })?
687 .to_string(),
688 key.to_string(),
689 ),
690 };
691 edges.push(EdgeInfo {
692 edge_type: edge_type.clone(),
693 src_key,
694 dst_key,
695 derived: false,
696 });
697 }
698 }
699 }
700 edges.sort_by(|a, b| {
701 a.edge_type
702 .cmp(&b.edge_type)
703 .then(a.src_key.cmp(&b.src_key))
704 .then(a.dst_key.cmp(&b.dst_key))
705 });
706 edges.dedup();
707 Ok(edges)
708}
709
710fn neighborhood_masked_from(
711 key: &str,
712 state: &FrozenOverlay,
713 base: &Option<Arc<MappedBase>>,
714 depth: u32,
715 edge_types: Option<&[&str]>,
716 dir: Dir,
717 mask: &NodeMask,
718) -> Option<ResultSet> {
719 let start_id = state.ids.get(key)?;
720 let view = make_view(state, base, Some(&mask.visible));
721 let resolved: Option<Vec<u32>> = edge_types.map(|names| {
722 names
723 .iter()
724 .filter_map(|name| view.syms.get(name))
725 .collect()
726 });
727 let nb = neighborhood(&view, start_id, depth, resolved.as_deref(), dir);
728 let mut rs = ResultSet::new(vec!["key".into(), "label".into(), "depth".into()]);
729 let mut visited: Vec<(u32, u32)> = Vec::with_capacity(nb.nodes.len() + 1);
731 visited.push((start_id, 0));
732 for (nid, d) in &nb.nodes {
733 let k = view.key_of(*nid);
734 let lbl = view
735 .label_of(*nid)
736 .expect("real nodes always have a label; u32::MAX sentinel cannot occur");
737 rs.push_row(vec![
738 Some(Value::Str(k.to_string())),
739 Some(Value::Str(lbl.to_string())),
740 Some(Value::Int(*d as i64)),
741 ]);
742 visited.push((*nid, *d));
743 }
744 if mask.mode() == crate::mask::MaskMode::Stub {
752 let raw_view = make_view(state, base, None);
753 let mut seen: HashSet<u32> = visited.iter().map(|(id, _)| *id).collect();
754 for (node_id, node_depth) in &visited {
755 if *node_depth >= depth {
756 continue;
757 }
758 for e in expand(&raw_view, *node_id, resolved.as_deref(), dir) {
759 let nbr = if e.src == *node_id { e.dst } else { e.src };
760 if !mask.contains_id(nbr) && seen.insert(nbr) {
761 if let Some(k) = state.ids.key_of(nbr) {
762 rs.push_row(vec![
763 Some(Value::Str(k.to_string())),
764 None,
765 Some(Value::Int((*node_depth + 1) as i64)),
766 ]);
767 }
768 }
769 }
770 }
771 }
772 Some(rs)
773}