use std::collections::{BTreeMap, HashSet};
use std::sync::{Arc, OnceLock};
use core_query::cypher::{execute, is_write_tokens, lex, parse, plan, Params};
use core_query::{expand, neighborhood, Dir, GraphView, ResultSet};
use core_storage::fulltext::FulltextIndex;
use core_storage::v8::seam::{ColumnsView, TopologyView};
use core_storage::v8::MappedBase;
use core_storage::wal::WalRecord;
use core_storage::{
ColumnStore, Direction, EdgeProps, EdgePropsView, GraphError, IdMap, Interner, Result,
Topology, Value,
};
use crate::db::{EdgeInfo, NodeInfo};
use crate::mask::{NodeMask, RoleMaskCache};
use crate::roles::RoleDef;
pub const FOLD_EVERY_K: usize = 64;
pub struct CommitDelta {
pub records: Vec<WalRecord>,
pub derived_inserts: Vec<(u32, u32, u32)>,
pub derived_deletes: Vec<(u32, u32, u32)>,
}
#[derive(Clone)]
pub struct FrozenOverlay {
pub ids: IdMap,
pub syms: Interner,
pub topo: Topology,
pub props: ColumnStore,
pub labels: Vec<u32>,
pub edge_props: EdgeProps,
pub roles: Option<Vec<RoleDef>>,
pub fulltext: FulltextIndex,
}
pub struct ReaderSnapshot {
pub frozen: Arc<FrozenOverlay>,
pub base: Option<Arc<MappedBase>>,
pub deltas: Vec<Arc<CommitDelta>>,
pub version: u64,
role_masks: Arc<RoleMaskCache>,
cache: OnceLock<std::result::Result<FrozenOverlay, String>>,
}
fn build_tv<'a>(topo: &'a Topology, base: &'a Option<Arc<MappedBase>>) -> TopologyView<'a> {
match base {
None => TopologyView::owned(topo),
Some(b) => {
let csr = b.topology().expect("base CSR CRC already verified at open");
TopologyView::with_base(topo, csr)
}
}
}
fn build_cv<'a>(props: &'a ColumnStore, base: &'a Option<Arc<MappedBase>>) -> ColumnsView<'a> {
match base {
None => ColumnsView::owned(props),
Some(b) => {
let cols = b
.columns()
.expect("base columns CRC already verified at open");
let strings = b
.string_table()
.transpose()
.expect("base strings CRC already verified at open");
ColumnsView::with_base_cached(props, cols, b.mixed_cache()).with_shared_strings(strings)
}
}
}
fn build_epv<'a>(
edge_props: &'a EdgeProps,
base: &'a Option<Arc<MappedBase>>,
) -> EdgePropsView<'a> {
match base {
None => EdgePropsView::owned(edge_props),
Some(b) => {
let archived = b
.edge_props_section()
.expect("base edge_props CRC already verified at open");
EdgePropsView::with_base(edge_props, archived)
}
}
}
fn make_view<'a>(
state: &'a FrozenOverlay,
base: &'a Option<Arc<MappedBase>>,
mask: Option<&'a HashSet<u32>>,
) -> GraphView<'a> {
GraphView {
ids: &state.ids,
syms: &state.syms,
labels: &state.labels,
props: build_cv(&state.props, base),
topo: build_tv(&state.topo, base),
edge_props: build_epv(&state.edge_props, base),
mask,
prop_index: None,
}
}
#[allow(clippy::too_many_arguments)]
fn apply_one(
ids: &mut IdMap,
syms: &mut Interner,
topo: &mut Topology,
props: &mut ColumnStore,
edge_props: &mut EdgeProps,
labels: &mut Vec<u32>,
fulltext: &mut FulltextIndex,
rec: &WalRecord,
) -> Result<()> {
match rec {
WalRecord::Intern { id, text } => {
let got = syms.intern(text);
if got != *id {
return Err(GraphError::Corrupt {
detail: format!(
"mvcc delta intern mismatch for {text:?}: expected {id}, got {got}"
),
});
}
}
WalRecord::InsertNodeId {
label,
key,
props: node_props,
} => {
let node_id = ids.try_insert(key)?;
if labels.len() <= node_id as usize {
labels.resize(node_id as usize + 1, u32::MAX);
}
labels[node_id as usize] = *label;
let label_str = syms
.resolve(*label)
.ok_or_else(|| GraphError::Corrupt {
detail: format!("mvcc delta: unknown label sym {label}"),
})?
.to_string();
for (field_sym, value) in node_props {
let field = syms
.resolve(*field_sym)
.ok_or_else(|| GraphError::Corrupt {
detail: format!("mvcc delta: unknown field sym {field_sym}"),
})?
.to_string();
props.set(node_id, &field, value.clone());
if fulltext.is_enabled(&label_str, &field) {
fulltext.add_tokens(node_id, &field, value);
}
}
}
WalRecord::InsertNode {
label,
key,
props: node_props,
} => {
let label_sym = syms.intern(label);
let node_id = ids.try_insert(key)?;
if labels.len() <= node_id as usize {
labels.resize(node_id as usize + 1, u32::MAX);
}
labels[node_id as usize] = label_sym;
for (field, value) in node_props {
props.set(node_id, field, value.clone());
if fulltext.is_enabled(label, field) {
fulltext.add_tokens(node_id, field, value);
}
}
}
WalRecord::SetPropId { id, field, value } => {
if let Some(field_str) = syms.resolve(*field).map(str::to_string) {
props.set(*id, &field_str, value.clone());
if let Some(&label_sym) = labels.get(*id as usize) {
if let Some(label_str) = syms.resolve(label_sym) {
if fulltext.is_enabled(label_str, &field_str) {
fulltext.add_tokens(*id, &field_str, value);
}
}
}
}
}
WalRecord::SetProp { key, field, value } => {
if let Some(node_id) = ids.get(key) {
props.set(node_id, field, value.clone());
if let Some(&label_sym) = labels.get(node_id as usize) {
if let Some(label_str) = syms.resolve(label_sym) {
if fulltext.is_enabled(label_str, field) {
fulltext.add_tokens(node_id, field, value);
}
}
}
}
}
WalRecord::RemoveProp { key, field } => {
if let Some(node_id) = ids.get(key) {
props.remove(node_id, field);
props.record_prop_tombstone(node_id, field);
fulltext.remove_node_field(node_id, field);
}
}
WalRecord::DeleteNode { key } => {
if let Some(node_id) = ids.delete(key) {
props.remove_all(node_id);
fulltext.remove_node(node_id);
if let Some(slot) = labels.get_mut(node_id as usize) {
*slot = u32::MAX;
}
let etypes: Vec<u32> = topo.etypes().collect();
let mut doomed = Vec::new();
for et in &etypes {
for &dst in topo.neighbors(*et, Direction::Out, node_id).as_ref() {
doomed.push((*et, node_id, dst));
}
for &src in topo.neighbors(*et, Direction::In, node_id).as_ref() {
doomed.push((*et, src, node_id));
}
}
for (et, s, d) in doomed {
topo.remove_edge(et, s, d);
edge_props.remove_edge(et, s, d);
}
}
}
WalRecord::InsertEdgeId { etype, src, dst } => {
topo.add_edge(*etype, *src, *dst);
}
WalRecord::InsertEdge {
edge_type,
src_key,
dst_key,
} => {
let etype = syms.intern(edge_type);
if let (Some(src), Some(dst)) = (ids.get(src_key), ids.get(dst_key)) {
topo.add_edge(etype, src, dst);
}
}
WalRecord::DeleteEdge {
edge_type,
src_key,
dst_key,
} => {
if let Some(etype) = syms.get(edge_type) {
if let (Some(src), Some(dst)) = (ids.get(src_key), ids.get(dst_key)) {
topo.remove_edge(etype, src, dst);
}
}
}
WalRecord::EnableFulltext { label, field } => {
fulltext.enable(label, field);
}
WalRecord::DisableFulltext { label, field } => {
fulltext.disable(label, field);
}
WalRecord::Batch(inner) => {
for r in inner {
apply_one(ids, syms, topo, props, edge_props, labels, fulltext, r)?;
}
}
WalRecord::CreateRule { .. }
| WalRecord::DeleteRule { .. }
| WalRecord::RebuildRule { .. }
| WalRecord::CreateView { .. }
| WalRecord::DeleteView { .. }
| WalRecord::EnableIndex { .. }
| WalRecord::DisableIndex { .. }
| WalRecord::DerivedEdgeAdded { .. }
| WalRecord::DerivedEdgeRetracted { .. } => {}
WalRecord::RenameNode { old_key, new_key } => {
if ids.get(old_key).is_some() {
ids.rename(old_key, new_key).map_err(|_| GraphError::Corrupt {
detail: format!("mvcc delta RenameNode {old_key}→{new_key} failed"),
})?;
}
}
}
Ok(())
}
fn mask_for_role_from(
state: &FrozenOverlay,
base: &Option<Arc<MappedBase>>,
role: &str,
) -> Result<NodeMask> {
let roles = state.roles.as_ref().ok_or_else(|| GraphError::Corrupt {
detail: "roles.json was corrupt at open; fix the file and re-open".into(),
})?;
let def = roles
.iter()
.find(|r| r.name == role)
.ok_or_else(|| GraphError::KeyNotFound {
key: format!("role:{role}"),
})?;
let mut visible = HashSet::new();
for key in &def.keys {
if let Some(id) = state.ids.get(key) {
visible.insert(id);
}
}
let props = def
.visible_where
.as_ref()
.map(|_| build_cv(&state.props, base));
for label_name in &def.labels {
if let Some(sym) = state.syms.get(label_name) {
for (i, &s) in state.labels.iter().enumerate() {
if s != sym {
continue;
}
let id = i as u32;
match (&def.visible_where, &props) {
(Some(pred), Some(view)) => {
let value = view.get(id, &pred.field).map(|vr| vr.into_value());
if pred.holds(value.as_ref()) {
visible.insert(id);
}
}
_ => {
visible.insert(id);
}
}
}
}
}
if def.namespaces.is_some() {
let cv = build_cv(&state.props, base);
visible.retain(|&id| {
let value = cv.get(id, core_storage::NS_PROP).map(|vr| vr.into_value());
def.sees_namespace(core_storage::namespace_of_value(value.as_ref()))
});
}
Ok(NodeMask::from_ids(visible))
}
impl ReaderSnapshot {
fn materialize(&self) -> Result<FrozenOverlay> {
let mut w = (*self.frozen).clone();
for delta in &self.deltas {
for rec in &delta.records {
apply_one(
&mut w.ids,
&mut w.syms,
&mut w.topo,
&mut w.props,
&mut w.edge_props,
&mut w.labels,
&mut w.fulltext,
rec,
)?;
}
for &(etype, src, dst) in &delta.derived_inserts {
w.topo.add_edge(etype, src, dst);
}
for &(etype, src, dst) in &delta.derived_deletes {
w.topo.remove_edge(etype, src, dst);
}
}
if !self.deltas.is_empty() {
let cv = build_cv(&w.props, &self.base);
w.fulltext.rebuild_all(&w.ids, &w.labels, &w.syms, cv);
}
Ok(w)
}
pub(crate) fn new(
frozen: Arc<FrozenOverlay>,
base: Option<Arc<MappedBase>>,
deltas: Vec<Arc<CommitDelta>>,
version: u64,
role_masks: Arc<RoleMaskCache>,
) -> Self {
Self {
frozen,
base,
deltas,
version,
role_masks,
cache: OnceLock::new(),
}
}
fn effective(&self) -> Result<&FrozenOverlay> {
if self.deltas.is_empty() {
return Ok(&self.frozen);
}
let cached = self
.cache
.get_or_init(|| self.materialize().map_err(|e| e.to_string()));
cached
.as_ref()
.map_err(|e| GraphError::Corrupt { detail: e.clone() })
}
pub fn mask_for_role(&self, role: &str) -> Result<NodeMask> {
self.role_masks
.get_or_build(role, self.version, || {
mask_for_role_from(self.effective()?, &self.base, role)
})
.map(|m| (*m).clone())
}
pub fn mask_for_namespace(&self, namespace: &str) -> Result<NodeMask> {
let state = self.effective()?;
let cv = build_cv(&state.props, &self.base);
let mut visible = HashSet::new();
for (i, &sym) in state.labels.iter().enumerate() {
if sym == u32::MAX {
continue; }
let id = i as u32;
let value = cv.get(id, core_storage::NS_PROP).map(|vr| vr.into_value());
if core_storage::namespace_of_value(value.as_ref()) == namespace {
visible.insert(id);
}
}
Ok(NodeMask::from_ids(visible))
}
pub fn resolve_key(&self, key: &str) -> Option<u32> {
self.effective().ok()?.ids.get(key)
}
pub fn query(&self, cypher: &str, params: &BTreeMap<String, Value>) -> Result<ResultSet> {
let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
detail: format!("lex: {e}"),
})?;
let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
detail: format!("parse: {e}"),
})?;
let ops = plan(&ast).map_err(|e| GraphError::QueryError {
detail: format!("plan: {e}"),
})?;
let state = self.effective()?;
let view = make_view(state, &self.base, None);
execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
detail: format!("execute: {e}"),
})
}
pub fn query_masked(
&self,
cypher: &str,
params: &BTreeMap<String, Value>,
mask: &NodeMask,
) -> Result<ResultSet> {
let tokens = lex(cypher).map_err(|e| GraphError::QueryError {
detail: format!("lex: {e}"),
})?;
if is_write_tokens(&tokens) {
return Err(GraphError::QueryError {
detail: "masked queries are read-only".into(),
});
}
let ast = parse(&tokens).map_err(|e| GraphError::QueryError {
detail: format!("parse: {e}"),
})?;
let ops = plan(&ast).map_err(|e| GraphError::QueryError {
detail: format!("plan: {e}"),
})?;
let state = self.effective()?;
let view = make_view(state, &self.base, Some(&mask.visible));
execute(&view, &ops, &Params(params)).map_err(|e| GraphError::QueryError {
detail: format!("execute: {e}"),
})
}
pub fn node_info(&self, key: &str) -> Option<NodeInfo> {
node_info_from(key, self.effective().ok()?, &self.base)
}
pub fn node_edges(&self, key: &str) -> Result<Vec<EdgeInfo>> {
node_edges_from(key, self.effective()?, &self.base)
}
pub fn neighborhood_masked(
&self,
key: &str,
depth: u32,
edge_types: Option<&[&str]>,
dir: Dir,
mask: &NodeMask,
) -> Option<ResultSet> {
neighborhood_masked_from(
key,
self.effective().ok()?,
&self.base,
depth,
edge_types,
dir,
mask,
)
}
}
fn node_info_from(
key: &str,
state: &FrozenOverlay,
base: &Option<Arc<MappedBase>>,
) -> Option<NodeInfo> {
let id = state.ids.get(key)?;
let label_sym = *state.labels.get(id as usize)?;
if label_sym == u32::MAX {
return None;
}
let label = state.syms.resolve(label_sym)?.to_string();
let cv = build_cv(&state.props, base);
let mut props = BTreeMap::new();
for field in cv.field_names() {
if let Some(vr) = cv.get(id, &field) {
props.insert(field, vr.into_value());
}
}
Some(NodeInfo {
key: key.to_string(),
label,
props,
})
}
fn node_edges_from(
key: &str,
state: &FrozenOverlay,
base: &Option<Arc<MappedBase>>,
) -> Result<Vec<EdgeInfo>> {
let id = state
.ids
.get(key)
.ok_or_else(|| GraphError::KeyNotFound { key: key.into() })?;
let tv = build_tv(&state.topo, base);
let mut edges = Vec::new();
for etype in tv.etypes() {
let edge_type = state
.syms
.resolve(etype)
.ok_or_else(|| GraphError::Corrupt {
detail: format!("reader: topology etype {etype} not in interner"),
})?
.to_string();
for dir in [Direction::Out, Direction::In] {
for &nbr in tv.neighbors(etype, dir, id).as_ref() {
let (src_key, dst_key) = match dir {
Direction::Out => (
key.to_string(),
state
.ids
.key_of(nbr)
.ok_or_else(|| GraphError::Corrupt {
detail: format!("topology id {nbr} has no key"),
})?
.to_string(),
),
Direction::In => (
state
.ids
.key_of(nbr)
.ok_or_else(|| GraphError::Corrupt {
detail: format!("topology id {nbr} has no key"),
})?
.to_string(),
key.to_string(),
),
};
edges.push(EdgeInfo {
edge_type: edge_type.clone(),
src_key,
dst_key,
derived: false,
});
}
}
}
edges.sort_by(|a, b| {
a.edge_type
.cmp(&b.edge_type)
.then(a.src_key.cmp(&b.src_key))
.then(a.dst_key.cmp(&b.dst_key))
});
edges.dedup();
Ok(edges)
}
fn neighborhood_masked_from(
key: &str,
state: &FrozenOverlay,
base: &Option<Arc<MappedBase>>,
depth: u32,
edge_types: Option<&[&str]>,
dir: Dir,
mask: &NodeMask,
) -> Option<ResultSet> {
let start_id = state.ids.get(key)?;
let view = make_view(state, base, Some(&mask.visible));
let resolved: Option<Vec<u32>> = edge_types.map(|names| {
names
.iter()
.filter_map(|name| view.syms.get(name))
.collect()
});
let nb = neighborhood(&view, start_id, depth, resolved.as_deref(), dir);
let mut rs = ResultSet::new(vec!["key".into(), "label".into(), "depth".into()]);
let mut visited: Vec<(u32, u32)> = Vec::with_capacity(nb.nodes.len() + 1);
visited.push((start_id, 0));
for (nid, d) in &nb.nodes {
let k = view.key_of(*nid);
let lbl = view
.label_of(*nid)
.expect("real nodes always have a label; u32::MAX sentinel cannot occur");
rs.push_row(vec![
Some(Value::Str(k.to_string())),
Some(Value::Str(lbl.to_string())),
Some(Value::Int(*d as i64)),
]);
visited.push((*nid, *d));
}
if mask.mode() == crate::mask::MaskMode::Stub {
let raw_view = make_view(state, base, None);
let mut seen: HashSet<u32> = visited.iter().map(|(id, _)| *id).collect();
for (node_id, node_depth) in &visited {
if *node_depth >= depth {
continue;
}
for e in expand(&raw_view, *node_id, resolved.as_deref(), dir) {
let nbr = if e.src == *node_id { e.dst } else { e.src };
if !mask.contains_id(nbr) && seen.insert(nbr) {
if let Some(k) = state.ids.key_of(nbr) {
rs.push_row(vec![
Some(Value::Str(k.to_string())),
None,
Some(Value::Int((*node_depth + 1) as i64)),
]);
}
}
}
}
}
Some(rs)
}