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..a9c3899a3 --- /dev/null +++ b/.claude/board/entries/2026-10-03-dir-sim-soa-quack.md @@ -0,0 +1,52 @@ +# 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. + +## 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/.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) | 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..a54c5ea1c --- /dev/null +++ b/crates/lance-graph-dir-sim/src/exec.rs @@ -0,0 +1,100 @@ +//! 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, materialize_rows, 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") +} + +/// 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)); + } + Kept { + bits, + n_rows: planes.n_rows, + } +} + +/// `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..59e18c85b --- /dev/null +++ b/crates/lance-graph-dir-sim/src/lib.rs @@ -0,0 +1,79 @@ +//! # 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 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 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 { + 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..e7e9cb85a --- /dev/null +++ b/crates/lance-graph-dir-sim/src/rule.rs @@ -0,0 +1,163 @@ +//! 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. +/// +/// 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, + /// 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..062928b0c --- /dev/null +++ b/crates/lance-graph-dir-sim/src/store.rs @@ -0,0 +1,352 @@ +//! 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 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 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(|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)) + } + + /// 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..26f84735a --- /dev/null +++ b/crates/lance-graph-dir-sim/src/validate.rs @@ -0,0 +1,212 @@ +//! 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::{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 keep(&p, &planes, &foreign).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 keep(&p, &planes, &foreign).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 = keep(&p, &base, &Foreign::NONE) + .rows() + .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..ee282e058 --- /dev/null +++ b/crates/lance-graph-dir-sim/tests/sim.rs @@ -0,0 +1,694 @@ +//! 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::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 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 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(); + 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 { + subtree(&view, p) + .unwrap() + .rows() + .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(); +}