From 2b892e15de3d290e676ea62ea01df7b5c6d2add2 Mon Sep 17 00:00:00 2001 From: Stefan Date: Thu, 1 Oct 2026 21:07:58 +0200 Subject: [PATCH] fix: record affected documents at change time Edits queue their path; before the next snapshot the engine walks dependents with a frozen copy of every recorded view (each dependent on the graph stack) and stamps affected documents with a sequence. Consumers read them via changes_since(cursor), so queries in between cannot absorb a change. Graph dictionary views now cover the combined import list, matching first-wins resolution. --- bindings/nodejs/index.d.ts | 6 +- bindings/nodejs/src/policy.rs | 19 +++-- core/engine/src/workspace/affected.rs | 30 +++++-- core/engine/src/workspace/db.rs | 95 ++++++++++++++++------ core/engine/src/workspace/graph/queries.rs | 6 +- core/engine/src/workspace/mod.rs | 4 +- core/engine/src/workspace/reads.rs | 67 ++++++++------- 7 files changed, 149 insertions(+), 78 deletions(-) diff --git a/bindings/nodejs/index.d.ts b/bindings/nodejs/index.d.ts index 0b3f4eb0..41b61b3e 100644 --- a/bindings/nodejs/index.d.ts +++ b/bindings/nodejs/index.d.ts @@ -118,6 +118,10 @@ export type PolicySlotState = | 'path'; export type PolicySlotRole = 'unary' | 'condition' | 'value' | 'path'; /** One enum value: `label` for display, `source` the ready-to-splice ZEN literal (null when unquotable). */ +export interface PolicyChanges { + cursor: number + paths: Array +} export interface PolicyValueOption { value: string; label: string; @@ -400,7 +404,7 @@ export declare class Workspace { isGraph(path: string): boolean uncheckedNodes(path: string): Array paths(): Array - affectedBy(paths: Array): Array + changesSince(cursor: number): PolicyChanges updateBlock(req: PolicyUpdateBlockRequest): void removeBlock(req: PolicyRemoveBlockRequest): boolean diagnostics(policyPath: string, maxDiagnostics?: number | undefined | null): Array diff --git a/bindings/nodejs/src/policy.rs b/bindings/nodejs/src/policy.rs index 421bb128..ccad4648 100644 --- a/bindings/nodejs/src/policy.rs +++ b/bindings/nodejs/src/policy.rs @@ -12,6 +12,12 @@ use zen_engine::workspace; type ResolverRef = FunctionRef, Option>; +#[napi(object)] +pub struct PolicyChanges { + pub cursor: i64, + pub paths: Vec, +} + #[napi(object)] pub struct PolicyExpressionCursor { pub policy_path: String, @@ -627,13 +633,12 @@ impl Workspace { } #[napi] - pub fn affected_by(&self, paths: Vec) -> Vec { - let paths: Vec<&str> = paths.iter().map(String::as_str).collect(); - self.inner - .affected_by(&paths) - .into_iter() - .map(|p| p.to_string()) - .collect() + pub fn changes_since(&self, cursor: i64) -> PolicyChanges { + let (next, paths) = self.inner.changes_since(cursor.max(0) as u64); + PolicyChanges { + cursor: next as i64, + paths: paths.into_iter().map(|p| p.to_string()).collect(), + } } #[napi] diff --git a/core/engine/src/workspace/affected.rs b/core/engine/src/workspace/affected.rs index f083bb21..7c6e289c 100644 --- a/core/engine/src/workspace/affected.rs +++ b/core/engine/src/workspace/affected.rs @@ -2,17 +2,31 @@ use std::sync::Arc; use crate::workspace::db::Db; use crate::workspace::document_dependencies::DependencyKind; +use crate::workspace::reads::ReadView; impl Db { - pub fn affected_by(&self, paths: &[&str]) -> Vec> { - let changed: Vec> = paths.iter().map(|path| Arc::from(*path)).collect(); + pub(crate) fn walk(&self, changed: &[Arc]) -> Vec> { + let frozen = self.frozen_views(); self.document_dependencies() - .affected(&changed, |user, dependency, kind| match kind { - DependencyKind::Import => true, - DependencyKind::Dictionaries | DependencyKind::Signature => { - match self.recorded_view(user, dependency, kind) { - Some(view) => !self.view_holds(dependency, &view), - None => true, + .affected(changed, |user, dependency, kind| { + let recorded = frozen + .get(&(user.clone(), dependency.clone())) + .into_iter() + .flatten() + .find(|view| { + matches!( + (kind, view), + (DependencyKind::Dictionaries, ReadView::Dictionaries(_)) + | (DependencyKind::Signature, ReadView::Signature(_)) + ) + }); + match (kind, recorded) { + (DependencyKind::Import, _) | (_, None) => true, + (_, Some(view)) => { + self.graph_stack.borrow_mut().push(user.clone()); + let holds = self.view_holds(dependency, view); + self.graph_stack.borrow_mut().pop(); + !holds } } }) diff --git a/core/engine/src/workspace/db.rs b/core/engine/src/workspace/db.rs index ca69c220..d6baf353 100644 --- a/core/engine/src/workspace/db.rs +++ b/core/engine/src/workspace/db.rs @@ -25,7 +25,7 @@ use crate::policy::queries::scope::{ VariableTypeScope, }; use crate::policy::raw::PolicyDocument; -use crate::workspace::document_dependencies::{DependencyIndex, DependencyKind}; +use crate::workspace::document_dependencies::DependencyIndex; use crate::workspace::graph::function::{ FunctionKey, FunctionResolutionRequest, FunctionTypeResolver, ResolvedFunction, }; @@ -45,6 +45,14 @@ pub(crate) struct GraphDeps { functions: Vec<(FunctionKey, u64)>, } +#[derive(Default)] +pub(crate) struct ChangeLog { + pending: Vec>, + settling: bool, + sequence: u64, + marked: HashMap, u64>, +} + #[derive(Default)] pub(crate) struct DepFrame { docs: HashSet>, @@ -197,6 +205,7 @@ pub struct Db { pub(crate) graph_stack: RefCell>>, graph_dep_frames: RefCell>, dependencies: DependencyIndex, + changes: RefCell, graph_fn_frames: RefCell>>, function_types: RefCell>, function_requests: RefCell>, @@ -227,6 +236,7 @@ impl Db { graph_stack: RefCell::new(Vec::new()), graph_dep_frames: RefCell::new(Vec::new()), dependencies: DependencyIndex::default(), + changes: RefCell::new(ChangeLog::default()), graph_fn_frames: RefCell::new(Vec::new()), function_types: RefCell::new(HashMap::default()), function_requests: RefCell::new(Vec::new()), @@ -239,6 +249,7 @@ impl Db { pub fn set_document(&mut self, path: Arc, doc: Arc) { self.dependencies.set(path.clone(), &doc); + self.changes.borrow_mut().pending.push(path.clone()); self.inputs.borrow_mut().documents.insert(path, doc); self.invalidate_snapshot(); } @@ -249,6 +260,7 @@ impl Db { pub fn remove_document(&mut self, path: &str) -> bool { self.dependencies.remove(path); + self.changes.borrow_mut().pending.push(Arc::from(path)); let existed = self.inputs.borrow_mut().documents.remove(path).is_some(); if existed { self.invalidate_snapshot(); @@ -256,36 +268,61 @@ impl Db { existed } - pub(crate) fn recorded_view( - &self, - user: &Arc, - dependency: &Arc, - kind: DependencyKind, - ) -> Option { + pub(crate) fn frozen_views(&self) -> HashMap<(Arc, Arc), Vec> { let cache = self.cache.graphs.borrow(); - let (deps, _) = cache.get(user)?; let inputs = self.inputs.borrow(); - let current = deps.docs.iter().any(|(doc, stamp)| { - doc == user - && matches!( - (stamp, inputs.documents.get(user)), - (Some(stamp), Some(now)) if Arc::ptr_eq(stamp, now) - ) - }); - if !current { - return None; - } - deps.views - .iter() - .find(|(doc, view)| { - doc == dependency + let mut out: HashMap<(Arc, Arc), Vec> = HashMap::default(); + for (user, (deps, _)) in cache.iter() { + let current = deps.docs.iter().any(|(doc, stamp)| { + doc == user && matches!( - (kind, view), - (DependencyKind::Dictionaries, ReadView::Dictionaries(_)) - | (DependencyKind::Signature, ReadView::Signature(_)) + (stamp, inputs.documents.get(user)), + (Some(stamp), Some(now)) if Arc::ptr_eq(stamp, now) ) - }) - .map(|(_, view)| view.clone()) + }); + if !current { + continue; + } + for (dependency, view) in &deps.views { + out.entry((user.clone(), dependency.clone())) + .or_default() + .push(view.clone()); + } + } + out + } + + pub(crate) fn settle(&self) { + if self.changes.borrow().settling || self.changes.borrow().pending.is_empty() { + return; + } + let pending = { + let mut changes = self.changes.borrow_mut(); + changes.settling = true; + std::mem::take(&mut changes.pending) + }; + let affected = self.walk(&pending); + let mut changes = self.changes.borrow_mut(); + changes.settling = false; + changes.sequence += 1; + let sequence = changes.sequence; + for path in affected { + changes.marked.insert(path, sequence); + } + } + + pub(crate) fn changes_since(&self, cursor: u64) -> (u64, Vec>) { + self.snapshot(); + self.settle(); + let changes = self.changes.borrow(); + let mut paths: Vec> = changes + .marked + .iter() + .filter(|(_, &sequence)| sequence > cursor) + .map(|(path, _)| path.clone()) + .collect(); + paths.sort(); + (changes.sequence, paths) } pub(crate) fn document_dependencies(&self) -> &DependencyIndex { @@ -482,6 +519,10 @@ impl Db { } pub fn snapshot(&self) -> Arc { + if let Some(s) = self.snapshot.borrow().clone() { + return s; + } + self.settle(); if let Some(s) = self.snapshot.borrow().clone() { return s; } diff --git a/core/engine/src/workspace/graph/queries.rs b/core/engine/src/workspace/graph/queries.rs index c2a4f145..3bb75b73 100644 --- a/core/engine/src/workspace/graph/queries.rs +++ b/core/engine/src/workspace/graph/queries.rs @@ -79,11 +79,9 @@ impl Db { imports: &[Arc], ) -> HashMap, VariableType> { let mut out = HashMap::new(); + let view = ReadView::Dictionaries(self.dictionary_view(imports)); for import in imports { - self.graph_dep_record_view( - import, - ReadView::Dictionaries(self.dictionary_view(import)), - ); + self.graph_dep_record_view(import, view.clone()); } for entry in self.graph_dictionary_blocks(imports) { out.insert(entry.ir.name.clone(), entry.ir.enum_type()); diff --git a/core/engine/src/workspace/mod.rs b/core/engine/src/workspace/mod.rs index c5c2794b..a5519af2 100644 --- a/core/engine/src/workspace/mod.rs +++ b/core/engine/src/workspace/mod.rs @@ -149,8 +149,8 @@ impl Workspace { self.db.all_diagnostics() } - pub fn affected_by(&self, paths: &[&str]) -> Vec> { - self.db.affected_by(paths) + pub fn changes_since(&self, cursor: u64) -> (u64, Vec>) { + self.db.changes_since(cursor) } pub fn document_uses(&self, path: &str) -> Vec { diff --git a/core/engine/src/workspace/reads.rs b/core/engine/src/workspace/reads.rs index 9c5961c8..a371e0e2 100644 --- a/core/engine/src/workspace/reads.rs +++ b/core/engine/src/workspace/reads.rs @@ -1,8 +1,5 @@ -use std::collections::VecDeque; use std::sync::Arc; -use ahash::HashSet; - use crate::policy::ir::DictionaryIr; use crate::workspace::db::Db; use crate::workspace::graph::SignatureResolution; @@ -13,43 +10,55 @@ pub(crate) enum ReadView { Signature(SignatureResolution), } -#[derive(Clone, PartialEq)] -pub(crate) enum DictionaryView { +#[derive(Clone, Copy, PartialEq)] +pub(crate) enum ImportKind { Missing, Graph, - Policy(Vec<(Arc, Arc, Arc)>), + Policy, +} + +#[derive(Clone, PartialEq)] +pub(crate) struct DictionaryView { + imports: Vec<(Arc, ImportKind)>, + entries: Vec<(Arc, Arc, Arc)>, } impl Db { - pub(crate) fn dictionary_view(&self, path: &Arc) -> DictionaryView { + pub(crate) fn dictionary_view(&self, imports: &[Arc]) -> DictionaryView { let snap = self.snapshot(); - if !snap.all_parsed.contains_key(path) { - return match snap.graphs.contains_key(path) { - true => DictionaryView::Graph, - false => DictionaryView::Missing, - }; + DictionaryView { + imports: imports + .iter() + .map(|import| { + let kind = match ( + snap.all_parsed.contains_key(import), + snap.graphs.contains_key(import), + ) { + (true, _) => ImportKind::Policy, + (false, true) => ImportKind::Graph, + (false, false) => ImportKind::Missing, + }; + (import.clone(), kind) + }) + .collect(), + entries: self + .graph_dictionary_blocks(imports) + .into_iter() + .map(|entry| (entry.policy_path, entry.block_id, entry.ir)) + .collect(), } - let mut visited: HashSet> = HashSet::default(); - let mut queue: VecDeque> = VecDeque::from([path.clone()]); - let mut entries = Vec::new(); - while let Some(current) = queue.pop_front() { - if !visited.insert(current.clone()) { - continue; - } - let Some(parsed) = snap.all_parsed.get(¤t) else { - continue; - }; - for block in &parsed.policy.dictionaries { - entries.push((current.clone(), block.id.clone(), block.ir.clone())); - } - queue.extend(parsed.policy.imports().iter().cloned()); - } - DictionaryView::Policy(entries) } pub(crate) fn view_holds(&self, path: &Arc, view: &ReadView) -> bool { match view { - ReadView::Dictionaries(recorded) => self.dictionary_view(path) == *recorded, + ReadView::Dictionaries(recorded) => { + let imports: Vec> = recorded + .imports + .iter() + .map(|(import, _)| import.clone()) + .collect(); + self.dictionary_view(&imports) == *recorded + } ReadView::Signature(recorded) => { self.graph_dep_frame_push(path); let current = self.decision_signature(path);