1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
//! The work one derivation pass does: the two sweeps, the watch-set reconcile,
//! the live `bl` projection and the ops tail (DESIGN §7.2, §5.1 #2/#4, §4.2).
//!
//! Split from [`super::derive`] per §12's budget — that file is the *pass*
//! (what is dirty, what is due, what gets published), this one is the *work*:
//! the sweeps, the fetches, and — since bl-4b28 — re-deriving one root, which
//! is the work every one of them ends in.
//! All of it runs on the derivation worker, never on the frame: the ball fetch
//! is a directory walk per project and the full sweep re-derives every
//! workspace, and those are exactly the costs that used to stall the window
//! (bl-ee0a).
use super::super::desired_watches;
use super::super::drift::{self, Drift};
use super::super::snapshot::growth_between;
use super::Deriver;
use crate::budgets::Scope;
use crate::fs_watcher::RootKind;
use crate::opslog::{self, OpRow};
use crate::projects;
use crate::projects::join;
use crate::watch::Mark;
use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
impl Deriver {
/// Re-enumerate the workspace set, reconcile the watch set to it, prune
/// snapshots for vanished workspaces, and mark newly-appeared workspaces for
/// an initial derive (§7.2 cheap enumeration; §7.3 re-primed-clone rebuild).
///
/// `mark` is why this reconcile is running. Under [`Mark::Sweep`] a
/// workspace whose *membership changed here* appeared or vanished with no
/// enumeration-root event announcing it — drift, returned as
/// [`Drift::Unenumerated`]. Keying on the membership delta (not on "has no
/// snapshot yet") is what keeps this a one-shot finding: a workspace that
/// simply fails to derive stays un-snapshotted forever and must not
/// re-accuse the watcher every 2 s.
///
/// …and only where an announcement was *possible*: the delta is filtered to
/// the enumeration roots [`announceable`](Self::announceable) says were armed
/// before this pass re-armed anything. On a pristine world the names root
/// does not exist until the start flow's own `create_dir_all` founds it
/// (§8.1 `EnsureWorkspace`), so the first workspace is born under a directory
/// nothing was watching — the watcher dropped no event, it was never given
/// one, and accusing it painted a healthy first run red forever (bl-f726).
pub(super) fn reconcile(&mut self, mark: Mark) -> Vec<Drift> {
let before: BTreeSet<PathBuf> = self.workspaces.iter().map(|w| w.path.clone()).collect();
let announceable = self.announceable();
self.workspaces = crate::binding::workspaces(&self.roots.yog_data, &self.roots.lernie_data);
let desired = desired_watches(&self.roots, &self.workspaces);
crate::state::lock_watchset(&self.watches).reconcile(&desired);
let known: BTreeSet<PathBuf> = self.workspaces.iter().map(|w| w.path.clone()).collect();
if before != known {
self.changed = true;
}
self.trees.retain(|path, _| known.contains(path));
let mut missing = Vec::new();
for w in &self.workspaces {
if !self.trees.contains_key(&w.path) {
missing.push((w.path.clone(), mark));
}
}
self.schedule.mark(missing);
// The workspace set is a join axis (§3.5): a freshly named workspace (the
// start flow's `lernie new`) that lands via a NamesRoot event must re-bind
// the balls at once, else the just-claimed ball renders claimed-elsewhere
// until the 15 s sweep. Rebuild the join over the already-fetched balls.
let cloned: Vec<PathBuf> = self.balls_by_project.keys().cloned().collect();
self.rebuild_join(&cloned);
if mark != Mark::Sweep {
return Vec::new();
}
before
.symmetric_difference(&known)
// The three roots are flat (§3.1), so a workspace's parent *is* the
// enumeration root that would have announced it.
.filter(|ws| ws.parent().is_some_and(|root| announceable.contains(root)))
.cloned()
.map(Drift::Unenumerated)
.collect()
}
/// The enumeration roots a watcher is armed on **right now** — read before
/// [`reconcile`](Self::reconcile) re-arms, so it answers "was an announcement
/// possible over the interval this delta happened in?".
///
/// Derived from `desired_watches` rather than restating the root list: that
/// function is the single home of which root is which kind (§7.1), and a
/// second copy here would be the fourth place the names root is spelled.
fn announceable(&self) -> BTreeSet<PathBuf> {
let set = crate::state::lock_watchset(&self.watches);
desired_watches(&self.roots, &[])
.into_iter()
.filter(|(root, kind)| {
matches!(kind, RootKind::NamesRoot | RootKind::WorkspacesRoot)
&& set.watches(root, *kind)
})
.map(|(root, _)| root)
.collect()
}
/// The 2 s cheap sweep (§7.2): reconcile, then the *targeted* liveness
/// re-probe — for each workspace holding a Live/InFlight agent, evict its
/// agents' cached lock observations (§10 eager refresh, so a silently
/// released flock is caught) and mark it for re-derivation.
///
/// The liveness half is a poll of **process state**, not of the filesystem:
/// a released flock emits no event for any watcher to drop, so its roots are
/// marked [`Mark::Poll`] and a change under them is the poll working, never
/// drift. That separation is why the two sweeps are justified separately
/// (§7.2).
pub(super) fn cheap_sweep(&mut self) -> Vec<Drift> {
let found = self.reconcile(Mark::Sweep);
self.reprobe_live();
found
}
/// The 15 s full sweep (§7.2): reconcile, re-fetch every project's balls and
/// the ops tail (the fetch cadence's floor), and mark every workspace
/// [`Mark::Sweep`] — so a dropped inotify event costs ≤15 s of latency,
/// never divergence, **and is named** when the re-derivation proves one
/// happened.
pub(super) fn full_sweep(&mut self) -> Vec<Drift> {
let found = self.reconcile(Mark::Sweep);
self.refresh_balls();
self.refresh_ops();
// The §5.1 #35 windows ride the fetch cadence's floor for the same
// reason the balls do: one hand-edited world-global file, re-read on
// the sweep rather than watched (`adopt_windows`).
self.adopt_windows();
let all: Vec<(PathBuf, Mark)> = self
.workspaces
.iter()
.map(|w| (w.path.clone(), Mark::Sweep))
.collect();
self.schedule.mark(all);
found
}
/// Re-derive one workspace through the held probe stack, replacing its
/// snapshot iff it actually changed (`GitTree: PartialEq` suppresses no-op
/// repaints, §7.2) and recording what grew (§7.2 growth, bl-ee0a). A read
/// failure keeps the last good snapshot.
pub(super) fn rederive(&mut self, workspace: &Path) -> bool {
// The `steps/` fold first, and **outside the tree's equality gate**
// (bl-9dd4): a step's `response.json` growing is spend that changed
// while every git ref stood still, so a fold behind `old == tree` would
// freeze the spend column at whatever it read when the refs last moved.
let billed = self.refold_bills(workspace);
let Ok(tree) = self.probes.derive(workspace) else {
return billed;
};
let old = self.trees.get(workspace);
if old == Some(&tree) {
return billed;
}
self.growth.extend(growth_between(workspace, old, &tree));
self.trees.insert(workspace.to_path_buf(), tree);
self.changed = true;
true
}
/// Re-walk one workspace's `steps/` tree, replacing its bills iff they
/// actually changed — the same no-op-repaint discipline the tree read
/// keeps, so a quiet workspace publishes nothing.
fn refold_bills(&mut self, workspace: &Path) -> bool {
let bills = crate::budgets::bills(workspace, &Scope::Workspace);
if self.bills.get(workspace) == Some(&bills) {
return false;
}
self.bills.insert(workspace.to_path_buf(), bills);
self.changed = true;
true
}
/// Append this pass's findings to `ops.jsonl` and re-read the tail, so the
/// §11 activity chip carries the drift count on the very snapshot it was
/// found in. Silent when nothing was found — the normal case, and the one
/// the design is aiming at. The `ts` comes from the injected clock (§4.2),
/// which is why this crate still reads no wall clock of its own.
pub(super) fn report_drift(&mut self, found: &[Drift]) {
if found.is_empty() {
return;
}
let root = self.roots.yog_state.clone();
let cwd = root.to_string_lossy().into_owned();
for entry in drift::entries(&self.clock.stamp(), &cwd, found) {
let _ = opslog::append(&root, &entry);
}
self.refresh_ops();
}
/// Re-fetch every *visible* project's live ball list (§5.1 #2) and rebuild
/// the §3.5 join. The fetch-cadence entry point: run on the clones-root
/// dirtiness or the 15 s full sweep (§7.2), never per frame. Nested-delivery
/// clones are never walked (§5.1 #1, bl-e3e7).
pub(super) fn refresh_balls(&mut self) {
let all = projects::enumerate(&self.roots.balls_clones);
self.projects = all.iter().map(|p| p.path.clone()).collect();
let visible: Vec<PathBuf> = projects::visible(&all)
.into_iter()
.map(|p| p.path.clone())
.collect();
// Key only the projects that list cleanly; a cloned project absent from
// the map is unlistable → an orphaned row in the join (§3.5).
self.balls_by_project = visible
.iter()
.filter_map(|p| Some((p.clone(), self.balls.live(p).ok()?)))
.collect();
self.rebuild_join(&visible);
}
/// Re-fetch **one** project after a dispatched `bl` verb (§15 Y16): its live
/// *and* closed balls (the delivered-row source, §5.1 #4 — closed is fetched
/// only here, never on the cadence), then rebuild the join and re-read the
/// ops tail. The frame reaches this by marking the project's clone dir
/// dirty; it is the same routing every other root gets, not a second
/// channel.
pub(super) fn refetch_project(&mut self, project: &Path) {
// Forgiving here (unlike the cadence fetch): a verb just ran against this
// project, so an unreadable listing is transient noise, not the §3.5
// orphan signal — show no balls and let the next sweep re-derive.
let live = self.balls.live(project).unwrap_or_default();
let closed = self.balls.closed(project).unwrap_or_default();
self.balls_by_project.insert(project.to_path_buf(), live);
self.closed_by_project.insert(project.to_path_buf(), closed);
let cloned: Vec<PathBuf> = self.balls_by_project.keys().cloned().collect();
self.rebuild_join(&cloned);
self.refresh_ops();
}
/// Rebuild the join rows from the cached live + closed balls and the
/// enumerated workspaces (§3.5). The binding is claimant = workspace name
/// (§3.2); no operator identity enters. Pure over the pre-fetched caches.
fn rebuild_join(&mut self, cloned: &[PathBuf]) {
let rows = join::join(
&self.projects,
cloned,
&self.balls_by_project,
&self.closed_by_project,
&self.workspaces,
);
if rows != self.join_rows {
self.join_rows = rows;
self.changed = true;
}
}
/// Re-read the `ops.jsonl` tail (§4.2). Run on the yog-state dirtiness, the
/// full sweep, and after any dispatched verb.
///
/// Each line is projected through [`opslog::detached::fold`] first, so a
/// detached driver's captured stderr — the only evidence a fired prompt died
/// after launch (§8.1, §13.3) — reaches the row from its sink file on *this*
/// pass. The fold is read-time by construction: the sink stays the
/// authority, `ops.jsonl` is never rewritten, and a driver still writing
/// surfaces more on the next pass.
pub(super) fn refresh_ops(&mut self) {
let root = self.roots.yog_state.clone();
let rows: Vec<OpRow> = opslog::tail(&root, opslog::OPS_TAIL)
.iter()
.map(|entry| OpRow::from(&opslog::detached::fold(&root, entry)))
.collect();
if rows != self.ops {
self.ops = rows;
self.changed = true;
}
}
}