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.
This commit is contained in:
Stefan
2026-10-01 21:15:43 +02:00
parent 03aa88c41b
commit 2b892e15de
7 changed files with 149 additions and 78 deletions
+5 -1
View File
@@ -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<string>
}
export interface PolicyValueOption {
value: string;
label: string;
@@ -400,7 +404,7 @@ export declare class Workspace {
isGraph(path: string): boolean
uncheckedNodes(path: string): Array<string>
paths(): Array<string>
affectedBy(paths: Array<string>): Array<string>
changesSince(cursor: number): PolicyChanges
updateBlock(req: PolicyUpdateBlockRequest): void
removeBlock(req: PolicyRemoveBlockRequest): boolean
diagnostics(policyPath: string, maxDiagnostics?: number | undefined | null): Array<PolicyDiagnostic>
+12 -7
View File
@@ -12,6 +12,12 @@ use zen_engine::workspace;
type ResolverRef = FunctionRef<FnArgs<(String, Value)>, Option<String>>;
#[napi(object)]
pub struct PolicyChanges {
pub cursor: i64,
pub paths: Vec<String>,
}
#[napi(object)]
pub struct PolicyExpressionCursor {
pub policy_path: String,
@@ -627,13 +633,12 @@ impl Workspace {
}
#[napi]
pub fn affected_by(&self, paths: Vec<String>) -> Vec<String> {
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]
+22 -8
View File
@@ -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<Arc<str>> {
let changed: Vec<Arc<str>> = paths.iter().map(|path| Arc::from(*path)).collect();
pub(crate) fn walk(&self, changed: &[Arc<str>]) -> Vec<Arc<str>> {
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
}
}
})
+68 -27
View File
@@ -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<Arc<str>>,
settling: bool,
sequence: u64,
marked: HashMap<Arc<str>, u64>,
}
#[derive(Default)]
pub(crate) struct DepFrame {
docs: HashSet<Arc<str>>,
@@ -197,6 +205,7 @@ pub struct Db {
pub(crate) graph_stack: RefCell<Vec<Arc<str>>>,
graph_dep_frames: RefCell<Vec<DepFrame>>,
dependencies: DependencyIndex,
changes: RefCell<ChangeLog>,
graph_fn_frames: RefCell<Vec<HashMap<FunctionKey, u64>>>,
function_types: RefCell<HashMap<FunctionKey, ResolvedFunction>>,
function_requests: RefCell<Vec<FunctionResolutionRequest>>,
@@ -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<str>, doc: Arc<DecisionContent>) {
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<str>,
dependency: &Arc<str>,
kind: DependencyKind,
) -> Option<ReadView> {
pub(crate) fn frozen_views(&self) -> HashMap<(Arc<str>, Arc<str>), Vec<ReadView>> {
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<str>, Arc<str>), Vec<ReadView>> = 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<Arc<str>>) {
self.snapshot();
self.settle();
let changes = self.changes.borrow();
let mut paths: Vec<Arc<str>> = 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<Snapshot> {
if let Some(s) = self.snapshot.borrow().clone() {
return s;
}
self.settle();
if let Some(s) = self.snapshot.borrow().clone() {
return s;
}
+2 -4
View File
@@ -79,11 +79,9 @@ impl Db {
imports: &[Arc<str>],
) -> HashMap<Arc<str>, 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());
+2 -2
View File
@@ -149,8 +149,8 @@ impl Workspace {
self.db.all_diagnostics()
}
pub fn affected_by(&self, paths: &[&str]) -> Vec<Arc<str>> {
self.db.affected_by(paths)
pub fn changes_since(&self, cursor: u64) -> (u64, Vec<Arc<str>>) {
self.db.changes_since(cursor)
}
pub fn document_uses(&self, path: &str) -> Vec<document_dependencies::Dependency> {
+38 -29
View File
@@ -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<str>, Arc<str>, Arc<DictionaryIr>)>),
Policy,
}
#[derive(Clone, PartialEq)]
pub(crate) struct DictionaryView {
imports: Vec<(Arc<str>, ImportKind)>,
entries: Vec<(Arc<str>, Arc<str>, Arc<DictionaryIr>)>,
}
impl Db {
pub(crate) fn dictionary_view(&self, path: &Arc<str>) -> DictionaryView {
pub(crate) fn dictionary_view(&self, imports: &[Arc<str>]) -> 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<Arc<str>> = HashSet::default();
let mut queue: VecDeque<Arc<str>> = 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(&current) 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<str>, view: &ReadView) -> bool {
match view {
ReadView::Dictionaries(recorded) => self.dictionary_view(path) == *recorded,
ReadView::Dictionaries(recorded) => {
let imports: Vec<Arc<str>> = 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);