use core_query::visible::VisibleSet;
use core_storage::fs::Fs;
use core_storage::Result;
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use crate::db::GraphDb;
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
pub enum MaskMode {
#[default]
Omit,
Stub,
}
#[derive(Clone, Debug)]
pub struct NodeMask {
pub(crate) visible: VisibleSet,
mode: MaskMode,
}
impl NodeMask {
pub fn from_keys<'a, F: Fs>(db: &GraphDb<F>, keys: impl IntoIterator<Item = &'a str>) -> Self {
let visible = keys.into_iter().filter_map(|k| db.ids().get(k)).collect();
NodeMask {
visible,
mode: MaskMode::default(),
}
}
pub fn from_ids(ids: impl IntoIterator<Item = u32>) -> Self {
NodeMask {
visible: ids.into_iter().collect(),
mode: MaskMode::default(),
}
}
pub fn with_mode(self, mode: MaskMode) -> Self {
NodeMask { mode, ..self }
}
pub fn mode(&self) -> MaskMode {
self.mode
}
pub fn len(&self) -> usize {
self.visible.len()
}
pub fn is_empty(&self) -> bool {
self.visible.is_empty()
}
pub fn intersect(&self, other: &NodeMask) -> NodeMask {
NodeMask {
visible: self.visible.intersect(&other.visible),
mode: MaskMode::Omit,
}
}
#[inline]
pub fn contains_id(&self, id: u32) -> bool {
self.visible.contains(id)
}
pub fn contains_node<F: core_storage::fs::Fs>(
&self,
db: &crate::db::GraphDb<F>,
key: &str,
) -> bool {
db.ids()
.get(key)
.is_some_and(|id| self.visible.contains(id))
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(crate) struct StoreId(u64);
impl StoreId {
pub(crate) fn next() -> StoreId {
static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
StoreId(NEXT.fetch_add(1, std::sync::atomic::Ordering::Relaxed))
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(crate) struct StoreStamp {
store: StoreId,
commit_seq: u64,
}
impl StoreStamp {
fn of<F: Fs>(db: &GraphDb<F>) -> StoreStamp {
StoreStamp {
store: db.store_id(),
commit_seq: db.commit_seq(),
}
}
}
#[cfg(test)]
thread_local! {
static SCOPE_KEYS_RESOLVES: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
}
pub struct Scope {
roles: Vec<String>,
namespaces: Vec<String>,
keys: Option<Vec<String>>,
keys_cache: Mutex<Option<(StoreStamp, NodeMask)>>,
}
impl Scope {
pub fn new(
role: Option<String>,
namespace: Option<String>,
keys: Option<Vec<String>>,
) -> Result<Scope> {
if role.is_none() && namespace.is_none() && keys.is_none() {
return Err(core_storage::GraphError::QueryError {
detail: "a scope needs at least one of role, namespace or keys; \
an empty scope is refused rather than read as unscoped"
.into(),
});
}
Ok(Scope {
roles: role.into_iter().collect(),
namespaces: namespace.into_iter().collect(),
keys,
keys_cache: Mutex::new(None),
})
}
pub fn resolve<F: Fs>(&self, db: &GraphDb<F>) -> Result<NodeMask> {
self.resolve_with(db, true)
}
pub(crate) fn resolve_uncached<F: Fs>(&self, db: &GraphDb<F>) -> Result<NodeMask> {
self.resolve_with(db, false)
}
fn resolve_with<F: Fs>(&self, db: &GraphDb<F>, cached: bool) -> Result<NodeMask> {
let mut out: Option<NodeMask> = None;
let mut narrow = |mask: NodeMask| {
out = Some(match out.take() {
Some(acc) => acc.intersect(&mask),
None => mask,
});
};
for role in &self.roles {
narrow(db.mask_for_role(role)?);
}
for namespace in &self.namespaces {
narrow(db.mask_for_namespace(namespace));
}
if self.keys.is_some() {
narrow(self.resolve_keys(db, cached));
}
Ok(out
.expect("a Scope always has at least one leg")
.with_mode(MaskMode::Omit))
}
pub fn intersect(&self, other: &Scope) -> Scope {
let keys = match (&self.keys, &other.keys) {
(Some(a), Some(b)) => {
let b: HashSet<&str> = b.iter().map(String::as_str).collect();
Some(
a.iter()
.filter(|k| b.contains(k.as_str()))
.cloned()
.collect(),
)
}
(Some(a), None) => Some(a.clone()),
(None, b) => b.clone(),
};
Scope {
roles: [self.roles.clone(), other.roles.clone()].concat(),
namespaces: [self.namespaces.clone(), other.namespaces.clone()].concat(),
keys,
keys_cache: Mutex::new(None),
}
}
fn resolve_keys<F: Fs>(&self, db: &GraphDb<F>, cached: bool) -> NodeMask {
let keys = self.keys.as_deref().unwrap_or_default();
let stamp = StoreStamp::of(db);
if cached {
if let Ok(cache) = self.keys_cache.lock() {
if let Some((at, mask)) = cache.as_ref() {
if *at == stamp {
return mask.clone();
}
}
}
}
#[cfg(test)]
SCOPE_KEYS_RESOLVES.with(|c| c.set(c.get() + 1));
let mask = NodeMask::from_keys(db, keys.iter().map(String::as_str));
if cached {
if let Ok(mut cache) = self.keys_cache.lock() {
*cache = Some((stamp, mask.clone()));
}
}
mask
}
}
impl Clone for Scope {
fn clone(&self) -> Scope {
Scope {
roles: self.roles.clone(),
namespaces: self.namespaces.clone(),
keys: self.keys.clone(),
keys_cache: Mutex::new(self.keys_cache.lock().ok().and_then(|cache| cache.clone())),
}
}
}
impl std::fmt::Debug for Scope {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Scope")
.field("roles", &self.roles)
.field("namespaces", &self.namespaces)
.field("keys", &self.keys.as_ref().map(Vec::len))
.finish()
}
}
#[derive(Default)]
pub struct RoleMaskCache {
entries: Mutex<HashMap<String, (u64, Arc<NodeMask>)>>,
}
impl RoleMaskCache {
pub fn new() -> Self {
Self::default()
}
pub fn get_or_build(
&self,
role: &str,
version: u64,
build: impl FnOnce() -> Result<NodeMask>,
) -> Result<Arc<NodeMask>> {
if let Ok(entries) = self.entries.lock() {
if let Some((v, mask)) = entries.get(role) {
if *v == version {
return Ok(Arc::clone(mask));
}
}
}
let mask = Arc::new(build()?);
if let Ok(mut entries) = self.entries.lock() {
entries.insert(role.to_string(), (version, Arc::clone(&mask)));
}
Ok(mask)
}
pub fn clear(&self) {
if let Ok(mut entries) = self.entries.lock() {
entries.clear();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::roles::RoleDef;
use crate::schema::Schema;
fn tmp_dir(name: &str) -> std::path::PathBuf {
let d =
std::env::temp_dir().join(format!("graphdb-mask-unit-{}-{}", name, std::process::id()));
let _ = std::fs::remove_dir_all(&d);
d
}
#[test]
fn scope_with_no_legs_is_refused() {
let err = Scope::new(None, None, None).expect_err("an empty scope must not be built");
match err {
core_storage::GraphError::QueryError { detail } => {
for arg in ["role", "namespace", "keys"] {
assert!(
detail.contains(arg),
"the refusal must name `{arg}`; got {detail:?}"
);
}
}
other => panic!("expected QueryError, got {other:?}"),
}
}
#[test]
fn scope_legs_intersect_and_never_widen() {
let dir = tmp_dir("scope-intersect");
let mut db = GraphDb::open(&dir).unwrap();
db.insert_node("Doc", "a", vec![]).unwrap();
db.insert_node("Doc", "b", vec![]).unwrap();
db.insert_node("Secret", "c", vec![]).unwrap();
db.apply_schema(&Schema {
roles: vec![RoleDef {
name: "reader".into(),
keys: vec![],
labels: vec!["Doc".into()],
visible_where: None,
namespaces: None,
write: None,
}],
..Default::default()
})
.unwrap();
let scope = Scope::new(
Some("reader".into()),
None,
Some(vec!["b".into(), "c".into()]),
)
.unwrap();
let mask = scope.resolve(&db).unwrap();
let b = db.ids().get("b").unwrap();
assert_eq!(mask.len(), 1, "only `b` is in both legs");
assert!(mask.contains_id(b));
assert_eq!(mask.mode(), MaskMode::Omit);
}
#[test]
fn scope_keys_leg_sees_a_key_created_after_construction() {
let dir = tmp_dir("scope-late-key");
let mut db = GraphDb::open(&dir).unwrap();
db.insert_node("Doc", "early", vec![]).unwrap();
let scope = Scope::new(None, None, Some(vec!["late".into()])).unwrap();
assert!(
scope.resolve(&db).unwrap().is_empty(),
"`late` does not exist yet"
);
db.insert_node("Doc", "late", vec![]).unwrap();
let mask = scope.resolve(&db).unwrap();
let late = db.ids().get("late").unwrap();
assert!(
mask.contains_id(late),
"the key leg must be re-resolved after the write"
);
assert_eq!(mask.len(), 1);
}
#[test]
fn scope_keys_leg_is_cached_within_a_commit() {
let dir = tmp_dir("scope-keys-cache");
let mut db = GraphDb::open(&dir).unwrap();
db.insert_node("Doc", "a", vec![]).unwrap();
let scope = Scope::new(None, None, Some(vec!["a".into()])).unwrap();
let before = SCOPE_KEYS_RESOLVES.with(|c| c.get());
scope.resolve(&db).unwrap();
scope.resolve(&db).unwrap();
assert_eq!(
SCOPE_KEYS_RESOLVES.with(|c| c.get()) - before,
1,
"two resolves at one commit rebuild the key leg once"
);
db.insert_node("Doc", "b", vec![]).unwrap();
scope.resolve(&db).unwrap();
assert_eq!(
SCOPE_KEYS_RESOLVES.with(|c| c.get()) - before,
2,
"a write invalidates the cached key leg"
);
}
#[test]
fn scope_resolve_uncached_neither_reads_nor_fills_the_cache() {
let dir = tmp_dir("scope-uncached");
let mut db = GraphDb::open(&dir).unwrap();
db.insert_node("Doc", "a", vec![]).unwrap();
let scope = Scope::new(None, None, Some(vec!["a".into()])).unwrap();
let before = SCOPE_KEYS_RESOLVES.with(|c| c.get());
scope.resolve_uncached(&db).unwrap();
scope.resolve_uncached(&db).unwrap();
assert_eq!(
SCOPE_KEYS_RESOLVES.with(|c| c.get()) - before,
2,
"an uncached resolve never serves the cached entry"
);
scope.resolve(&db).unwrap();
scope.resolve(&db).unwrap();
assert_eq!(
SCOPE_KEYS_RESOLVES.with(|c| c.get()) - before,
3,
"and never fills it either: the first cached resolve still rebuilds"
);
}
#[test]
fn scope_keys_cache_never_crosses_stores() {
let dir_a = tmp_dir("scope-store-a");
let dir_b = tmp_dir("scope-store-b");
let mut a = GraphDb::open(&dir_a).unwrap();
a.insert_node("Doc", "filler", vec![]).unwrap();
a.insert_node("Doc", "target", vec![]).unwrap();
let mut b = GraphDb::open(&dir_b).unwrap();
b.insert_node("Doc", "other", vec![]).unwrap();
b.insert_node("Doc", "secret", vec![]).unwrap();
assert_eq!(
a.commit_seq(),
b.commit_seq(),
"the two stores must collide on commit_seq for this to test anything"
);
let scope = Scope::new(None, None, Some(vec!["target".into()])).unwrap();
let mask_a = scope.resolve(&a).unwrap();
assert!(mask_a.contains_id(a.ids().get("target").unwrap()));
let mask_b = scope.resolve(&b).unwrap();
for key in ["other", "secret"] {
let id = b.ids().get(key).unwrap();
assert!(
!mask_b.contains_id(id),
"store B's mask admits `{key}`, a node this scope never named"
);
}
assert!(
mask_b.is_empty(),
"store B has no `target`, so the key leg resolves to nothing there"
);
}
#[test]
fn scope_keys_cache_is_dropped_when_the_store_reloads() {
let dir = tmp_dir("scope-reload");
let mut w = GraphDb::open(&dir).unwrap();
w.insert_node("Doc", "filler", vec![]).unwrap();
w.insert_node("Doc", "target", vec![]).unwrap();
w.insert_node("Doc", "keep", vec![]).unwrap();
let mut r = GraphDb::open_with_options(
&dir,
crate::db::OpenOptions {
read_only: true,
..Default::default()
},
)
.unwrap();
let seq_before = r.commit_seq();
let scope = Scope::new(None, None, Some(vec!["target".into()])).unwrap();
assert_eq!(
scope.resolve(&r).unwrap().len(),
1,
"`target` is visible before the delete"
);
w.delete_node("target").unwrap();
w.snapshot().unwrap();
r.refresh().unwrap();
assert_eq!(
r.commit_seq(),
seq_before,
"the reseed must land back on the sequence the mask was cached at"
);
let mask = scope.resolve(&r).unwrap();
for key in ["filler", "keep"] {
let id = r.ids().get(key).unwrap();
assert!(
!mask.contains_id(id),
"after the reload the stale mask admits `{key}`, which the scope never named"
);
}
assert!(
mask.is_empty(),
"`target` is gone, so the key leg must resolve to nothing"
);
}
fn intersect_fixture(name: &str) -> (std::path::PathBuf, GraphDb<core_storage::fs::RealFs>) {
let dir = tmp_dir(name);
let mut db = GraphDb::open(&dir).unwrap();
db.insert_node("Doc", "a", vec![]).unwrap();
db.insert_node("Doc", "b", vec![]).unwrap();
db.insert_node("Secret", "c", vec![]).unwrap();
db.apply_schema(&Schema {
roles: vec![RoleDef {
name: "reader".into(),
keys: vec![],
labels: vec!["Doc".into()],
visible_where: None,
namespaces: None,
write: None,
}],
..Default::default()
})
.unwrap();
(dir, db)
}
fn visible_keys<F: Fs>(db: &GraphDb<F>, mask: &NodeMask, keys: &[&str]) -> Vec<String> {
keys.iter()
.filter(|k| mask.contains_node(db, k))
.map(|k| (*k).to_string())
.collect()
}
#[test]
fn intersect_of_two_key_legs_keeps_only_the_keys_in_both() {
let (_dir, db) = intersect_fixture("intersect-both-keys");
let left = Scope::new(None, None, Some(vec!["a".into(), "b".into()])).unwrap();
let right = Scope::new(None, None, Some(vec!["b".into(), "c".into()])).unwrap();
let mask = left.intersect(&right).resolve(&db).unwrap();
assert_eq!(visible_keys(&db, &mask, &["a", "b", "c"]), vec!["b"]);
let mask = right.intersect(&left).resolve(&db).unwrap();
assert_eq!(visible_keys(&db, &mask, &["a", "b", "c"]), vec!["b"]);
}
#[test]
fn intersect_carries_a_lone_key_leg_from_either_side() {
let (_dir, db) = intersect_fixture("intersect-one-key");
let keyed = Scope::new(None, None, Some(vec!["b".into()])).unwrap();
let roled = Scope::new(Some("reader".into()), None, None).unwrap();
let mask = keyed.intersect(&roled).resolve(&db).unwrap();
assert_eq!(visible_keys(&db, &mask, &["a", "b", "c"]), vec!["b"]);
let mask = roled.intersect(&keyed).resolve(&db).unwrap();
assert_eq!(visible_keys(&db, &mask, &["a", "b", "c"]), vec!["b"]);
}
#[test]
fn intersect_of_two_keyless_scopes_has_no_key_leg() {
let (_dir, db) = intersect_fixture("intersect-no-keys");
let roled = Scope::new(Some("reader".into()), None, None).unwrap();
let namespaced = Scope::new(None, Some("default".into()), None).unwrap();
let both = roled.intersect(&namespaced);
assert!(
both.keys.is_none(),
"an absent key leg must stay absent, not become `Some(vec![])`"
);
let mask = both.resolve(&db).unwrap();
assert_eq!(
visible_keys(&db, &mask, &["a", "b", "c"]),
vec!["a", "b"],
"the role leg still decides; the missing key leg narrows nothing"
);
}
#[test]
fn intersect_accumulates_role_and_namespace_legs() {
let (_dir, db) = intersect_fixture("intersect-legs");
let reader = Scope::new(Some("reader".into()), None, None).unwrap();
let elsewhere = Scope::new(None, Some("other".into()), None).unwrap();
let both = reader.intersect(&elsewhere);
assert_eq!(both.roles, vec!["reader".to_string()]);
assert_eq!(both.namespaces, vec!["other".to_string()]);
let mask = both.resolve(&db).unwrap();
assert!(
mask.is_empty(),
"every node is in the `default` namespace, so the two legs share nobody"
);
}
#[test]
fn intersect_can_never_widen_either_side() {
let (_dir, db) = intersect_fixture("intersect-never-widens");
let all = ["a", "b", "c"];
let scopes = || {
vec![
Scope::new(Some("reader".into()), None, None).unwrap(),
Scope::new(None, Some("default".into()), None).unwrap(),
Scope::new(None, None, Some(vec!["b".into(), "c".into()])).unwrap(),
Scope::new(None, None, Some(vec![])).unwrap(),
Scope::new(Some("reader".into()), None, Some(vec!["a".into()])).unwrap(),
]
};
for left in scopes() {
for right in scopes() {
let l = visible_keys(&db, &left.resolve(&db).unwrap(), &all);
let r = visible_keys(&db, &right.resolve(&db).unwrap(), &all);
let both = visible_keys(&db, &left.intersect(&right).resolve(&db).unwrap(), &all);
for key in &both {
assert!(
l.contains(key) && r.contains(key),
"{left:?} ∩ {right:?} sees `{key}`, which one side alone does not"
);
}
}
}
}
#[test]
fn a_version_change_rebuilds_and_clear_empties() {
let cache = RoleMaskCache::new();
let built = std::cell::Cell::new(0u32);
let build = |ids: Vec<u32>| {
built.set(built.get() + 1);
Ok(NodeMask::from_ids(ids))
};
let m = cache.get_or_build("r", 1, || build(vec![1])).unwrap();
assert_eq!(m.len(), 1);
assert_eq!(built.get(), 1);
let m = cache.get_or_build("r", 1, || build(vec![1, 2])).unwrap();
assert_eq!(m.len(), 1, "the memoised mask is returned unchanged");
assert_eq!(built.get(), 1);
let m = cache.get_or_build("r", 2, || build(vec![1, 2])).unwrap();
assert_eq!(m.len(), 2);
assert_eq!(built.get(), 2);
let m = cache.get_or_build("other", 2, || build(vec![9])).unwrap();
assert_eq!(m.len(), 1);
assert_eq!(built.get(), 3);
cache.clear();
let _ = cache.get_or_build("r", 2, || build(vec![1, 2])).unwrap();
assert_eq!(built.get(), 4, "clear drops the entry, so it rebuilds");
}
#[test]
fn a_failed_build_is_not_cached() {
let cache = RoleMaskCache::new();
assert!(cache
.get_or_build("r", 1, || Err(core_storage::GraphError::KeyNotFound {
key: "role:r".into()
}))
.is_err());
let m = cache
.get_or_build("r", 1, || Ok(NodeMask::from_ids(vec![7])))
.unwrap();
assert_eq!(m.len(), 1);
}
}