From d974261945360e9f88853f95c3a649756b6daa88 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:08:09 +0000 Subject: [PATCH 1/6] lance-graph-dir-sim: directory desired-state simulation over SoA + Quack Observed snapshot as SoA lanes (Guid128-sorted id lane = ordinal, kind/active bit planes, dictionary-id value/key lanes, packed OU-HHTL lane, sorted membership relation). A version is the shared Arc plus a delta-sized overlay; nothing copies the population. Invariants are Quack programs: edge integrity = anti-join (negated Semijoin -> MaskOp::Gather) over the kind planes, SMTP/UPN uniqueness = GroupReduce Count on the normalized-key id, folded over base + overlay. Population rules use folded GROUP BY sinks. Semantic diff (delta path / reconcile merge path), plan with basis preconditions, provenance audit. One-edge mutation measured at 853 B allocated at both 1k and 100k users. Vocabulary from OGAR ogar-dir-sim. No AD/Graph/Exchange/LDAP/PowerShell I/O. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- Cargo.toml | 6 + crates/lance-graph-dir-sim/Cargo.toml | 29 + crates/lance-graph-dir-sim/src/exec.rs | 61 ++ crates/lance-graph-dir-sim/src/lib.rs | 78 +++ crates/lance-graph-dir-sim/src/observe.rs | 59 ++ crates/lance-graph-dir-sim/src/rule.rs | 149 +++++ crates/lance-graph-dir-sim/src/snapshot.rs | 343 +++++++++++ crates/lance-graph-dir-sim/src/store.rs | 342 +++++++++++ crates/lance-graph-dir-sim/src/validate.rs | 213 +++++++ crates/lance-graph-dir-sim/src/view.rs | 258 ++++++++ crates/lance-graph-dir-sim/tests/alloc.rs | 90 +++ crates/lance-graph-dir-sim/tests/sim.rs | 654 +++++++++++++++++++++ 12 files changed, 2282 insertions(+) create mode 100644 crates/lance-graph-dir-sim/Cargo.toml create mode 100644 crates/lance-graph-dir-sim/src/exec.rs create mode 100644 crates/lance-graph-dir-sim/src/lib.rs create mode 100644 crates/lance-graph-dir-sim/src/observe.rs create mode 100644 crates/lance-graph-dir-sim/src/rule.rs create mode 100644 crates/lance-graph-dir-sim/src/snapshot.rs create mode 100644 crates/lance-graph-dir-sim/src/store.rs create mode 100644 crates/lance-graph-dir-sim/src/validate.rs create mode 100644 crates/lance-graph-dir-sim/src/view.rs create mode 100644 crates/lance-graph-dir-sim/tests/alloc.rs create mode 100644 crates/lance-graph-dir-sim/tests/sim.rs diff --git a/Cargo.toml b/Cargo.toml index 119374056..9a9a6d82f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -110,6 +110,12 @@ exclude = [ # reason as lance-graph-ogar: path-deps the OGAR sibling. Verify via # `cargo test --manifest-path crates/lance-graph-report-ogar/Cargo.toml`. "crates/lance-graph-report-ogar", + # Directory desired-state simulation over SoA lanes + Quack (versions as + # shared snapshot + delta overlay; invariants as Semijoin / GroupReduce + # programs). EXCLUDED for the same reason: path-deps the OGAR sibling + # (ogar-dir-core / ogar-ad / ogar-dir-sim). Verify via + # `cargo test --manifest-path crates/lance-graph-dir-sim/Cargo.toml`. + "crates/lance-graph-dir-sim", # Cognitive Compilation loop (claude/cognitive-compilation-lance-graph-h8sgym). # The NEW piece is the Elixir-shaped template — thinking styles, JITson, and # i4-32D thinking-style vectors already exist; the template was the gap. Four diff --git a/crates/lance-graph-dir-sim/Cargo.toml b/crates/lance-graph-dir-sim/Cargo.toml new file mode 100644 index 000000000..5475c18ff --- /dev/null +++ b/crates/lance-graph-dir-sim/Cargo.toml @@ -0,0 +1,29 @@ +# lance-graph-dir-sim — simulate a directory's desired state over the SoA +# substrate. EXCLUDED from the lance-graph workspace (own [workspace] root) +# because it path-deps the OGAR sibling — the same shape as +# lance-graph-report-ogar / lance-graph-ogar. +# +# Direction: this crate depends on BOTH sides; neither depends on it. OGAR +# owns the semantic vocabulary (ogar-dir-sim: Change, provenance, Violation, +# ExecutionPlan) and the observation records (ogar-dir-core / ogar-ad); +# lance-graph owns execution: every population query lowers through +# lance-graph-quack onto lance-graph-mask-risc's one evaluator. +# +# Build/verify: cargo test --manifest-path crates/lance-graph-dir-sim/Cargo.toml + +[package] +name = "lance-graph-dir-sim" +version = "0.1.0" +edition = "2021" +publish = false +license = "Apache-2.0" +description = "Versioned directory graph over SoA lanes: observed snapshots shared structurally across simulated versions (parent + delta), pure population rules, invariants as Quack semijoin / GroupReduce programs, semantic diff and a side-effect-free ExecutionPlan. No AD / Graph / Exchange / LDAP / PowerShell I/O." + +[workspace] + +[dependencies] +lance-graph-quack = { path = "../lance-graph-quack" } +lance-graph-mask-risc = { path = "../lance-graph-mask-risc" } +ogar-dir-core = { path = "../../../OGAR/crates/ogar-dir-core" } +ogar-dir-sim = { path = "../../../OGAR/crates/ogar-dir-sim" } +ogar-ad = { path = "../../../OGAR/crates/ogar-ad" } diff --git a/crates/lance-graph-dir-sim/src/exec.rs b/crates/lance-graph-dir-sim/src/exec.rs new file mode 100644 index 000000000..099953691 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/exec.rs @@ -0,0 +1,61 @@ +//! The one place a Quack [`Query`] is lowered and run. +//! +//! Every population operation in this crate is a `Filter` + `Agg` over +//! borrowed lanes and resident planes, lowered by `lance-graph-quack` and +//! executed by `lance-graph-mask-risc`. Nothing here iterates rows, builds a +//! hash table, or produces joined tuples. The caller owns every buffer: +//! scratch is one tile per slot (`Scratch::for_program`), a `Keep` lands in a +//! `words_for(n_rows)` bitmap, a `GroupReduce` in a `K`-slot sink. + +use lance_graph_mask_risc::{ + execute_into, words_for, Foreign, Out, Planes, Program, Scratch, Terminal, Value, +}; +use lance_graph_quack::{lower, Agg, Col, Filter, GroupAddr, GroupAgg, Query}; + +/// Lower a directory query. Directory queries are fixed shapes built in this +/// crate, so a lowering failure is a programming error, not input. +pub(crate) fn program(filter: Filter, agg: Agg) -> Program { + lower(&Query { filter, agg }).expect("directory queries lower") +} + +fn run(p: &Program, planes: &Planes<'_>, foreign: &Foreign<'_>, out: Out<'_>) -> Value { + let mut scratch = Scratch::for_program(p, planes.n_rows).expect("scratch carves"); + execute_into(p, planes, foreign, &mut scratch, out).expect("directory program runs") +} + +/// Surviving rows as a `words_for(n_rows)` bitmap (`Agg::Rows` → `Keep`). +pub(crate) fn keep(p: &Program, planes: &Planes<'_>, foreign: &Foreign<'_>) -> Vec { + debug_assert!(matches!(p.terminal, Terminal::Keep { .. })); + let mut bits = vec![0u64; words_for(planes.n_rows)]; + if planes.n_rows > 0 { + run(p, planes, foreign, Out::Mask(&mut bits)); + } + bits +} + +/// `GROUP BY key COUNT(*)` over the rows `filter` keeps, ADDED into `sink` +/// (whose length is the group universe; keys past it — e.g. `NONE` — drop). +pub(crate) fn group_count_into( + filter: Filter, + key: Col, + planes: &Planes<'_>, + foreign: &Foreign<'_>, + sink: &mut [i64], +) { + if planes.n_rows == 0 || sink.is_empty() { + return; + } + let p = program( + filter, + Agg::GroupReduce { + key: GroupAddr::Local(key), + agg: GroupAgg::Count, + }, + ); + let mut part = vec![0i64; sink.len()]; + let v = run(&p, planes, foreign, Out::I64(&mut part)); + debug_assert_eq!(v, Value::GroupReduced); + for (s, x) in sink.iter_mut().zip(part) { + *s += x; + } +} diff --git a/crates/lance-graph-dir-sim/src/lib.rs b/crates/lance-graph-dir-sim/src/lib.rs new file mode 100644 index 000000000..986b270c6 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/lib.rs @@ -0,0 +1,78 @@ +//! # lance-graph-dir-sim — explore a directory future over the SoA substrate +//! +//! ```text +//! observed G0 ──rule──► G1 ──rule──► G2 ──validate──► "desired" ──diff(G0,G2)──► ExecutionPlan ──X +//! ``` +//! +//! OGAR owns the meaning (`ogar-dir-sim`: `Change`, provenance, `Violation`, +//! `ExecutionPlan`); this crate owns execution: +//! +//! * [`snapshot`] — one observation as SoA lanes and bit planes, `Guid128` +//! sorted so the ordinal is the index; strings in store dictionaries. +//! * [`view`] — a version = shared `Arc` + delta-sized overlay. +//! * [`rule`] — pure population rules over a borrowed [`View`]. +//! * [`validate`] — invariants as Quack programs (`Semijoin` anti-joins, +//! `GroupReduce` counts). +//! * [`store`] — append-only versions, tags, diff, plan, audit. +//! * [`observe`] — `ogar-ad` records → observation. +//! +//! No network, process or file I/O. Nothing writes to AD, Entra, Exchange, +//! LDAP or PowerShell. `#![forbid(unsafe_code)]`. + +#![forbid(unsafe_code)] + +mod exec; +pub mod observe; +pub mod rule; +pub mod snapshot; +pub mod store; +pub mod validate; +pub mod view; + +pub use rule::{member_counts, GrantGroup, ImplyGroup, Rule, SetPrimarySmtp}; +pub use snapshot::{ + pack_ou, BuildError, Dict, Dicts, NodeKind, Observation, ObservedNode, Snapshot, NONE, +}; +pub use store::{Rejection, SimError, VersionStore}; +pub use view::{ApplyError, View}; + +use lance_graph_mask_risc::{Foreign, LaneRef, Planes}; +use lance_graph_quack::{Cmp, Col, Filter, Mask}; +use ogar_dir_core::OuHhtl; + +/// A subtree prefix deeper than the packed lane (4 levels) can address. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct SubtreeTooDeep(pub usize); + +/// Nodes located in the OU subtree `prefix` (ancestor-or-self), as a node +/// bitmap — one ternary match on the packed OU lane (`Cmp::MatchU64`), no DN +/// strings. Prefixes up to depth 4 are exact; deeper ones are refused. +pub fn subtree(v: &View<'_>, prefix: &OuHhtl) -> Result, SubtreeTooDeep> { + let d = prefix.depth(); + if d > 4 { + return Err(SubtreeTooDeep(d)); + } + let care = if d == 0 { 0 } else { u64::MAX << (64 - 16 * d) }; + let s = v.snap; + let lanes = [LaneRef::U64(&s.ou_hi)]; + let masks: [&[u64]; 1] = [&s.ou_present]; + let planes = Planes { + n_rows: s.len(), + masks: &masks, + lanes: &lanes, + }; + let p = exec::program( + Filter::and([ + Filter::plane(Mask(0)), + Filter::cmp( + Col(0), + Cmp::MatchU64 { + pattern: pack_ou(prefix), + care, + }, + ), + ]), + lance_graph_quack::Agg::Rows, + ); + Ok(exec::keep(&p, &planes, &Foreign::NONE)) +} diff --git a/crates/lance-graph-dir-sim/src/observe.rs b/crates/lance-graph-dir-sim/src/observe.rs new file mode 100644 index 000000000..a9d831a59 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/observe.rs @@ -0,0 +1,59 @@ +//! `ogar-ad` records (OGAR PR #313) → an [`Observation`]. +//! +//! The ingestion boundary: values are read out of the record's value pool +//! once and handed to [`Snapshot::build`](crate::Snapshot::build) for +//! interning. "Active" is derived from `userAccountControl` bit `0x2` +//! (ACCOUNTDISABLE); the primary SMTP is the `SMTP:` proxy; the OU-HHTL is +//! taken as-is (never a DN string). Memberships are relations that `ogar-ad` +//! records do not carry; the caller adds observed ones. + +use crate::snapshot::{NodeKind, Observation, ObservedNode}; +use ogar_ad::{AdKind, SCHEMA_V1}; +use ogar_dir_core::{DirRecord, ValuePool}; + +const UAC_ACCOUNTDISABLE: u32 = 0x2; + +fn slot(name: &str) -> usize { + SCHEMA_V1 + .iter() + .find(|d| d.name == name) + .map(|d| d.slot as usize) + .expect("ogar-ad schema v1") +} + +/// Users and groups among `records` (other kinds are skipped). +pub fn from_ad(records: &[DirRecord], pool: &ValuePool) -> Observation { + let text = |r: &DirRecord, name: &str| { + r.str_ref(slot(name)) + .and_then(|s| pool.get(s)) + .and_then(|b| std::str::from_utf8(b).ok()) + .map(str::to_string) + }; + let mut obs = Observation::default(); + for r in records { + let kind = match r.object_kind() { + k if k == AdKind::User as u16 => NodeKind::User, + k if k == AdKind::Group as u16 => NodeKind::Group, + _ => continue, + }; + let primary_smtp = r + .str_ref(slot("proxyAddresses")) + .and_then(|s| pool.get_multi(s)) + .and_then(|vs| { + vs.into_iter() + .filter_map(|v| std::str::from_utf8(v).ok()) + .find_map(|v| v.strip_prefix("SMTP:").map(str::to_string)) + }); + obs.nodes.push(( + r.node_guid(), + ObservedNode { + kind, + active: r.num(0).is_none_or(|uac| uac & UAC_ACCOUNTDISABLE == 0), + upn: text(r, "userPrincipalName"), + primary_smtp, + ou: r.ou_hhtl(), + }, + )); + } + obs +} diff --git a/crates/lance-graph-dir-sim/src/rule.rs b/crates/lance-graph-dir-sim/src/rule.rs new file mode 100644 index 000000000..745addcc3 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/rule.rs @@ -0,0 +1,149 @@ +//! Pure rules: `G(n+1) = R(G(n), evidence)`. +//! +//! A rule receives a [`View`] — borrowed lanes of one version — and returns +//! the [`Change`]s it proposes. It has no I/O handle and no way to mutate the +//! view. Rules select **populations**: a population is a bitmap over the +//! version's node ordinals, produced by Quack programs, never a loop that +//! runs a workflow per user. + +use crate::exec::group_count_into; +use crate::view::View; +use lance_graph_mask_risc::{Foreign, LaneRef, Planes}; +use lance_graph_quack::{Cmp, Col, Filter, Mask}; +use ogar_dir_core::Guid128; +use ogar_dir_sim::{Attribute, Change, EvidenceRef, RuleId}; + +/// A pure graph transformation. +pub trait Rule { + /// Identity recorded in provenance. + fn id(&self) -> RuleId; + /// Proposed changes; must be deterministic in `(version, evidence)`. + fn propose(&self, v: &View<'_>, evidence: &[EvidenceRef]) -> Vec; +} + +/// Per-node membership count in `group` for one version: +/// `GROUP BY member COUNT(*) WHERE group = g` over the live base relation, +/// plus the overlay's added rows (delta-sized). One `K = |nodes|` sink; no +/// joined rows. +pub fn member_counts(v: &View<'_>, group: &Guid128) -> Vec { + let s = v.snap; + let mut counts = vec![0i64; s.len()]; + let Some(g) = s.ordinal(group) else { + return counts; + }; + let live = v.live_rows(); + let lanes = [LaneRef::U32(&s.m_user), LaneRef::U32(&s.m_group)]; + let masks: [&[u64]; 1] = [&live]; + let planes = Planes { + n_rows: s.membership_rows(), + masks: &masks, + lanes: &lanes, + }; + group_count_into( + Filter::and([Filter::plane(Mask(0)), Filter::cmp(Col(1), Cmp::EqU32(g))]), + Col(0), + &planes, + &Foreign::NONE, + &mut counts, + ); + for &(uo, go) in v.ov.added.values() { + if go == g && (uo as usize) < counts.len() { + counts[uo as usize] += 1; + } + } + counts +} + +/// Grant `group` to an explicit, evidence-sized list of users. Already +/// members are skipped, so the proposal is exactly the net change. The work +/// is proportional to the request, not to the directory. +pub struct GrantGroup { + /// Rule identity. + pub rule: RuleId, + /// The group. + pub group: Guid128, + /// Who should receive it. + pub to: Vec, +} + +impl Rule for GrantGroup { + fn id(&self) -> RuleId { + self.rule + } + fn propose(&self, v: &View<'_>, _: &[EvidenceRef]) -> Vec { + let mut to = self.to.clone(); + to.sort(); + to.dedup(); + to.into_iter() + .filter(|u| !v.is_member(u, &self.group)) + .map(|user| Change::AddMembership { + user, + group: self.group, + }) + .collect() + } +} + +/// Population rule: every active member of `source` is a member of `target`. +/// `active ∧ count(source) > 0 ∧ count(target) = 0`, from two folded +/// `GROUP BY` sinks and the resident active-user plane — no per-user workflow. +pub struct ImplyGroup { + /// Rule identity. + pub rule: RuleId, + /// Membership that implies… + pub source: Guid128, + /// …membership here. + pub target: Guid128, +} + +impl Rule for ImplyGroup { + fn id(&self) -> RuleId { + self.rule + } + fn propose(&self, v: &View<'_>, _: &[EvidenceRef]) -> Vec { + let (cs, ct) = ( + member_counts(v, &self.source), + member_counts(v, &self.target), + ); + let active = v.active_users(); + (0..v.len()) + .filter(|&o| active[o / 64] >> (o % 64) & 1 == 1 && cs[o] > 0 && ct[o] == 0) + .filter_map(|o| v.guid(o as u32)) + .map(|user| Change::AddMembership { + user, + group: self.target, + }) + .collect() + } +} + +/// Compare-and-set one user's primary SMTP against the version it reads. +pub struct SetPrimarySmtp { + /// Rule identity. + pub rule: RuleId, + /// User. + pub user: Guid128, + /// New address (raw). + pub to: String, +} + +impl Rule for SetPrimarySmtp { + fn id(&self) -> RuleId { + self.rule + } + fn propose(&self, v: &View<'_>, _: &[EvidenceRef]) -> Vec { + let Some(o) = v.ordinal(&self.user) else { + return Vec::new(); + }; + let from = v.attr(o, Attribute::PrimarySmtp); + if from == Some(self.to.as_str()) { + return Vec::new(); + } + vec![Change::SetAttribute { + node: self.user, + attribute: Attribute::PrimarySmtp, + from: from.map(str::to_string), + to: Some(self.to.clone()), + }] + } +} diff --git a/crates/lance-graph-dir-sim/src/snapshot.rs b/crates/lance-graph-dir-sim/src/snapshot.rs new file mode 100644 index 000000000..1e3cead12 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/snapshot.rs @@ -0,0 +1,343 @@ +//! One observed directory state as SoA lanes — never as objects. +//! +//! | lane / plane | type | meaning | +//! |-------------------------|-------------|---------------------------------------------| +//! | `ids` | `[Guid128]` | sorted; a node's **ordinal** is its index | +//! | `user`, `group` | bit planes | node kind | +//! | `active_user` | bit plane | user ∧ enabled — the recipient population | +//! | `upn_val` / `smtp_val` | `[u32]` | raw value id in [`Dicts::values`] | +//! | `upn_key` / `smtp_key` | `[u32]` | normalized key id in [`Dicts::keys`] | +//! | `ou_hi`, `ou_present` | `[u64]`, plane | OU-HHTL levels 0..4 packed (subtree = prefix match) | +//! | `m_user`, `m_group` | `[u32]` | membership relation, sorted by (user, group) | +//! +//! **Ordinal ≠ identity.** An ordinal is valid only inside one snapshot and +//! is never stored in provenance, diffs or plans; those carry [`Guid128`]. +//! Ordinals come from sorting by `Guid128`, so they — and every dictionary +//! id assigned while building — are independent of ingestion order. +//! +//! [`Observation`] is the ingestion boundary: it owns its strings once, and +//! [`Snapshot::build`] interns them into the store's dictionaries. After +//! that, execution works on ids and masks; a string is resolved only to +//! report evidence or to compare a compare-and-set value. + +use lance_graph_mask_risc::words_for; +use ogar_dir_core::{Guid128, OuHhtl}; +use ogar_dir_sim::normalize; +use std::collections::BTreeMap; + +/// "No value" / "unresolved endpoint" sentinel. Out of range for every +/// lane, so mask-risc's zero-fallback treats it as matching nothing. +pub const NONE: u32 = u32::MAX; + +/// Kind of a directory node. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum NodeKind { + /// User (a recipient when active with a primary SMTP). + User, + /// Group. + Group, +} + +/// Append-only string dictionary. Ids are assignment order; nothing is +/// iterated in hash order. +#[derive(Debug, Default, Clone, PartialEq, Eq)] +pub struct Dict { + strings: Vec, + index: BTreeMap, +} + +impl Dict { + /// Id of `s`, interning it if new. + pub fn intern(&mut self, s: &str) -> u32 { + if let Some(&i) = self.index.get(s) { + return i; + } + let i = u32::try_from(self.strings.len()).expect("dictionary fits u32"); + self.strings.push(s.to_string()); + self.index.insert(s.to_string(), i); + i + } + /// Id of `s` if interned. + pub fn get(&self, s: &str) -> Option { + self.index.get(s).copied() + } + /// The string behind an id. + pub fn resolve(&self, id: u32) -> Option<&str> { + self.strings.get(id as usize).map(String::as_str) + } + /// Number of entries (the group universe of a key `GROUP BY`). + pub fn len(&self) -> usize { + self.strings.len() + } + /// True if empty. + pub fn is_empty(&self) -> bool { + self.strings.is_empty() + } +} + +/// Value storage shared by every snapshot and version of one store. +#[derive(Debug, Default)] +pub struct Dicts { + /// Raw values exactly as observed. + pub values: Dict, + /// Normalized comparison keys. + pub keys: Dict, +} + +impl Dicts { + /// `(value id, key id)` for an optional raw value. + pub fn intern_attr(&mut self, s: Option<&str>) -> (u32, u32) { + match s { + None => (NONE, NONE), + Some(s) => (self.values.intern(s), self.keys.intern(&normalize(s))), + } + } + /// `(value id, key id)` without interning; `None` if not yet interned. + pub fn lookup_attr(&self, s: Option<&str>) -> Option<(u32, u32)> { + match s { + None => Some((NONE, NONE)), + Some(s) => Some((self.values.get(s)?, self.keys.get(&normalize(s))?)), + } + } +} + +/// One node as read from a source (ingestion staging). +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct ObservedNode { + /// Kind. + pub kind: NodeKind, + /// Enabled. + pub active: bool, + /// Raw UPN. + pub upn: Option, + /// Raw primary SMTP. + pub primary_smtp: Option, + /// OU location, if known. + pub ou: Option, +} + +impl ObservedNode { + /// Active user with UPN and primary SMTP. + pub fn user(upn: &str, smtp: &str) -> Self { + Self { + kind: NodeKind::User, + active: true, + upn: Some(upn.into()), + primary_smtp: Some(smtp.into()), + ou: None, + } + } + /// Group. + pub fn group() -> Self { + Self { + kind: NodeKind::Group, + active: true, + upn: None, + primary_smtp: None, + ou: None, + } + } +} + +/// What a source reported. Order of either list is irrelevant. +#[derive(Clone, Debug, Default)] +pub struct Observation { + /// Nodes. + pub nodes: Vec<(Guid128, ObservedNode)>, + /// `(user, group)` memberships; endpoints need not exist (dangling + /// observations are representable so validation can report them). + pub members: Vec<(Guid128, Guid128)>, +} + +/// Why an observation could not become a snapshot. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum BuildError { + /// The same Guid128 was observed twice (which one wins would depend on + /// ingestion order, so it is refused rather than guessed). + DuplicateNode(Guid128), + /// More than `u32::MAX - 1` nodes or memberships. + TooLarge, +} + +pub(crate) fn bit(words: &[u64], i: u32) -> bool { + i != NONE + && words + .get(i as usize / 64) + .is_some_and(|w| w >> (i % 64) & 1 == 1) +} +pub(crate) fn set_bit(words: &mut [u64], i: usize) { + words[i / 64] |= 1 << (i % 64); +} +pub(crate) fn clear_bit(words: &mut [u64], i: usize) { + words[i / 64] &= !(1 << (i % 64)); +} +pub(crate) fn ones(n: usize) -> Vec { + let mut w = vec![u64::MAX; words_for(n)]; + if !n.is_multiple_of(64) { + if let Some(last) = w.last_mut() { + *last = (1u64 << (n % 64)) - 1; + } + } + w +} + +/// OU-HHTL levels 0..4 packed into one `u64` (level 0 in the top 16 bits), +/// so "subtree of a prefix of depth ≤ 4" is one ternary match. +pub fn pack_ou(h: &OuHhtl) -> u64 { + (u64::from(h.0[0]) << 48) + | (u64::from(h.0[1]) << 32) + | (u64::from(h.0[2]) << 16) + | u64::from(h.0[3]) +} + +/// One observed state. Immutable once built; versions share it by `Arc`. +#[derive(Debug, PartialEq, Eq)] +pub struct Snapshot { + pub(crate) ids: Vec, + pub(crate) user: Vec, + pub(crate) group: Vec, + pub(crate) active_user: Vec, + pub(crate) upn_val: Vec, + pub(crate) upn_key: Vec, + pub(crate) smtp_val: Vec, + pub(crate) smtp_key: Vec, + pub(crate) ou_hi: Vec, + pub(crate) ou_present: Vec, + pub(crate) m_user: Vec, + pub(crate) m_group: Vec, + pub(crate) m_all: Vec, + /// Membership rows with an endpoint that resolved to no node: row → + /// the observed identities (rare; kept so evidence never loses identity). + pub(crate) m_unresolved: BTreeMap, +} + +impl Snapshot { + /// Build from an observation, interning strings into `d`. + pub fn build(mut obs: Observation, d: &mut Dicts) -> Result { + obs.nodes.sort_by_key(|n| n.0); + if let Some(w) = obs.nodes.windows(2).find(|w| w[0].0 == w[1].0) { + return Err(BuildError::DuplicateNode(w[0].0)); + } + let n = obs.nodes.len(); + if n >= NONE as usize || obs.members.len() >= NONE as usize { + return Err(BuildError::TooLarge); + } + let mut s = Self { + ids: Vec::with_capacity(n), + user: vec![0; words_for(n)], + group: vec![0; words_for(n)], + active_user: vec![0; words_for(n)], + upn_val: Vec::with_capacity(n), + upn_key: Vec::with_capacity(n), + smtp_val: Vec::with_capacity(n), + smtp_key: Vec::with_capacity(n), + ou_hi: Vec::with_capacity(n), + ou_present: vec![0; words_for(n)], + m_user: Vec::new(), + m_group: Vec::new(), + m_all: Vec::new(), + m_unresolved: BTreeMap::new(), + }; + for (i, (id, node)) in obs.nodes.iter().enumerate() { + s.ids.push(*id); + match node.kind { + NodeKind::User => { + set_bit(&mut s.user, i); + if node.active { + set_bit(&mut s.active_user, i); + } + } + NodeKind::Group => set_bit(&mut s.group, i), + } + let (uv, uk) = d.intern_attr(node.upn.as_deref()); + let (sv, sk) = d.intern_attr(node.primary_smtp.as_deref()); + s.upn_val.push(uv); + s.upn_key.push(uk); + s.smtp_val.push(sv); + s.smtp_key.push(sk); + s.ou_hi.push(node.ou.as_ref().map_or(0, pack_ou)); + if node.ou.is_some() { + set_bit(&mut s.ou_present, i); + } + } + let mut rows: Vec<(u32, u32, Guid128, Guid128)> = obs + .members + .iter() + .map(|(u, g)| { + ( + s.ordinal(u).unwrap_or(NONE), + s.ordinal(g).unwrap_or(NONE), + *u, + *g, + ) + }) + .collect(); + rows.sort(); + rows.dedup(); + for (r, (uo, go, ug, gg)) in rows.into_iter().enumerate() { + if uo == NONE || go == NONE { + s.m_unresolved.insert(r as u32, (ug, gg)); + } + s.m_user.push(uo); + s.m_group.push(go); + } + s.m_all = ones(s.m_user.len()); + Ok(s) + } + + /// Node count. + pub fn len(&self) -> usize { + self.ids.len() + } + /// True if no nodes. + pub fn is_empty(&self) -> bool { + self.ids.is_empty() + } + /// Membership row count. + pub fn membership_rows(&self) -> usize { + self.m_user.len() + } + /// Dense execution ordinal of an identity (binary search over the sorted id lane). + pub fn ordinal(&self, g: &Guid128) -> Option { + self.ids.binary_search(g).ok().map(|i| i as u32) + } + /// Identity of an ordinal. + pub fn guid(&self, o: u32) -> Option { + self.ids.get(o as usize).copied() + } + /// Identities of membership row `r`. + pub(crate) fn member_guids(&self, r: u32) -> (Guid128, Guid128) { + if let Some(p) = self.m_unresolved.get(&r) { + return *p; + } + ( + self.ids[self.m_user[r as usize] as usize], + self.ids[self.m_group[r as usize] as usize], + ) + } + /// Row of membership `(user, group)` in the base relation, if observed. + pub(crate) fn member_row(&self, user: &Guid128, group: &Guid128) -> Option { + match (self.ordinal(user), self.ordinal(group)) { + (Some(uo), Some(go)) => { + // Rows are sorted by the full (user, group, …) tuple, so the + // (user, group) prefix is monotonic over EVERY row — including + // half-resolved ones (`NONE` sorts after every ordinal). + let (mut lo, mut hi) = (0usize, self.m_user.len()); + while lo < hi { + let mid = (lo + hi) / 2; + match (self.m_user[mid], self.m_group[mid]).cmp(&(uo, go)) { + std::cmp::Ordering::Less => lo = mid + 1, + std::cmp::Ordering::Greater => hi = mid, + std::cmp::Ordering::Equal => return Some(mid as u32), + } + } + None + } + _ => self + .m_unresolved + .iter() + .find(|(_, p)| **p == (*user, *group)) + .map(|(r, _)| *r), + } + } +} diff --git a/crates/lance-graph-dir-sim/src/store.rs b/crates/lance-graph-dir-sim/src/store.rs new file mode 100644 index 000000000..77d364758 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/store.rs @@ -0,0 +1,342 @@ +//! The version history — and the audit trail. There is no other log. +//! +//! * An observation is a root: an immutable [`Snapshot`] behind an `Arc`. +//! * A simulated version stores only `parent`, `origin` (rule + evidence) and +//! `delta`. Its state is the root snapshot plus the overlay obtained by +//! folding the lineage's deltas — work and memory proportional to the +//! accumulated delta, never a copy of the population. +//! * Versions, snapshots and dictionaries are append-only. Simulating, +//! validating or rejecting a version cannot alter another. + +use crate::rule::Rule; +use crate::snapshot::{BuildError, Dicts, Observation, Snapshot}; +use crate::validate::validate; +use crate::view::{ApplyError, Overlay, View}; +use ogar_dir_core::Guid128; +use ogar_dir_sim::{ + Attribute, Change, EvidenceRef, ExecutionPlan, Origin, PlanError, Version, VersionId, + Violation, TAG_DESIRED, TAG_OBSERVED, +}; +use std::collections::{BTreeMap, BTreeSet}; +use std::sync::Arc; + +/// Simulation failure. No version is created. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum SimError { + /// No such version. + UnknownVersion(VersionId), + /// The rule proposed nothing (no new version). + EmptyProposal(ogar_dir_sim::RuleId), + /// The proposal does not apply to the parent. + Apply(ApplyError), + /// The two versions do not share a node set (diff unsupported in this slice). + NodeSetChanged, +} + +/// Refusal to make a version desired. The version stays as evidence. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Rejection { + /// The rejected version. + pub version: VersionId, + /// Why. + pub violations: Vec, +} + +/// Append-only version store. +#[derive(Debug, Default)] +pub struct VersionStore { + dicts: Dicts, + roots: BTreeMap>, + versions: Vec, + tags: BTreeMap, + verdicts: BTreeMap>, +} + +impl VersionStore { + /// Empty store. + pub fn new() -> Self { + Self::default() + } + + fn next_id(&self) -> VersionId { + VersionId(self.versions.len() as u64) + } + + /// Record an observation as a new root; tags it `"observed"`. + pub fn observe( + &mut self, + source: &str, + observed_at_ms: i64, + obs: Observation, + ) -> Result { + let snap = Snapshot::build(obs, &mut self.dicts)?; + let id = self.next_id(); + self.versions.push(Version { + id, + parent: None, + origin: Origin::Observed { + source: source.into(), + observed_at_ms, + }, + delta: Vec::new(), + }); + self.roots.insert(id, Arc::new(snap)); + self.tags.insert(TAG_OBSERVED.into(), id); + Ok(id) + } + + /// Provenance record. + pub fn version(&self, v: VersionId) -> Option<&Version> { + self.versions.get(v.0 as usize) + } + + /// Root first. + pub fn lineage(&self, v: VersionId) -> Result, SimError> { + let mut path = Vec::new(); + let mut cur = Some(v); + while let Some(c) = cur { + path.push(c); + cur = self.version(c).ok_or(SimError::UnknownVersion(c))?.parent; + } + path.reverse(); + Ok(path) + } + + /// The shared snapshot under `v`. + pub fn snapshot(&self, v: VersionId) -> Result<&Arc, SimError> { + let root = self.lineage(v)?[0]; + self.roots.get(&root).ok_or(SimError::UnknownVersion(root)) + } + + /// A coherent read of `v`: shared snapshot + folded overlay. + pub fn view(&self, v: VersionId) -> Result, SimError> { + let path = self.lineage(v)?; + let snap = self + .roots + .get(&path[0]) + .ok_or(SimError::UnknownVersion(path[0]))?; + let mut view = View::new(snap, &self.dicts, Overlay::default()); + for id in &path[1..] { + for c in &self.versions[id.0 as usize].delta { + view.apply(c).map_err(SimError::Apply)?; + } + } + Ok(view) + } + + /// Run a pure rule against `parent`, recording a hypothetical version. + pub fn simulate( + &mut self, + parent: VersionId, + rule: &dyn Rule, + evidence: &[EvidenceRef], + ) -> Result { + let delta = rule.propose(&self.view(parent)?, evidence); + if delta.is_empty() { + return Err(SimError::EmptyProposal(rule.id())); + } + for c in &delta { + if let Change::SetAttribute { to, .. } = c { + self.dicts.intern_attr(to.as_deref()); + } + } + let mut view = self.view(parent)?; + for c in &delta { + view.apply(c).map_err(SimError::Apply)?; + } + let id = self.next_id(); + self.versions.push(Version { + id, + parent: Some(parent), + origin: Origin::Simulated { + rule: rule.id(), + evidence: evidence.to_vec(), + }, + delta, + }); + Ok(id) + } + + /// Validate `v` and record the verdict. Empty = valid. + pub fn validate(&mut self, v: VersionId) -> Result, SimError> { + let violations = validate(&self.view(v)?); + self.verdicts.insert(v, violations.clone()); + Ok(violations) + } + + /// Recorded verdict. + pub fn verdict(&self, v: VersionId) -> Option<&[Violation]> { + self.verdicts.get(&v).map(Vec::as_slice) + } + + /// Make `v` desired — only if valid. On failure the tag is untouched. + pub fn promote_desired(&mut self, v: VersionId) -> Result<(), Rejection> { + let violations = self.validate(v).map_err(|_| Rejection { + version: v, + violations: Vec::new(), + })?; + if !violations.is_empty() { + return Err(Rejection { + version: v, + violations, + }); + } + self.tags.insert(TAG_DESIRED.into(), v); + Ok(()) + } + + /// Version a tag points to. + pub fn tag(&self, name: &str) -> Option { + self.tags.get(name).copied() + } + + /// Semantic difference `a → b`, sorted. Versions over the same snapshot + /// compare only their overlays' touched keys (delta-sized); versions over + /// different snapshots (reconciliation) merge the two sorted relations. + pub fn diff(&self, a: VersionId, b: VersionId) -> Result, SimError> { + let (va, vb) = (self.view(a)?, self.view(b)?); + let mut out = if std::ptr::eq(va.snap, vb.snap) { + diff_shared(&va, &vb) + } else { + diff_full(&va, &vb)? + }; + out.sort(); + Ok(out) + } + + /// Plan for the current desired version, from its observed basis. + pub fn plan(&self, target: VersionId) -> Result { + if self.tag(TAG_DESIRED) != Some(target) { + return Err(PlanError::NotDesired(target)); + } + let basis = self + .lineage(target) + .map_err(|_| PlanError::UnknownVersion(target))?[0]; + let diff = self + .diff(basis, target) + .map_err(|_| PlanError::UnknownVersion(target))?; + Ok(ExecutionPlan::from_diff(basis, target, diff)) + } + + /// Why does `v` contain membership `(user, group)`? The lineage from + /// the observed root to the version whose delta last added it; a chain + /// of length 1 means it was observed. `None` if `v` lacks it. + pub fn explain_membership( + &self, + v: VersionId, + user: &Guid128, + group: &Guid128, + ) -> Option> { + if !self.view(v).ok()?.is_member(user, group) { + return None; + } + let path = self.lineage(v).ok()?; + let add = Change::AddMembership { + user: *user, + group: *group, + }; + let at = path + .iter() + .rposition(|id| self.versions[id.0 as usize].delta.contains(&add)) + .unwrap_or(0); + Some( + path[..=at] + .iter() + .map(|id| &self.versions[id.0 as usize]) + .collect(), + ) + } +} + +fn set_change( + node: Guid128, + attribute: Attribute, + from: Option<&str>, + to: Option<&str>, +) -> Option { + (from != to).then(|| Change::SetAttribute { + node, + attribute, + from: from.map(str::to_string), + to: to.map(str::to_string), + }) +} + +/// Same snapshot: only keys either overlay touched can differ. +fn diff_shared(a: &View<'_>, b: &View<'_>) -> Vec { + let mut pairs: BTreeSet<(Guid128, Guid128)> = + a.ov.added + .keys() + .chain(b.ov.added.keys()) + .copied() + .collect(); + for v in [a, b] { + if let Some(rm) = &v.ov.removed { + for r in lance_graph_mask_risc::materialize_rows(rm, v.snap.membership_rows()) { + pairs.insert(v.snap.member_guids(r as u32)); + } + } + } + let mut out = Vec::new(); + for (u, g) in pairs { + match (a.is_member(&u, &g), b.is_member(&u, &g)) { + (false, true) => out.push(Change::AddMembership { user: u, group: g }), + (true, false) => out.push(Change::RemoveMembership { user: u, group: g }), + _ => {} + } + } + for attr in [Attribute::Upn, Attribute::PrimarySmtp] { + let touched: BTreeSet = match attr { + Attribute::Upn => a.ov.upn.keys().chain(b.ov.upn.keys()).copied().collect(), + Attribute::PrimarySmtp => a.ov.smtp.keys().chain(b.ov.smtp.keys()).copied().collect(), + }; + for o in touched { + out.extend(set_change( + a.snap.ids[o as usize], + attr, + a.attr(o, attr), + b.attr(o, attr), + )); + } + } + out +} + +/// Different snapshots (actual vs desired): a merge of two sorted relations, +/// O(n + m) — the reconciliation path, not the simulation hot path. +fn diff_full(a: &View<'_>, b: &View<'_>) -> Result, SimError> { + if a.snap.ids != b.snap.ids { + return Err(SimError::NodeSetChanged); + } + let effective = |v: &View<'_>| -> BTreeSet<(Guid128, Guid128)> { + let live = v.live_rows(); + lance_graph_mask_risc::materialize_rows(&live, v.snap.membership_rows()) + .into_iter() + .map(|r| v.snap.member_guids(r as u32)) + .chain(v.ov.added.keys().copied()) + .collect() + }; + let (ea, eb) = (effective(a), effective(b)); + let mut out: Vec = eb + .difference(&ea) + .map(|(u, g)| Change::AddMembership { + user: *u, + group: *g, + }) + .collect(); + out.extend(ea.difference(&eb).map(|(u, g)| Change::RemoveMembership { + user: *u, + group: *g, + })); + for o in 0..a.snap.len() as u32 { + for attr in [Attribute::Upn, Attribute::PrimarySmtp] { + out.extend(set_change( + a.snap.ids[o as usize], + attr, + a.attr(o, attr), + b.attr(o, attr), + )); + } + } + Ok(out) +} diff --git a/crates/lance-graph-dir-sim/src/validate.rs b/crates/lance-graph-dir-sim/src/validate.rs new file mode 100644 index 000000000..5c34f9466 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/validate.rs @@ -0,0 +1,213 @@ +//! Invariants as Quack programs over the version's lanes. +//! +//! | invariant | relational form | lowering | +//! |----------------------|---------------------------------------------------------|----------------------------------------| +//! | edge integrity | `members ANTI JOIN users ∪ members ANTI JOIN groups` | `Rows(¬Semijoin(user,·) ∨ ¬Semijoin(group,·))` — `MaskOp::Gather` over the node kind planes | +//! | unique SMTP / UPN | `GROUP BY key HAVING count > 1` over active users | `GroupReduce Count` keyed on the normalized-key dictionary id | +//! +//! No user objects are built. The integrity check reads the two membership +//! lanes and two kind planes — it cannot see a string. Uniqueness reads one +//! key lane and one plane; strings are resolved only for the (few) keys that +//! actually collide, to report them. +//! +//! Materialisations, all at the evidence boundary and bounded by the number +//! of violations: the offending membership rows (`materialize_rows` of the +//! kept mask) and each duplicate key's owner rows. + +use crate::exec::{group_count_into, keep, program}; +use crate::snapshot::bit; +use crate::view::View; +use lance_graph_mask_risc::{ + materialize_rows, Foreign, ForeignPlane as FPlane, LaneRef, Planes, Program, +}; +use lance_graph_quack::{Agg, Cmp, Col, Filter, ForeignPlane, Mask}; +use ogar_dir_core::Guid128; +use ogar_dir_sim::{normalize, Attribute, Endpoint, Violation}; + +/// The edge-integrity program over a membership relation: lane 0 = user +/// ordinal, lane 1 = group ordinal, plane 0 = live rows; foreign plane 0 = +/// the node table's user plane, 1 = its group plane. An anti-join on each +/// side; an unresolved endpoint (`NONE`) is out of range and gathers 0. +pub fn dangling_program() -> Program { + program( + Filter::and([ + Filter::plane(Mask(0)), + Filter::or([ + Filter::negate(Filter::semijoin(Col(0), ForeignPlane(0))), + Filter::negate(Filter::semijoin(Col(1), ForeignPlane(1))), + ]), + ]), + Agg::Rows, + ) +} + +/// Dangling memberships of a version (base rows still live + added rows). +pub fn dangling(v: &View<'_>) -> Vec { + let s = v.snap; + let fps = [ + FPlane { + words: &s.user, + rows: s.len(), + }, + FPlane { + words: &s.group, + rows: s.len(), + }, + ]; + let foreign = Foreign { + planes: &fps, + lanes: &[], + }; + let p = dangling_program(); + let side = |uo: u32| { + if bit(&s.user, uo) { + Endpoint::Group + } else { + Endpoint::User + } + }; + let mut out = Vec::new(); + + let live = v.live_rows(); + let lanes = [LaneRef::U32(&s.m_user), LaneRef::U32(&s.m_group)]; + let masks: [&[u64]; 1] = [&live]; + let planes = Planes { + n_rows: s.membership_rows(), + masks: &masks, + lanes: &lanes, + }; + for r in materialize_rows(&keep(&p, &planes, &foreign), planes.n_rows) { + let (user, group) = s.member_guids(r as u32); + out.push(Violation::DanglingMembership { + user, + group, + missing: side(s.m_user[r]), + }); + } + + // The overlay relation: delta-sized lanes, same program. + let (au, ag): (Vec, Vec) = v.ov.added.values().copied().unzip(); + let ids: Vec<(Guid128, Guid128)> = v.ov.added.keys().copied().collect(); + let all = crate::snapshot::ones(au.len()); + let lanes = [LaneRef::U32(&au), LaneRef::U32(&ag)]; + let masks: [&[u64]; 1] = [&all]; + let planes = Planes { + n_rows: au.len(), + masks: &masks, + lanes: &lanes, + }; + for r in materialize_rows(&keep(&p, &planes, &foreign), planes.n_rows) { + let (user, group) = ids[r]; + out.push(Violation::DanglingMembership { + user, + group, + missing: side(au[r]), + }); + } + out.sort(); + out +} + +/// Active users sharing a normalized value of `a`. +pub fn duplicates(v: &View<'_>, a: Attribute) -> Vec { + let s = v.snap; + let keys = &v.dicts.keys; + let k = keys.len(); + if k == 0 { + return Vec::new(); + } + let base_key = match a { + Attribute::Upn => &s.upn_key, + Attribute::PrimarySmtp => &s.smtp_key, + }; + let over = match a { + Attribute::Upn => &v.ov.upn, + Attribute::PrimarySmtp => &v.ov.smtp, + }; + let owners_plane = v.live_owners(a); + let mut counts = vec![0i64; k]; + + // Base: GROUP BY key COUNT(*) over the users still owning their observed value. + let lanes = [LaneRef::U32(base_key)]; + let masks: [&[u64]; 1] = [&owners_plane]; + let base = Planes { + n_rows: s.len(), + masks: &masks, + lanes: &lanes, + }; + group_count_into( + Filter::plane(Mask(0)), + Col(0), + &base, + &Foreign::NONE, + &mut counts, + ); + + // Overrides: the same GROUP BY over the delta rows, admitted by a + // semijoin of the overridden ordinal against the active-user plane. + let (oo, ok): (Vec, Vec) = over.iter().map(|(o, (_, key))| (*o, *key)).unzip(); + let all = crate::snapshot::ones(oo.len()); + let fps = [FPlane { + words: &s.active_user, + rows: s.len(), + }]; + let foreign = Foreign { + planes: &fps, + lanes: &[], + }; + let lanes = [LaneRef::U32(&oo), LaneRef::U32(&ok)]; + let masks: [&[u64]; 1] = [&all]; + let delta = Planes { + n_rows: oo.len(), + masks: &masks, + lanes: &lanes, + }; + group_count_into( + Filter::and([ + Filter::plane(Mask(0)), + Filter::semijoin(Col(0), ForeignPlane(0)), + ]), + Col(1), + &delta, + &foreign, + &mut counts, + ); + + let mut out = Vec::new(); + for (key, _) in counts.iter().enumerate().filter(|(_, c)| **c > 1) { + let key = key as u32; + // Owners: one gated equality over the base key lane (evidence boundary). + let p = program( + Filter::and([Filter::plane(Mask(0)), Filter::cmp(Col(0), Cmp::EqU32(key))]), + Agg::Rows, + ); + let mut owners: Vec = materialize_rows(&keep(&p, &base, &Foreign::NONE), s.len()) + .into_iter() + .map(|o| s.ids[o]) + .collect(); + owners.extend( + over.iter() + .filter(|(o, (_, kk))| *kk == key && bit(&s.active_user, **o)) + .map(|(o, _)| s.ids[*o as usize]), + ); + owners.sort(); + let value = keys.resolve(key).map(normalize).unwrap_or_default(); + out.push(match a { + Attribute::Upn => Violation::DuplicateUpn { upn: value, owners }, + Attribute::PrimarySmtp => Violation::DuplicateSmtp { + address: value, + owners, + }, + }); + } + out +} + +/// Every invariant, sorted. Empty = valid. +pub fn validate(v: &View<'_>) -> Vec { + let mut out = dangling(v); + out.extend(duplicates(v, Attribute::PrimarySmtp)); + out.extend(duplicates(v, Attribute::Upn)); + out.sort(); + out +} diff --git a/crates/lance-graph-dir-sim/src/view.rs b/crates/lance-graph-dir-sim/src/view.rs new file mode 100644 index 000000000..c53866719 --- /dev/null +++ b/crates/lance-graph-dir-sim/src/view.rs @@ -0,0 +1,258 @@ +//! A version as `shared base snapshot + overlay`. +//! +//! Structural sharing: every version derived from one observation borrows +//! the same [`Snapshot`] (held by `Arc` in the store). What a version adds is +//! an [`Overlay`] whose size is proportional to its accumulated delta: +//! +//! * added membership rows (a tiny second relation), +//! * a removed-rows bitmap over the base relation — allocated only if +//! something was removed, +//! * per-attribute override maps `ordinal → (value id, key id)`. +//! +//! Queries run over the base lanes gated by "still live" planes, plus over +//! the overlay rows, and fold the two results (a union is a sum of counts or +//! an OR of masks). The base is never copied. + +use crate::snapshot::{bit, clear_bit, set_bit, Dicts, Snapshot, NONE}; +use lance_graph_mask_risc::words_for; +use ogar_dir_core::Guid128; +use ogar_dir_sim::{Attribute, Change}; +use std::collections::BTreeMap; + +/// Why a change does not apply to a version. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum ApplyError { + /// `SetAttribute` on a node that does not exist. + UnknownNode(Guid128), + /// `SetAttribute`'s `from` no longer holds (stale compare-and-set). + Stale { + /// Node. + node: Guid128, + /// Attribute. + attribute: Attribute, + /// What the change expected. + expected: Option, + /// What the version holds. + actual: Option, + }, + /// The membership already holds (add) / does not hold (remove). + NoOp(Change), + /// A value in the change was never interned (store bug, not input). + Uninterned, +} + +/// The delta-sized part of a version. +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub(crate) struct Overlay { + /// Added memberships: identity pair → ordinals (`NONE` if unresolved). + pub(crate) added: BTreeMap<(Guid128, Guid128), (u32, u32)>, + /// Removed base membership rows; `None` until the first removal. + pub(crate) removed: Option>, + /// UPN overrides: ordinal → (value id, key id), `NONE` = cleared. + pub(crate) upn: BTreeMap, + /// Primary-SMTP overrides. + pub(crate) smtp: BTreeMap, +} + +impl Overlay { + /// Number of delta entries held (memberships + overrides). + pub(crate) fn delta_len(&self) -> usize { + self.added.len() + + self + .removed + .as_ref() + .map_or(0, |r| r.iter().map(|w| w.count_ones() as usize).sum()) + + self.upn.len() + + self.smtp.len() + } + fn overrides(&self, a: Attribute) -> &BTreeMap { + match a { + Attribute::Upn => &self.upn, + Attribute::PrimarySmtp => &self.smtp, + } + } +} + +/// A coherent read of one version: shared snapshot + overlay. +#[derive(Debug)] +pub struct View<'s> { + pub(crate) snap: &'s Snapshot, + pub(crate) dicts: &'s Dicts, + pub(crate) ov: Overlay, +} + +impl<'s> View<'s> { + pub(crate) fn new(snap: &'s Snapshot, dicts: &'s Dicts, ov: Overlay) -> Self { + Self { snap, dicts, ov } + } + + /// The shared snapshot (identity check for structural sharing). + pub fn snapshot(&self) -> &'s Snapshot { + self.snap + } + /// Delta entries this version holds over its snapshot. + pub fn delta_len(&self) -> usize { + self.ov.delta_len() + } + /// Dense execution ordinal of an identity. + pub fn ordinal(&self, g: &Guid128) -> Option { + self.snap.ordinal(g) + } + /// Identity of an ordinal. + pub fn guid(&self, o: u32) -> Option { + self.snap.guid(o) + } + /// Node count (the population width). + pub fn len(&self) -> usize { + self.snap.len() + } + /// True if the version has no nodes. + pub fn is_empty(&self) -> bool { + self.snap.is_empty() + } + /// Active users as a resident bit plane (borrowed). + pub fn active_users(&self) -> &'s [u64] { + &self.snap.active_user + } + + /// Effective membership, by identity. `O(log n)` + overlay lookup. + pub fn is_member(&self, user: &Guid128, group: &Guid128) -> bool { + if self.ov.added.contains_key(&(*user, *group)) { + return true; + } + self.snap + .member_row(user, group) + .is_some_and(|r| !self.ov.removed.as_ref().is_some_and(|rm| bit(rm, r))) + } + + /// Effective raw value of an attribute — the one place a value is + /// resolved to a string (compare-and-set, evidence, plan). + pub fn attr(&self, node: u32, a: Attribute) -> Option<&'s str> { + let val = match self.ov.overrides(a).get(&node) { + Some((v, _)) => *v, + None => match a { + Attribute::Upn => *self.snap.upn_val.get(node as usize)?, + Attribute::PrimarySmtp => *self.snap.smtp_val.get(node as usize)?, + }, + }; + (val != NONE) + .then(|| self.dicts.values.resolve(val)) + .flatten() + } + + /// The base membership rows still live in this version. + pub(crate) fn live_rows(&self) -> std::borrow::Cow<'s, [u64]> { + match &self.ov.removed { + None => std::borrow::Cow::Borrowed(&self.snap.m_all), + Some(rm) => std::borrow::Cow::Owned( + self.snap + .m_all + .iter() + .zip(rm) + .map(|(a, r)| a & !r) + .collect(), + ), + } + } + + /// `active_user` minus the nodes whose attribute `a` is overridden — + /// the base rows that still own their observed value. + pub(crate) fn live_owners(&self, a: Attribute) -> std::borrow::Cow<'s, [u64]> { + let ov = self.ov.overrides(a); + if ov.is_empty() { + return std::borrow::Cow::Borrowed(&self.snap.active_user); + } + let mut p = self.snap.active_user.clone(); + for o in ov.keys() { + clear_bit(&mut p, *o as usize); + } + std::borrow::Cow::Owned(p) + } + + /// Apply one change to the overlay. Pure with respect to everything but + /// `self.ov`; the snapshot is never touched. + pub(crate) fn apply(&mut self, c: &Change) -> Result<(), ApplyError> { + match c { + Change::AddMembership { user, group } => { + if self.is_member(user, group) { + return Err(ApplyError::NoOp(c.clone())); + } + match self.snap.member_row(user, group) { + Some(r) => { + let rm = self + .ov + .removed + .as_mut() + .expect("absent base row was removed"); + clear_bit(rm, r as usize); + if rm.iter().all(|w| *w == 0) { + self.ov.removed = None; + } + } + None => { + let ords = ( + self.ordinal(user).unwrap_or(NONE), + self.ordinal(group).unwrap_or(NONE), + ); + self.ov.added.insert((*user, *group), ords); + } + } + } + Change::RemoveMembership { user, group } => { + if self.ov.added.remove(&(*user, *group)).is_some() { + return Ok(()); + } + match self.snap.member_row(user, group) { + Some(r) if self.is_member(user, group) => { + let n = self.snap.membership_rows(); + let rm = self.ov.removed.get_or_insert_with(|| vec![0; words_for(n)]); + set_bit(rm, r as usize); + } + _ => return Err(ApplyError::NoOp(c.clone())), + } + } + Change::SetAttribute { + node, + attribute, + from, + to, + } => { + let o = self.ordinal(node).ok_or(ApplyError::UnknownNode(*node))?; + let actual = self.attr(o, *attribute); + if actual != from.as_deref() { + return Err(ApplyError::Stale { + node: *node, + attribute: *attribute, + expected: from.clone(), + actual: actual.map(str::to_string), + }); + } + let ids = self + .dicts + .lookup_attr(to.as_deref()) + .ok_or(ApplyError::Uninterned)?; + let base = match attribute { + Attribute::Upn => { + (self.snap.upn_val[o as usize], self.snap.upn_key[o as usize]) + } + Attribute::PrimarySmtp => ( + self.snap.smtp_val[o as usize], + self.snap.smtp_key[o as usize], + ), + }; + let map = match attribute { + Attribute::Upn => &mut self.ov.upn, + Attribute::PrimarySmtp => &mut self.ov.smtp, + }; + // Net effect only: setting a value back to the observed one + // removes the override instead of recording a no-op change. + if ids.0 == base.0 { + map.remove(&o); + } else { + map.insert(o, ids); + } + } + } + Ok(()) + } +} diff --git a/crates/lance-graph-dir-sim/tests/alloc.rs b/crates/lance-graph-dir-sim/tests/alloc.rs new file mode 100644 index 000000000..bd9066ad3 --- /dev/null +++ b/crates/lance-graph-dir-sim/tests/alloc.rs @@ -0,0 +1,90 @@ +//! Allocation instrumentation: a one-edge simulated mutation must cost work +//! proportional to the delta, not to the directory. Detects an accidentally +//! population-copying or row-exploding version path. +//! +//! Own test binary (the counting allocator is process-global). + +use lance_graph_dir_sim::*; +use ogar_dir_core::Guid128; +use ogar_dir_sim::{EvidenceRef, RuleId}; +use std::alloc::{GlobalAlloc, Layout, System}; +use std::sync::atomic::{AtomicUsize, Ordering}; + +struct Counting; +static BYTES: AtomicUsize = AtomicUsize::new(0); +unsafe impl GlobalAlloc for Counting { + unsafe fn alloc(&self, l: Layout) -> *mut u8 { + BYTES.fetch_add(l.size(), Ordering::Relaxed); + // SAFETY: forwards to the system allocator with the caller's layout. + unsafe { System.alloc(l) } + } + unsafe fn dealloc(&self, p: *mut u8, l: Layout) { + // SAFETY: `p` was returned by `System.alloc` with this layout. + unsafe { System.dealloc(p, l) } + } +} +#[global_allocator] +static A: Counting = Counting; + +fn guid(i: u32) -> Guid128 { + let mut b = [0u8; 16]; + b[..4].copy_from_slice(&i.to_be_bytes()); + b[15] = 1; + Guid128(b) +} + +/// Bytes allocated by `simulate` of one membership for a directory of `n` users. +fn one_edge(n: u32) -> (usize, usize) { + let (employees, exchange) = (guid(u32::MAX - 1), guid(u32::MAX)); + let mut obs = Observation::default(); + for i in 0..n { + obs.nodes.push(( + guid(i), + ObservedNode::user(&format!("u{i}@example.test"), &format!("u{i}@example.test")), + )); + obs.members.push((guid(i), employees)); + } + obs.nodes.push((employees, ObservedNode::group())); + obs.nodes.push((exchange, ObservedNode::group())); + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 0, obs).unwrap(); + let rule = GrantGroup { + rule: RuleId { + name: "ExchangeAccess", + version: 1, + }, + group: exchange, + to: vec![guid(7)], + }; + let ev = vec![EvidenceRef("REQ-1".into())]; + + let before = BYTES.load(Ordering::Relaxed); + let g1 = st.simulate(g0, &rule, &ev).unwrap(); + let simulate = BYTES.load(Ordering::Relaxed) - before; + + let before = BYTES.load(Ordering::Relaxed); + let d = st.diff(g0, g1).unwrap(); + let diff = BYTES.load(Ordering::Relaxed) - before; + assert_eq!(d.len(), 1); + (simulate, diff) +} + +#[test] +fn one_edge_mutation_is_delta_sized_not_population_sized() { + let (s_small, d_small) = one_edge(1_000); + let (s_large, d_large) = one_edge(100_000); + eprintln!( + "simulate: {s_small} B @1k users, {s_large} B @100k users; diff: {d_small} B / {d_large} B" + ); + // A copy of even one u32 lane at 100k users is 400 KB; a node-bitmap is + // 12.5 KB. Delta-sized work stays far below both and does not grow 100x. + assert!( + s_large < 4_096, + "simulate allocated {s_large} B at 100k users" + ); + assert!(d_large < 4_096, "diff allocated {d_large} B at 100k users"); + assert!( + s_large <= s_small + 256 && d_large <= d_small + 256, + "allocation grew with population" + ); +} diff --git a/crates/lance-graph-dir-sim/tests/sim.rs b/crates/lance-graph-dir-sim/tests/sim.rs new file mode 100644 index 000000000..e3967ed78 --- /dev/null +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -0,0 +1,654 @@ +//! observe → simulate → validate → diff → plan, over the SoA substrate. + +use lance_graph_dir_sim::validate::{dangling, dangling_program, validate}; +use lance_graph_dir_sim::*; +use lance_graph_mask_risc::{materialize_rows, MaskOp}; +use ogar_dir_core::{Guid128, OuHhtl}; +use ogar_dir_sim::*; +use std::sync::Arc; + +fn g(n: u8) -> Guid128 { + Guid128([n; 16]) +} +const ALICE: u8 = 0xA1; +const BOB: u8 = 0xB0; +const EMPLOYEES: u8 = 0xE0; +const EXCHANGE: u8 = 0xEC; + +const EXCHANGE_ACCESS: RuleId = RuleId { + name: "ExchangeAccess", + version: 1, +}; +const EMPLOYEES_GET_EXCHANGE: RuleId = RuleId { + name: "EmployeesGetExchange", + version: 1, +}; +const RENAME_MAIL: RuleId = RuleId { + name: "RenameMail", + version: 1, +}; + +fn observed() -> Observation { + Observation { + nodes: vec![ + ( + g(ALICE), + ObservedNode::user("alice@example.test", "alice@example.test"), + ), + ( + g(BOB), + ObservedNode::user("bob@example.test", "bob@example.test"), + ), + (g(EMPLOYEES), ObservedNode::group()), + (g(EXCHANGE), ObservedNode::group()), + ], + members: vec![(g(ALICE), g(EMPLOYEES)), (g(BOB), g(EMPLOYEES))], + } +} +fn grant_alice() -> GrantGroup { + GrantGroup { + rule: EXCHANGE_ACCESS, + group: g(EXCHANGE), + to: vec![g(ALICE)], + } +} +fn imply() -> ImplyGroup { + ImplyGroup { + rule: EMPLOYEES_GET_EXCHANGE, + source: g(EMPLOYEES), + target: g(EXCHANGE), + } +} +fn ev(s: &str) -> Vec { + vec![EvidenceRef(s.into())] +} +fn chain_from(obs: Observation) -> (VersionStore, VersionId, VersionId, VersionId) { + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 1_000, obs).unwrap(); + let g1 = st.simulate(g0, &grant_alice(), &ev("REQ-1")).unwrap(); + let g2 = st.simulate(g1, &imply(), &ev("POLICY-7")).unwrap(); + (st, g0, g1, g2) +} +fn chain() -> (VersionStore, VersionId, VersionId, VersionId) { + chain_from(observed()) +} + +// 1. G0 is unchanged by simulation. +#[test] +fn t01_g0_immutable() { + let mut fresh = VersionStore::new(); + let only = fresh.observe("lab", 1_000, observed()).unwrap(); + let (st, g0, ..) = chain(); + assert_eq!(**st.snapshot(g0).unwrap(), **fresh.snapshot(only).unwrap()); + assert_eq!(st.view(g0).unwrap().delta_len(), 0); + assert!(!st.view(g0).unwrap().is_member(&g(ALICE), &g(EXCHANGE))); +} + +// 2 + 10. G1 has parent G0; provenance names the rule and the evidence. +#[test] +fn t02_t10_child_with_provenance() { + let (st, g0, g1, _) = chain(); + let v = st.version(g1).unwrap(); + assert_eq!(v.parent, Some(g0)); + assert_eq!( + v.origin, + Origin::Simulated { + rule: EXCHANGE_ACCESS, + evidence: ev("REQ-1") + } + ); + assert_eq!( + v.delta, + vec![Change::AddMembership { + user: g(ALICE), + group: g(EXCHANGE) + }] + ); + assert!(st.view(g1).unwrap().is_member(&g(ALICE), &g(EXCHANGE))); +} + +// 3. A population rule chains from G1 and adds only Bob. +#[test] +fn t03_population_rule_chains() { + let (mut st, g0, g1, g2) = chain(); + assert_eq!(st.version(g2).unwrap().parent, Some(g1)); + assert_eq!( + st.version(g2).unwrap().delta, + vec![Change::AddMembership { + user: g(BOB), + group: g(EXCHANGE) + }] + ); + assert_eq!(st.lineage(g2).unwrap(), vec![g0, g1, g2]); + assert_eq!( + st.simulate(g2, &imply(), &[]), + Err(SimError::EmptyProposal(EMPLOYEES_GET_EXCHANGE)) + ); +} + +// 4. diff(G0, G1) is one semantic change. +#[test] +fn t04_diff_one_change() { + let (st, g0, g1, g2) = chain(); + assert_eq!( + st.diff(g0, g1).unwrap(), + vec![Change::AddMembership { + user: g(ALICE), + group: g(EXCHANGE) + }] + ); + assert_eq!(st.diff(g0, g2).unwrap().len(), 2); + assert_eq!( + st.diff(g2, g0).unwrap().len(), + 2, + "reverse diff removes both" + ); + assert!(st.diff(g1, g1).unwrap().is_empty()); +} + +// 5. Valid memberships pass. +#[test] +fn t05_valid_passes() { + let (mut st, g0, _, g2) = chain(); + assert!(st.validate(g0).unwrap().is_empty()); + assert!(st.validate(g2).unwrap().is_empty()); + st.promote_desired(g2).unwrap(); + assert_eq!(st.tag(TAG_DESIRED), Some(g2)); +} + +// 6. Dangling memberships fail — in the base relation (observed) and in the +// overlay (simulated), on both endpoint sides, with identity preserved. +#[test] +fn t06_dangling_fails() { + let ghost = g(0x66); + let mut obs = observed(); + obs.members.push((ghost, g(EMPLOYEES))); // unknown user + obs.members.push((g(EMPLOYEES), g(EXCHANGE))); // a group as member + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 0, obs).unwrap(); + let bad = st + .simulate( + g0, + &GrantGroup { + rule: EXCHANGE_ACCESS, + group: ghost, + to: vec![g(BOB)], + }, + &[], + ) + .unwrap(); + let v = st.validate(bad).unwrap(); + assert_eq!( + v, + vec![ + Violation::DanglingMembership { + user: ghost, + group: g(EMPLOYEES), + missing: Endpoint::User + }, + Violation::DanglingMembership { + user: g(BOB), + group: ghost, + missing: Endpoint::Group + }, + Violation::DanglingMembership { + user: g(EMPLOYEES), + group: g(EXCHANGE), + missing: Endpoint::User + }, + ] + ); +} + +// 7. Duplicate UPN fails (normalized); a disabled owner does not count. +#[test] +fn t07_duplicate_upn() { + let mut obs = observed(); + obs.nodes.push(( + g(0x77), + ObservedNode::user(" ALICE@example.test", "other@example.test"), + )); + let mut st = VersionStore::new(); + let v0 = st.observe("lab", 0, obs.clone()).unwrap(); + assert_eq!( + st.validate(v0).unwrap(), + vec![Violation::DuplicateUpn { + upn: "alice@example.test".into(), + owners: vec![g(0x77), g(ALICE)] + }] + ); + obs.nodes.last_mut().unwrap().1.active = false; + let v1 = st.observe("lab", 1, obs).unwrap(); + assert!(st.validate(v1).unwrap().is_empty()); +} + +// 8 + 9 + rejected future: Bob.primarySMTP = alice@… is constructible, +// detected through the overlay, refused, kept — and nothing else moves. +#[test] +fn t08_t09_smtp_collision_rejected() { + let (mut st, g0, g1, g2) = chain(); + st.promote_desired(g2).unwrap(); + let before: Vec<_> = [g0, g1, g2] + .iter() + .map(|v| st.diff(g0, *v).unwrap()) + .collect(); + let snap_before = Arc::clone(st.snapshot(g0).unwrap()); + + let rule = SetPrimarySmtp { + rule: RENAME_MAIL, + user: g(BOB), + to: "Alice@Example.test".into(), + }; + let g3 = st + .simulate(g2, &rule, &ev("REQ-2")) + .expect("hypothetical future is constructible"); + let rej = st.promote_desired(g3).unwrap_err(); + assert_eq!( + rej.violations, + vec![Violation::DuplicateSmtp { + address: "alice@example.test".into(), + owners: vec![g(ALICE), g(BOB)] + }] + ); + assert_eq!(st.tag(TAG_DESIRED), Some(g2)); + let after: Vec<_> = [g0, g1, g2] + .iter() + .map(|v| st.diff(g0, *v).unwrap()) + .collect(); + assert_eq!(before, after); + assert!( + Arc::ptr_eq(&snap_before, st.snapshot(g3).unwrap()), + "the base was never copied" + ); + assert_eq!(st.verdict(g3), Some(rej.violations.as_slice())); + let bob = st.view(g3).unwrap().ordinal(&g(BOB)).unwrap(); + assert_eq!( + st.view(g3).unwrap().attr(bob, Attribute::PrimarySmtp), + Some("Alice@Example.test") + ); +} + +// A stale compare-and-set creates no version. +#[test] +fn stale_change_creates_no_version() { + struct Stale; + impl Rule for Stale { + fn id(&self) -> RuleId { + RENAME_MAIL + } + fn propose(&self, _: &View<'_>, _: &[EvidenceRef]) -> Vec { + vec![Change::SetAttribute { + node: g(BOB), + attribute: Attribute::PrimarySmtp, + from: Some("not-it@example.test".into()), + to: Some("x@example.test".into()), + }] + } + } + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 0, observed()).unwrap(); + assert!(matches!( + st.simulate(g0, &Stale, &[]), + Err(SimError::Apply(ApplyError::Stale { .. })) + )); + assert!(st.version(VersionId(1)).is_none()); +} + +// 11. The diff becomes a plan with basis preconditions and no transport. +#[test] +fn t11_plan() { + let (mut st, g0, _, g2) = chain(); + assert_eq!(st.plan(g2), Err(PlanError::NotDesired(g2))); + st.promote_desired(g2).unwrap(); + let plan = st.plan(g2).unwrap(); + assert_eq!((plan.basis, plan.target), (g0, g2)); + assert_eq!( + plan.ops, + vec![ + PlannedOp { + op: Operation::AddGroupMember { + group: g(EXCHANGE), + member: g(ALICE) + }, + precondition: Precondition::NotMember + }, + PlannedOp { + op: Operation::AddGroupMember { + group: g(EXCHANGE), + member: g(BOB) + }, + precondition: Precondition::NotMember + }, + ] + ); + let mut s2 = VersionStore::new(); + let b0 = s2.observe("lab", 0, observed()).unwrap(); + let b1 = s2 + .simulate( + b0, + &SetPrimarySmtp { + rule: RENAME_MAIL, + user: g(BOB), + to: "robert@example.test".into(), + }, + &[], + ) + .unwrap(); + s2.promote_desired(b1).unwrap(); + assert_eq!( + s2.plan(b1).unwrap().ops, + vec![PlannedOp { + op: Operation::SetAttribute { + object: g(BOB), + attribute: Attribute::PrimarySmtp, + value: Some("robert@example.test".into()) + }, + precondition: Precondition::AttributeEquals(Some("bob@example.test".into())), + }] + ); +} + +// 12. Same input + same rules ⇒ same desired state, diff and plan. +#[test] +fn t12_deterministic() { + let run = || { + let (mut st, g0, _, g2) = chain(); + st.promote_desired(g2).unwrap(); + (st.diff(g0, g2).unwrap(), st.plan(g2).unwrap()) + }; + assert_eq!(run(), run()); +} + +// 13. Guid128 survives the dense-ordinal mapping, including GUIDs that +// differ only in their upper 64 bits. +#[test] +fn t13_guid_ordinal_round_trip() { + let a = Guid128::parse("00000000-0000-0001-0123-456789abcdef").unwrap(); + let b = Guid128::parse("00000000-0000-0002-0123-456789abcdef").unwrap(); + let mut obs = observed(); + obs.nodes.push((b, ObservedNode::group())); + obs.nodes.push((a, ObservedNode::group())); + let mut st = VersionStore::new(); + let v = st.observe("lab", 0, obs).unwrap(); + let view = st.view(v).unwrap(); + for id in [a, b, g(ALICE), g(BOB), g(EMPLOYEES), g(EXCHANGE)] { + let o = view.ordinal(&id).unwrap(); + assert_eq!(view.guid(o), Some(id)); + } + assert_ne!(view.ordinal(&a), view.ordinal(&b)); + assert_eq!(view.ordinal(&g(0x99)), None); +} + +// 14. Membership validation reads only membership lanes and kind planes: it +// finds the dangling edge in a directory with no strings at all. +#[test] +fn t14_membership_validation_reads_no_attributes() { + let node = |kind| ObservedNode { + kind, + active: true, + upn: None, + primary_smtp: None, + ou: None, + }; + let obs = Observation { + nodes: vec![(g(1), node(NodeKind::User)), (g(2), node(NodeKind::Group))], + members: vec![(g(1), g(2)), (g(1), g(3))], + }; + let mut st = VersionStore::new(); + let v = st.observe("lab", 0, obs).unwrap(); + let view = st.view(v).unwrap(); + assert_eq!( + dangling(&view), + vec![Violation::DanglingMembership { + user: g(1), + group: g(3), + missing: Endpoint::Group + }] + ); +} + +// 15. Edge integrity is an anti-join lowered to MaskOp::Gather (Quack's +// semijoin), not a hand-written hash join. +#[test] +fn t15_dangling_uses_semijoin() { + let p = dangling_program(); + let gathers = p + .ops + .iter() + .filter(|op| matches!(op, MaskOp::Gather { .. })) + .count(); + assert_eq!(gathers, 2, "one semijoin per endpoint side: {:?}", p.ops); +} + +// 16. A one-edge mutation shares the snapshot and holds one delta entry. +#[test] +fn t16_structural_sharing() { + let (st, g0, g1, g2) = chain(); + let s0 = st.snapshot(g0).unwrap(); + assert!(Arc::ptr_eq(s0, st.snapshot(g1).unwrap())); + assert!(Arc::ptr_eq(s0, st.snapshot(g2).unwrap())); + assert_eq!(st.view(g1).unwrap().delta_len(), 1); + assert_eq!(st.view(g2).unwrap().delta_len(), 2); +} + +// 17. Semantic output is invariant under input ordering. +#[test] +fn t17_ordering_invariance() { + let mut shuffled = observed(); + shuffled.nodes.reverse(); + shuffled.nodes.swap(0, 2); + shuffled.members.reverse(); + shuffled.members.push((g(ALICE), g(EMPLOYEES))); // duplicate observation + let outputs = |obs| { + let (mut st, g0, _, g2) = chain_from(obs); + st.promote_desired(g2).unwrap(); + let rule = SetPrimarySmtp { + rule: RENAME_MAIL, + user: g(BOB), + to: "alice@example.test".into(), + }; + let g3 = st.simulate(g2, &rule, &[]).unwrap(); + ( + st.diff(g0, g2).unwrap(), + st.plan(g2).unwrap().ops, + st.validate(g3).unwrap(), + ) + }; + assert_eq!(outputs(observed()), outputs(shuffled)); +} + +// Audit: why is Alice in ExchangeUsers? +#[test] +fn audit_chain() { + let (st, g0, g1, g2) = chain(); + let why = st.explain_membership(g2, &g(ALICE), &g(EXCHANGE)).unwrap(); + assert_eq!(why.iter().map(|v| v.id).collect::>(), vec![g0, g1]); + assert!(matches!(why[0].origin, Origin::Observed { .. })); + match &why[1].origin { + Origin::Simulated { rule, evidence } => { + assert_eq!(rule.to_string(), "ExchangeAccess/v1"); + assert_eq!(evidence, &ev("REQ-1")); + } + o => panic!("unexpected {o:?}"), + } + assert_eq!( + st.explain_membership(g2, &g(BOB), &g(EXCHANGE)) + .unwrap() + .last() + .unwrap() + .id, + g2 + ); + assert_eq!( + st.explain_membership(g2, &g(ALICE), &g(EMPLOYEES)) + .unwrap() + .len(), + 1 + ); + assert!(st.explain_membership(g0, &g(ALICE), &g(EXCHANGE)).is_none()); +} + +// Convergence: re-observing the desired state diffs empty (full-merge path). +#[test] +fn converged_observation_diffs_empty() { + let (mut st, _, _, g2) = chain(); + st.promote_desired(g2).unwrap(); + let mut actual = observed(); + actual + .members + .extend([(g(ALICE), g(EXCHANGE)), (g(BOB), g(EXCHANGE))]); + let o = st.observe("lab", 2_000, actual).unwrap(); + assert!(st.diff(o, g2).unwrap().is_empty()); + let mut drifted = observed(); + drifted.members.push((g(ALICE), g(EXCHANGE))); + let d = st.observe("lab", 3_000, drifted).unwrap(); + assert_eq!( + st.diff(d, g2).unwrap(), + vec![Change::AddMembership { + user: g(BOB), + group: g(EXCHANGE) + }] + ); +} + +// Removal and re-add net out; removals reach the plan as RemoveGroupMember. +#[test] +fn removal_round_trip() { + struct Drop; + impl Rule for Drop { + fn id(&self) -> RuleId { + RuleId { + name: "Leaver", + version: 1, + } + } + fn propose(&self, _: &View<'_>, _: &[EvidenceRef]) -> Vec { + vec![Change::RemoveMembership { + user: g(BOB), + group: g(EMPLOYEES), + }] + } + } + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 0, observed()).unwrap(); + let r1 = st.simulate(g0, &Drop, &[]).unwrap(); + assert!(!st.view(r1).unwrap().is_member(&g(BOB), &g(EMPLOYEES))); + let r2 = st + .simulate( + r1, + &GrantGroup { + rule: EXCHANGE_ACCESS, + group: g(EMPLOYEES), + to: vec![g(BOB)], + }, + &[], + ) + .unwrap(); + assert!(st.diff(g0, r2).unwrap().is_empty()); + st.promote_desired(r1).unwrap(); + assert_eq!( + st.plan(r1).unwrap().ops, + vec![PlannedOp { + op: Operation::RemoveGroupMember { + group: g(EMPLOYEES), + member: g(BOB) + }, + precondition: Precondition::IsMember + }] + ); +} + +// HHTL as executable geometry: subtree selection is one prefix match. +#[test] +fn ou_subtree_is_a_prefix_match() { + let ou = |l: &[u16]| { + let mut h = OuHhtl::ROOT; + h.0[..l.len()].copy_from_slice(l); + h + }; + let mut obs = observed(); + obs.nodes[0].1.ou = Some(ou(&[1, 1, 1])); // Stuttgart/Infrastructure/Exchange + obs.nodes[1].1.ou = Some(ou(&[2])); // Berlin + obs.nodes.push(( + g(0x55), + ObservedNode { + ou: Some(ou(&[1, 2])), + ..ObservedNode::user("c@x", "c@x") + }, + )); + let mut st = VersionStore::new(); + let v = st.observe("lab", 0, obs).unwrap(); + let view = st.view(v).unwrap(); + let pick = |p: &OuHhtl| -> Vec { + materialize_rows(&subtree(&view, p).unwrap(), view.len()) + .into_iter() + .map(|o| view.guid(o as u32).unwrap()) + .collect() + }; + assert_eq!(pick(&ou(&[1])), vec![g(0x55), g(ALICE)]); + assert_eq!(pick(&ou(&[1, 1])), vec![g(ALICE)]); + assert_eq!(pick(&ou(&[2])), vec![g(BOB)]); + assert_eq!(pick(&OuHhtl::ROOT).len(), 3, "root = every located node"); + assert_eq!( + subtree(&view, &ou(&[1, 1, 1, 1, 1])), + Err(SubtreeTooDeep(5)) + ); +} + +// Observation from OGAR PR #313 records. +#[test] +fn observe_from_ogar_ad() { + use ogar_dir_core::{OuDictionary, ValuePool}; + let ldif = "dn: CN=Alice,OU=Staff,DC=example,DC=test\nobjectGUID:: 4AQlP4lP0xGaDAMF6CwzAQ==\nobjectClass: user\nuserPrincipalName: alice@example.test\nproxyAddresses: smtp:a@legacy.test\nproxyAddresses: SMTP:alice@example.test\nuserAccountControl: 514\n\ndn: CN=Employees,OU=Groups,DC=example,DC=test\nobjectGUID:: 1MOyoQAAAECAAAAAAAC+7w==\nobjectClass: group\n"; + let (mut d, mut p) = (OuDictionary::new(), ValuePool::new()); + let recs: Vec<_> = ogar_ad::ldif::parse(ldif) + .unwrap() + .iter() + .map(|e| { + ogar_ad::encode(e, Guid128::NIL, &mut d, &mut p, 0) + .unwrap() + .record + }) + .collect(); + let obs = observe::from_ad(&recs, &p); + assert_eq!( + obs.nodes[0].1.primary_smtp.as_deref(), + Some("alice@example.test") + ); + assert!(!obs.nodes[0].1.active, "UAC 514 = disabled"); + assert!(obs.nodes[0].1.ou.is_some()); + assert_eq!(obs.nodes[1].1.kind, NodeKind::Group); + let mut st = VersionStore::new(); + let v = st.observe("ogar-ad:ldif", 0, obs).unwrap(); + assert!(validate(&st.view(v).unwrap()).is_empty()); +} + +// Overrides REPLACE the observed value: renaming the second owner away from +// an observed collision must clear it (the base owner plane drops overridden +// users; otherwise the old value would still be counted). +#[test] +fn renaming_away_resolves_an_observed_collision() { + let carol = g(0xC0); + let mut obs = observed(); + obs.nodes.push(( + carol, + ObservedNode::user("carol@example.test", "BOB@example.test"), + )); + let mut st = VersionStore::new(); + let g0 = st.observe("lab", 0, obs).unwrap(); + assert_eq!( + st.validate(g0).unwrap(), + vec![Violation::DuplicateSmtp { + address: "bob@example.test".into(), + owners: vec![g(BOB), carol] + }] + ); + let fix = SetPrimarySmtp { + rule: RENAME_MAIL, + user: carol, + to: "carol@example.test".into(), + }; + let g1 = st.simulate(g0, &fix, &[]).unwrap(); + assert!(st.validate(g1).unwrap().is_empty()); + st.promote_desired(g1).unwrap(); +} From f6306af971b001bf196936587780fd6f43f7df3f Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:09:44 +0000 Subject: [PATCH 2/6] board: entry for lance-graph-dir-sim (directory simulation over SoA + Quack) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- .../entries/2026-10-03-dir-sim-soa-quack.md | 34 +++++++++++++++++++ .claude/board/entries/README.md | 3 +- 2 files changed, 36 insertions(+), 1 deletion(-) create mode 100644 .claude/board/entries/2026-10-03-dir-sim-soa-quack.md diff --git a/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md b/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md new file mode 100644 index 000000000..d483694cc --- /dev/null +++ b/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md @@ -0,0 +1,34 @@ +# Directory desired-state simulation over SoA + Quack (2026-10-03) + +**Status:** MEASURED. The `crates/lance-graph-dir-sim` crate is excluded from +the workspace because it path-depends on the OGAR sibling, the same shape as +`lance-graph-report-ogar`. It consumes the semantic vocabulary in OGAR +`ogar-dir-sim` (OGAR PR #314). + +A directory version is a shared `Arc` of SoA lanes plus a +delta-sized overlay. Invariants are Quack programs: + +- **Edge integrity** is an anti-join: `negate(Semijoin)`, lowered to two + `MaskOp::Gather` over the node-kind planes. +- **SMTP / UPN uniqueness** is `GroupReduce Count` keyed on a normalized-key + dictionary id, folded over base + overlay. The overlay rows are admitted by a + `Semijoin` against the active-user plane. + +Population rules combine two folded `GROUP BY` sinks at the consumer. They +never pass one program's mask to another program's `Semijoin`. + +| quantity (one membership mutation, `tests/alloc.rs`) | 1k users | 100k users | +|---|---|---| +| bytes allocated by `simulate` | 853 | 853 | +| bytes allocated by same-root `diff` | 1,208 | 1,208 | + +The suite has 23 tests. Ten guards were disabled one at a time and each turned +its test red. + +Open for the operator: + +- **`VersionedGraph` (`u32` node ids, additions-only diff).** It cannot + persist 128-bit directory identity. The choice is to widen it upstream or + use a dedicated directory dataset. +- **Row-population masks.** There is still no shared row-population mask type + in the contract; this crate uses mask-risc bitmaps directly. diff --git a/.claude/board/entries/README.md b/.claude/board/entries/README.md index 64fd79d33..6c11f4df9 100644 --- a/.claude/board/entries/README.md +++ b/.claude/board/entries/README.md @@ -25,10 +25,11 @@ index row, (3) no duplicate entry id. Checks 1 and 2 are deliberately opposite directions; the stranding this convention prevents shows up in exactly one of them, never both. -177 entries, 2026-08-06 .. 2026-09-30. +178 entries, 2026-08-06 .. 2026-10-03. | date | entry id | finding | file | |---|---|---|---| +| 2026-10-03 | `dir-sim-soa-quack` | Directory simulation on SoA + Quack: one-edge mutation 853 B at 1k and 100k users | [2026-10-03-dir-sim-soa-quack.md](2026-10-03-dir-sim-soa-quack.md) | | 2026-09-30 | `deepnsm-v2-coverage-bands` | | [2026-09-30-deepnsm-v2-coverage-bands.md](2026-09-30-deepnsm-v2-coverage-bands.md) | | 2026-09-29 | `deepnsm-v2-counted-pick-tag-deltas` | | [2026-09-29-deepnsm-v2-counted-pick-tag-deltas.md](2026-09-29-deepnsm-v2-counted-pick-tag-deltas.md) | | 2026-09-26 | `deepnsm-v2-lexical-evidence-survives-routing` | | [2026-09-26-deepnsm-v2-lexical-evidence-survives-routing.md](2026-09-26-deepnsm-v2-lexical-evidence-survives-routing.md) | From 7c88dc7abcbdce81f61e4d6a9ec604e8ceacd089 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:14:18 +0000 Subject: [PATCH 3/6] dir-sim: seal kept masks so no program output can feed a Semijoin A program's kept rows now come back as Kept, whose only exit is rows() (the evidence boundary). It has no &[u64] view, so a ForeignPlane cannot be built from it: feeding one program's mask into another's Semijoin is a compile error, pinned by a compile_fail doctest with a passing twin. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- crates/lance-graph-dir-sim/src/exec.rs | 46 ++++++++++++++++++++-- crates/lance-graph-dir-sim/src/lib.rs | 3 +- crates/lance-graph-dir-sim/src/validate.rs | 8 ++-- crates/lance-graph-dir-sim/tests/sim.rs | 4 +- 4 files changed, 50 insertions(+), 11 deletions(-) diff --git a/crates/lance-graph-dir-sim/src/exec.rs b/crates/lance-graph-dir-sim/src/exec.rs index 099953691..543bcb0ed 100644 --- a/crates/lance-graph-dir-sim/src/exec.rs +++ b/crates/lance-graph-dir-sim/src/exec.rs @@ -8,7 +8,7 @@ //! `words_for(n_rows)` bitmap, a `GroupReduce` in a `K`-slot sink. use lance_graph_mask_risc::{ - execute_into, words_for, Foreign, Out, Planes, Program, Scratch, Terminal, Value, + execute_into, materialize_rows, words_for, Foreign, Out, Planes, Program, Scratch, Terminal, Value, }; use lance_graph_quack::{lower, Agg, Col, Filter, GroupAddr, GroupAgg, Query}; @@ -23,14 +23,52 @@ fn run(p: &Program, planes: &Planes<'_>, foreign: &Foreign<'_>, out: Out<'_>) -> execute_into(p, planes, foreign, &mut scratch, out).expect("directory program runs") } -/// Surviving rows as a `words_for(n_rows)` bitmap (`Agg::Rows` → `Keep`). -pub(crate) fn keep(p: &Program, planes: &Planes<'_>, foreign: &Foreign<'_>) -> Vec { +/// The rows one program kept. Its only exit is [`Kept::rows`], the evidence +/// boundary: it has no `&[u64]` view, so it cannot become another program's +/// `ForeignPlane`. A `Semijoin` gathers only from resident planes (node kinds, +/// the active-user plane), never from a population a program produced. +/// +/// ```compile_fail +/// use lance_graph_dir_sim::Kept; +/// use lance_graph_mask_risc::ForeignPlane; +/// fn feed(k: &Kept) -> ForeignPlane<'_> { +/// ForeignPlane { words: k, rows: 0 } +/// } +/// ``` +/// +/// The same imports and shape compile against a resident plane: +/// +/// ``` +/// use lance_graph_dir_sim::Kept; +/// use lance_graph_mask_risc::ForeignPlane; +/// fn feed<'a>(_k: &Kept, resident: &'a [u64]) -> ForeignPlane<'a> { +/// ForeignPlane { words: resident, rows: 0 } +/// } +/// ``` +#[derive(Debug, PartialEq, Eq)] +pub struct Kept { + bits: Vec, + n_rows: usize, +} + +impl Kept { + /// Kept row indices, ascending. Bounded by the number of survivors. + pub fn rows(&self) -> Vec { + materialize_rows(&self.bits, self.n_rows) + } +} + +/// Surviving rows (`Agg::Rows` → `Keep`), sealed in a [`Kept`]. +pub(crate) fn keep(p: &Program, planes: &Planes<'_>, foreign: &Foreign<'_>) -> Kept { debug_assert!(matches!(p.terminal, Terminal::Keep { .. })); let mut bits = vec![0u64; words_for(planes.n_rows)]; if planes.n_rows > 0 { run(p, planes, foreign, Out::Mask(&mut bits)); } - bits + Kept { + bits, + n_rows: planes.n_rows, + } } /// `GROUP BY key COUNT(*)` over the rows `filter` keeps, ADDED into `sink` diff --git a/crates/lance-graph-dir-sim/src/lib.rs b/crates/lance-graph-dir-sim/src/lib.rs index 986b270c6..4ee72f011 100644 --- a/crates/lance-graph-dir-sim/src/lib.rs +++ b/crates/lance-graph-dir-sim/src/lib.rs @@ -33,6 +33,7 @@ pub use rule::{member_counts, GrantGroup, ImplyGroup, Rule, SetPrimarySmtp}; pub use snapshot::{ pack_ou, BuildError, Dict, Dicts, NodeKind, Observation, ObservedNode, Snapshot, NONE, }; +pub use exec::Kept; pub use store::{Rejection, SimError, VersionStore}; pub use view::{ApplyError, View}; @@ -47,7 +48,7 @@ pub struct SubtreeTooDeep(pub usize); /// Nodes located in the OU subtree `prefix` (ancestor-or-self), as a node /// bitmap — one ternary match on the packed OU lane (`Cmp::MatchU64`), no DN /// strings. Prefixes up to depth 4 are exact; deeper ones are refused. -pub fn subtree(v: &View<'_>, prefix: &OuHhtl) -> Result, SubtreeTooDeep> { +pub fn subtree(v: &View<'_>, prefix: &OuHhtl) -> Result { let d = prefix.depth(); if d > 4 { return Err(SubtreeTooDeep(d)); diff --git a/crates/lance-graph-dir-sim/src/validate.rs b/crates/lance-graph-dir-sim/src/validate.rs index 5c34f9466..a243928f7 100644 --- a/crates/lance-graph-dir-sim/src/validate.rs +++ b/crates/lance-graph-dir-sim/src/validate.rs @@ -18,7 +18,7 @@ use crate::exec::{group_count_into, keep, program}; use crate::snapshot::bit; use crate::view::View; use lance_graph_mask_risc::{ - materialize_rows, Foreign, ForeignPlane as FPlane, LaneRef, Planes, Program, + Foreign, ForeignPlane as FPlane, LaneRef, Planes, Program, }; use lance_graph_quack::{Agg, Cmp, Col, Filter, ForeignPlane, Mask}; use ogar_dir_core::Guid128; @@ -76,7 +76,7 @@ pub fn dangling(v: &View<'_>) -> Vec { masks: &masks, lanes: &lanes, }; - for r in materialize_rows(&keep(&p, &planes, &foreign), planes.n_rows) { + for r in keep(&p, &planes, &foreign).rows() { let (user, group) = s.member_guids(r as u32); out.push(Violation::DanglingMembership { user, @@ -96,7 +96,7 @@ pub fn dangling(v: &View<'_>) -> Vec { masks: &masks, lanes: &lanes, }; - for r in materialize_rows(&keep(&p, &planes, &foreign), planes.n_rows) { + for r in keep(&p, &planes, &foreign).rows() { let (user, group) = ids[r]; out.push(Violation::DanglingMembership { user, @@ -181,7 +181,7 @@ pub fn duplicates(v: &View<'_>, a: Attribute) -> Vec { Filter::and([Filter::plane(Mask(0)), Filter::cmp(Col(0), Cmp::EqU32(key))]), Agg::Rows, ); - let mut owners: Vec = materialize_rows(&keep(&p, &base, &Foreign::NONE), s.len()) + let mut owners: Vec = keep(&p, &base, &Foreign::NONE).rows() .into_iter() .map(|o| s.ids[o]) .collect(); diff --git a/crates/lance-graph-dir-sim/tests/sim.rs b/crates/lance-graph-dir-sim/tests/sim.rs index e3967ed78..f878e69c8 100644 --- a/crates/lance-graph-dir-sim/tests/sim.rs +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -2,7 +2,7 @@ use lance_graph_dir_sim::validate::{dangling, dangling_program, validate}; use lance_graph_dir_sim::*; -use lance_graph_mask_risc::{materialize_rows, MaskOp}; +use lance_graph_mask_risc::MaskOp; use ogar_dir_core::{Guid128, OuHhtl}; use ogar_dir_sim::*; use std::sync::Arc; @@ -580,7 +580,7 @@ fn ou_subtree_is_a_prefix_match() { let v = st.observe("lab", 0, obs).unwrap(); let view = st.view(v).unwrap(); let pick = |p: &OuHhtl| -> Vec { - materialize_rows(&subtree(&view, p).unwrap(), view.len()) + subtree(&view, p).unwrap().rows() .into_iter() .map(|o| view.guid(o as u32).unwrap()) .collect() From e184d5b353364677f9c1fd066cb20b003653d1e9 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:15:00 +0000 Subject: [PATCH 4/6] dir-sim: document why ImplyGroup stays two programs, not ogar-loco Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- .../entries/2026-10-03-dir-sim-soa-quack.md | 18 ++++++++++++++++++ crates/lance-graph-dir-sim/src/exec.rs | 3 ++- crates/lance-graph-dir-sim/src/lib.rs | 2 +- crates/lance-graph-dir-sim/src/rule.rs | 14 ++++++++++++++ crates/lance-graph-dir-sim/src/validate.rs | 7 +++---- crates/lance-graph-dir-sim/tests/sim.rs | 4 +++- 6 files changed, 41 insertions(+), 7 deletions(-) diff --git a/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md b/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md index d483694cc..a9c3899a3 100644 --- a/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md +++ b/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md @@ -32,3 +32,21 @@ Open for the operator: use a dedicated directory dataset. - **Row-population masks.** There is still no shared row-population mask type in the contract; this crate uses mask-risc bitmaps directly. + +## Addendum (same day): no produced mask into a Semijoin, now a compile error + +- Kept rows from a program are sealed in `Kept`, which only exits through + `rows()`. Nothing can build a `ForeignPlane` from it. A `compile_fail` + doctest pins this, with a passing twin. Disable-run: giving `Kept` a slice + `Deref` turns the doctest red. +- Audit by reading: every `ForeignPlane` in the crate is resident (the user + plane, the group plane, the active-user plane). +- ndarray is mandatory. It comes in through `lance-graph-mask-risc`, a + non-optional path dependency (`cargo tree -i ndarray`). +- All builds and tests ran with `CARGO_PROFILE_DEV_DEBUG=0` and + `CARGO_INCREMENTAL=0`. +- `ogar-loco` was not adopted. The only loco→mask-risc dialect is a test-local + `FoldDialect` in `r2il-mask-abi-probe`. Its `GROUP_SUM` sink keeps only the + last fold, so it cannot express ImplyGroup's two keyed counts. +- OPEN: promote a fold dialect to a library, with a multi-sink `GROUP_SUM`. + Then rules can be loco program data. diff --git a/crates/lance-graph-dir-sim/src/exec.rs b/crates/lance-graph-dir-sim/src/exec.rs index 543bcb0ed..a54c5ea1c 100644 --- a/crates/lance-graph-dir-sim/src/exec.rs +++ b/crates/lance-graph-dir-sim/src/exec.rs @@ -8,7 +8,8 @@ //! `words_for(n_rows)` bitmap, a `GroupReduce` in a `K`-slot sink. use lance_graph_mask_risc::{ - execute_into, materialize_rows, words_for, Foreign, Out, Planes, Program, Scratch, Terminal, Value, + execute_into, materialize_rows, words_for, Foreign, Out, Planes, Program, Scratch, Terminal, + Value, }; use lance_graph_quack::{lower, Agg, Col, Filter, GroupAddr, GroupAgg, Query}; diff --git a/crates/lance-graph-dir-sim/src/lib.rs b/crates/lance-graph-dir-sim/src/lib.rs index 4ee72f011..59e18c85b 100644 --- a/crates/lance-graph-dir-sim/src/lib.rs +++ b/crates/lance-graph-dir-sim/src/lib.rs @@ -29,11 +29,11 @@ pub mod store; pub mod validate; pub mod view; +pub use exec::Kept; pub use rule::{member_counts, GrantGroup, ImplyGroup, Rule, SetPrimarySmtp}; pub use snapshot::{ pack_ou, BuildError, Dict, Dicts, NodeKind, Observation, ObservedNode, Snapshot, NONE, }; -pub use exec::Kept; pub use store::{Rejection, SimError, VersionStore}; pub use view::{ApplyError, View}; diff --git a/crates/lance-graph-dir-sim/src/rule.rs b/crates/lance-graph-dir-sim/src/rule.rs index 745addcc3..e7e9cb85a 100644 --- a/crates/lance-graph-dir-sim/src/rule.rs +++ b/crates/lance-graph-dir-sim/src/rule.rs @@ -87,6 +87,20 @@ impl Rule for GrantGroup { /// Population rule: every active member of `source` is a member of `target`. /// `active ∧ count(source) > 0 ∧ count(target) = 0`, from two folded /// `GROUP BY` sinks and the resident active-user plane — no per-user workflow. +/// +/// Why two programs and not one: "member of A and not member of B" is a +/// per-user fact over the membership relation, so a single program would +/// have to `Semijoin` users against a population another program produced. +/// That is the forbidden shape ([`crate::Kept`] makes it a compile error). +/// The two sinks are keyed by user ordinal and sized by the node universe, +/// never by membership rows, and are combined here at the consumer. +/// +/// Not an `ogar-loco` program, for now. The only loco dialect that lowers to +/// mask-risc (`FoldDialect`) lives inside a test file in +/// `r2il-mask-abi-probe`, not in a library. It also combines only scalar +/// folds; its `GROUP_SUM` sink is zero-filled per fold and keeps just the +/// last one, so it cannot hold both counts. Copying it here would make a +/// second arity table. That dialect has to become a library first. pub struct ImplyGroup { /// Rule identity. pub rule: RuleId, diff --git a/crates/lance-graph-dir-sim/src/validate.rs b/crates/lance-graph-dir-sim/src/validate.rs index a243928f7..26f84735a 100644 --- a/crates/lance-graph-dir-sim/src/validate.rs +++ b/crates/lance-graph-dir-sim/src/validate.rs @@ -17,9 +17,7 @@ use crate::exec::{group_count_into, keep, program}; use crate::snapshot::bit; use crate::view::View; -use lance_graph_mask_risc::{ - Foreign, ForeignPlane as FPlane, LaneRef, Planes, Program, -}; +use lance_graph_mask_risc::{Foreign, ForeignPlane as FPlane, LaneRef, Planes, Program}; use lance_graph_quack::{Agg, Cmp, Col, Filter, ForeignPlane, Mask}; use ogar_dir_core::Guid128; use ogar_dir_sim::{normalize, Attribute, Endpoint, Violation}; @@ -181,7 +179,8 @@ pub fn duplicates(v: &View<'_>, a: Attribute) -> Vec { Filter::and([Filter::plane(Mask(0)), Filter::cmp(Col(0), Cmp::EqU32(key))]), Agg::Rows, ); - let mut owners: Vec = keep(&p, &base, &Foreign::NONE).rows() + let mut owners: Vec = keep(&p, &base, &Foreign::NONE) + .rows() .into_iter() .map(|o| s.ids[o]) .collect(); diff --git a/crates/lance-graph-dir-sim/tests/sim.rs b/crates/lance-graph-dir-sim/tests/sim.rs index f878e69c8..ad1270ca2 100644 --- a/crates/lance-graph-dir-sim/tests/sim.rs +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -580,7 +580,9 @@ fn ou_subtree_is_a_prefix_match() { let v = st.observe("lab", 0, obs).unwrap(); let view = st.view(v).unwrap(); let pick = |p: &OuHhtl| -> Vec { - subtree(&view, p).unwrap().rows() + subtree(&view, p) + .unwrap() + .rows() .into_iter() .map(|o| view.guid(o as u32).unwrap()) .collect() From 0aefa9a7a05925828930d2e39099e8efb2e26013 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:33:03 +0000 Subject: [PATCH 5/6] dir-sim: base the execution plan on the latest observation plan() used the desired version's lineage root as its basis. After a re-observation that already carries part of the desired state, the plan then repeated those changes and their NotMember preconditions would fail on execution. The basis is now the version tagged observed (the lineage root only before any observation is tagged). Test: plan_is_based_on_the_latest_observation. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- crates/lance-graph-dir-sim/src/store.rs | 11 +++++++++-- crates/lance-graph-dir-sim/tests/sim.rs | 19 +++++++++++++++++++ 2 files changed, 28 insertions(+), 2 deletions(-) diff --git a/crates/lance-graph-dir-sim/src/store.rs b/crates/lance-graph-dir-sim/src/store.rs index 77d364758..af9b3dc81 100644 --- a/crates/lance-graph-dir-sim/src/store.rs +++ b/crates/lance-graph-dir-sim/src/store.rs @@ -204,14 +204,21 @@ impl VersionStore { Ok(out) } - /// Plan for the current desired version, from its observed basis. + /// Plan for the current desired version, from the LATEST observation. + /// + /// The basis is the version tagged observed, not the desired version's + /// own lineage root: after a re-observation the directory may already + /// carry part of the desired state, and the plan must cover only what is + /// still missing, or its `NotMember` preconditions fail on execution. + /// Before any observation is tagged, the lineage root is the basis. pub fn plan(&self, target: VersionId) -> Result { if self.tag(TAG_DESIRED) != Some(target) { return Err(PlanError::NotDesired(target)); } - let basis = self + let root = self .lineage(target) .map_err(|_| PlanError::UnknownVersion(target))?[0]; + let basis = self.tag(TAG_OBSERVED).unwrap_or(root); let diff = self .diff(basis, target) .map_err(|_| PlanError::UnknownVersion(target))?; diff --git a/crates/lance-graph-dir-sim/tests/sim.rs b/crates/lance-graph-dir-sim/tests/sim.rs index ad1270ca2..e0aaa8220 100644 --- a/crates/lance-graph-dir-sim/tests/sim.rs +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -489,6 +489,25 @@ fn audit_chain() { } // Convergence: re-observing the desired state diffs empty (full-merge path). +#[test] +fn plan_is_based_on_the_latest_observation() { + let (mut st, g0, _, g2) = chain(); + st.promote_desired(g2).unwrap(); + assert_eq!(st.plan(g2).unwrap().basis, g0); + let mut partly = observed(); + partly.members.push((g(ALICE), g(EXCHANGE))); + let o = st.observe("lab", 2_000, partly).unwrap(); + let plan = st.plan(g2).unwrap(); + assert_eq!(plan.basis, o); + assert_eq!( + plan.ops, + vec![PlannedOp::from(Change::AddMembership { + user: g(BOB), + group: g(EXCHANGE) + })] + ); +} + #[test] fn converged_observation_diffs_empty() { let (mut st, _, _, g2) = chain(); From 9a2cafaf209ae5f3075f35f55962a5ea67ae057a Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 07:40:14 +0000 Subject: [PATCH 6/6] dir-sim: report a changed node set from plan(), not UnknownVersion Basing the plan on the latest observation makes a cross-root diff possible; a node-set change there surfaced as UnknownVersion(target). plan() now maps NodeSetChanged and UnknownVersion separately. Test: plan_reports_a_changed_node_set. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01G22yT6htkcdyXsihxxXdrg --- crates/lance-graph-dir-sim/src/store.rs | 9 ++++++--- crates/lance-graph-dir-sim/tests/sim.rs | 19 +++++++++++++++++++ 2 files changed, 25 insertions(+), 3 deletions(-) diff --git a/crates/lance-graph-dir-sim/src/store.rs b/crates/lance-graph-dir-sim/src/store.rs index af9b3dc81..062928b0c 100644 --- a/crates/lance-graph-dir-sim/src/store.rs +++ b/crates/lance-graph-dir-sim/src/store.rs @@ -219,9 +219,12 @@ impl VersionStore { .lineage(target) .map_err(|_| PlanError::UnknownVersion(target))?[0]; let basis = self.tag(TAG_OBSERVED).unwrap_or(root); - let diff = self - .diff(basis, target) - .map_err(|_| PlanError::UnknownVersion(target))?; + let diff = self.diff(basis, target).map_err(|e| match e { + SimError::NodeSetChanged => PlanError::NodeSetChanged { basis, target }, + SimError::UnknownVersion(v) => PlanError::UnknownVersion(v), + // `diff` raises nothing else. + _ => PlanError::UnknownVersion(target), + })?; Ok(ExecutionPlan::from_diff(basis, target, diff)) } diff --git a/crates/lance-graph-dir-sim/tests/sim.rs b/crates/lance-graph-dir-sim/tests/sim.rs index e0aaa8220..ee282e058 100644 --- a/crates/lance-graph-dir-sim/tests/sim.rs +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -508,6 +508,25 @@ fn plan_is_based_on_the_latest_observation() { ); } +#[test] +fn plan_reports_a_changed_node_set() { + let (mut st, _, _, g2) = chain(); + st.promote_desired(g2).unwrap(); + let mut grown = observed(); + grown.nodes.push(( + g(99), + ObservedNode::user("new@example.test", "new@example.test"), + )); + let o = st.observe("lab", 2_000, grown).unwrap(); + assert_eq!( + st.plan(g2), + Err(PlanError::NodeSetChanged { + basis: o, + target: g2 + }) + ); +} + #[test] fn converged_observation_diffs_empty() { let (mut st, _, _, g2) = chain();