Skip to main content

lex_vcs/
op_log.rs

1//! Persistence + DAG queries for the operation log.
2//!
3//! Layout: `<root>/ops/<op_id>.json` — one canonical-JSON file
4//! per [`OperationRecord`]. Atomic writes via tempfile + rename.
5//! Idempotent: writing an existing op_id is a no-op (content
6//! addressing guarantees the bytes match).
7//!
8//! # Packfiles (#261 slice 1)
9//!
10//! Loose-file storage is fine to ~10k ops; past that the
11//! filesystem starts to thrash. [`OpLog::repack`] consolidates
12//! loose files into deterministic, content-addressed packfiles:
13//!
14//! - `<dir>/pack-<hash>.pack`: each record framed as `[8-byte BE
15//!   length][canonical JSON]`, ops sorted by op_id within the pack.
16//! - `<dir>/pack-<hash>.idx`: JSON map of `op_id` → byte offset
17//!   into the `.pack` (offset of the length header).
18//!
19//! Pack name is the SHA-256 of the sorted op_ids, newline-joined,
20//! so the same input set always produces the same pack hash —
21//! a re-run of `lex op repack` is a no-op.
22//!
23//! [`OpLog::get`] tries loose first, falls back to scanning all
24//! `.idx` files in the directory. The write path
25//! ([`OpLog::put`]) only ever writes loose; ops migrate into
26//! packs via the explicit [`OpLog::repack`] call.
27
28use crate::canonical::hash_bytes;
29use crate::operation::{OpId, OperationRecord};
30use std::collections::{BTreeMap, BTreeSet, VecDeque};
31use std::fs;
32use std::io::{self, Read, Seek, SeekFrom, Write};
33use std::path::{Path, PathBuf};
34
35pub struct OpLog {
36    dir: PathBuf,
37}
38
39impl OpLog {
40    pub fn open(root: &Path) -> io::Result<Self> {
41        let dir = root.join("ops");
42        fs::create_dir_all(&dir)?;
43        Ok(Self { dir })
44    }
45
46    fn path(&self, op_id: &OpId) -> PathBuf {
47        self.dir.join(format!("{op_id}.json"))
48    }
49
50    /// Persist a record. Idempotent on existing op_ids (the bytes
51    /// must match by content addressing).
52    ///
53    /// Crash safety: the tempfile's data is fsync'd before rename,
54    /// so a successful return implies a durable file at the final
55    /// path. The containing directory is not fsync'd; on a crash
56    /// between rename and the directory's metadata flush, the file
57    /// can be lost. For a content-addressed log this is acceptable
58    /// — a lost record can be re-derived from the same source — but
59    /// callers that *also* persist references to the op_id (e.g.
60    /// branch heads) should fsync those refs after `put` returns.
61    pub fn put(&self, rec: &OperationRecord) -> io::Result<()> {
62        let path = self.path(&rec.op_id);
63        if path.exists() {
64            return Ok(());
65        }
66        let bytes = serde_json::to_vec(rec)
67            .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
68        let tmp = path.with_extension("json.tmp");
69        let mut f = fs::File::create(&tmp)?;
70        f.write_all(&bytes)?;
71        f.sync_all()?;
72        fs::rename(&tmp, &path)?;
73        Ok(())
74    }
75
76    pub fn get(&self, op_id: &OpId) -> io::Result<Option<OperationRecord>> {
77        let path = self.path(op_id);
78        if path.exists() {
79            let bytes = fs::read(&path)?;
80            let rec: OperationRecord = serde_json::from_slice(&bytes)
81                .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
82            return Ok(Some(rec));
83        }
84        // Loose miss — scan packfiles. Each `.idx` is a tiny JSON
85        // map; small constant cost per pack. For larger stores we
86        // could maintain an in-memory cache keyed off pack mtimes,
87        // but slice 1 keeps it simple — measure before optimizing.
88        for pack_idx in self.list_pack_indices()? {
89            let idx = PackIndex::load(&pack_idx)?;
90            if let Some(&offset) = idx.ops.get(op_id) {
91                let pack_path = pack_idx.with_extension("pack");
92                return read_packed_op(&pack_path, offset).map(Some);
93            }
94        }
95        Ok(None)
96    }
97
98    /// Walk the directory for `pack-*.idx` files. Order is whatever
99    /// the filesystem gives us — `get` doesn't depend on it (op_ids
100    /// are unique by content addressing, so the right pack wins).
101    fn list_pack_indices(&self) -> io::Result<Vec<PathBuf>> {
102        let mut out = Vec::new();
103        for entry in fs::read_dir(&self.dir)? {
104            let entry = entry?;
105            let name = match entry.file_name().into_string() {
106                Ok(s) => s,
107                Err(_) => continue,
108            };
109            if name.starts_with("pack-") && name.ends_with(".idx") {
110                out.push(entry.path());
111            }
112        }
113        Ok(out)
114    }
115
116    /// Consolidate loose op records into a deterministic, content-
117    /// addressed packfile (#261 slice 1). Returns the number of
118    /// ops moved into the new pack.
119    ///
120    /// `threshold` is the minimum number of loose ops required to
121    /// trigger a repack — under that, returns `0` and leaves the
122    /// log alone. The idea: small stores stay loose; only repack
123    /// when the file count starts to matter.
124    ///
125    /// Determinism: the pack name is the SHA-256 of the sorted
126    /// op_ids (newline-joined), so two independent runs against the
127    /// same set of loose ops produce a byte-identical pack.
128    /// Re-running on an empty loose directory is a no-op.
129    ///
130    /// Crash safety: the `.pack.tmp` and `.idx.tmp` files are
131    /// fsync'd before rename; loose files are deleted only after
132    /// both renames succeed. A crash mid-repack leaves both loose
133    /// and partial-pack files; a subsequent `get` finds the loose
134    /// version, and a subsequent `repack` cleans up.
135    pub fn repack(&self, threshold: usize) -> io::Result<usize> {
136        let loose: Vec<(OpId, PathBuf)> = self.list_loose_files()?;
137        if loose.len() < threshold {
138            return Ok(0);
139        }
140        // Sort ops deterministically by op_id (lex order). The pack
141        // hash is the SHA-256 of those op_ids joined by newlines —
142        // same input → same name.
143        let mut ops: Vec<(OpId, Vec<u8>)> = Vec::with_capacity(loose.len());
144        for (op_id, path) in &loose {
145            let bytes = fs::read(path)?;
146            ops.push((op_id.clone(), bytes));
147        }
148        ops.sort_by(|a, b| a.0.cmp(&b.0));
149        let mut name_input = Vec::new();
150        for (id, _) in &ops {
151            name_input.extend_from_slice(id.as_bytes());
152            name_input.push(b'\n');
153        }
154        let pack_hash = hash_bytes(&name_input);
155        let pack_path = self.dir.join(format!("pack-{pack_hash}.pack"));
156        let idx_path = self.dir.join(format!("pack-{pack_hash}.idx"));
157        if pack_path.exists() && idx_path.exists() {
158            // Same input set — pack already exists. Just clean up
159            // the loose duplicates.
160            let count = ops.len();
161            for (_, path) in &loose {
162                let _ = fs::remove_file(path);
163            }
164            return Ok(count);
165        }
166
167        // Write `<pack>.pack.tmp` framed as [8-byte BE length][JSON]
168        // for each record; record offsets for the index.
169        let pack_tmp = pack_path.with_extension("pack.tmp");
170        let idx_tmp = idx_path.with_extension("idx.tmp");
171        let mut offsets: BTreeMap<OpId, u64> = BTreeMap::new();
172        {
173            let mut f = fs::File::create(&pack_tmp)?;
174            let mut cursor: u64 = 0;
175            for (op_id, bytes) in &ops {
176                offsets.insert(op_id.clone(), cursor);
177                let len = bytes.len() as u64;
178                f.write_all(&len.to_be_bytes())?;
179                f.write_all(bytes)?;
180                cursor += 8 + len;
181            }
182            f.sync_all()?;
183        }
184        // Write the index. JSON for inspectability and
185        // forward-compat (we can add fields without breaking
186        // readers).
187        let idx = PackIndex { version: 1, ops: offsets };
188        idx.save(&idx_tmp)?;
189
190        fs::rename(&pack_tmp, &pack_path)?;
191        fs::rename(&idx_tmp, &idx_path)?;
192
193        // Now safe to delete the loose files — pack is durable.
194        let count = ops.len();
195        for (_, path) in &loose {
196            let _ = fs::remove_file(path);
197        }
198        Ok(count)
199    }
200
201    /// Remove every op_id in `victims` from the log, across both
202    /// loose files and packfiles (#261 slice 2). Used by
203    /// `lex op gc` after a retention plan identifies which ops to
204    /// drop. Idempotent — calling twice with the same set is a
205    /// no-op on the second pass.
206    ///
207    /// Pack handling: any pack containing one or more victims is
208    /// rewritten to a new content-addressed pack with only the
209    /// surviving ops; the old pack and its index file are deleted.
210    /// A pack whose every op is a victim is deleted outright.
211    ///
212    /// Returns the count of ops actually removed (loose files
213    /// deleted + packed ops dropped). Pre-existing absences don't
214    /// contribute.
215    pub fn evict(&self, victims: &BTreeSet<OpId>) -> io::Result<usize> {
216        if victims.is_empty() {
217            return Ok(0);
218        }
219        let mut removed = 0;
220        // Loose files: just delete the matching `<op_id>.json`.
221        for (op_id, path) in self.list_loose_files()? {
222            if victims.contains(&op_id) {
223                match fs::remove_file(&path) {
224                    Ok(()) => removed += 1,
225                    Err(e) if e.kind() == io::ErrorKind::NotFound => {}
226                    Err(e) => return Err(e),
227                }
228            }
229        }
230        // Packs: rewrite each affected pack with only surviving ops.
231        for pack_idx in self.list_pack_indices()? {
232            let idx = PackIndex::load(&pack_idx)?;
233            let pack_path = pack_idx.with_extension("pack");
234            let touched = idx.ops.keys().any(|op_id| victims.contains(op_id));
235            if !touched {
236                continue;
237            }
238            // Read every surviving op out, then drop the old pack.
239            let mut survivors: Vec<(OpId, Vec<u8>)> = Vec::new();
240            for (op_id, &offset) in &idx.ops {
241                if victims.contains(op_id) {
242                    removed += 1;
243                    continue;
244                }
245                let rec = read_packed_op(&pack_path, offset)?;
246                let bytes = serde_json::to_vec(&rec)
247                    .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
248                survivors.push((op_id.clone(), bytes));
249            }
250            // Delete old pack + idx first; we re-emit a fresh
251            // (different-hash) pack from the survivors below.
252            let _ = fs::remove_file(&pack_path);
253            let _ = fs::remove_file(&pack_idx);
254            if survivors.is_empty() {
255                continue;
256            }
257            self.write_pack_from_survivors(survivors)?;
258        }
259        Ok(removed)
260    }
261
262    /// Helper: write a new content-addressed pack from
263    /// already-serialized op bytes. Same shape as
264    /// [`Self::repack`]'s output path; factored out so
265    /// [`Self::evict`] can reuse it.
266    fn write_pack_from_survivors(
267        &self,
268        mut ops: Vec<(OpId, Vec<u8>)>,
269    ) -> io::Result<()> {
270        ops.sort_by(|a, b| a.0.cmp(&b.0));
271        let mut name_input = Vec::new();
272        for (id, _) in &ops {
273            name_input.extend_from_slice(id.as_bytes());
274            name_input.push(b'\n');
275        }
276        let pack_hash = hash_bytes(&name_input);
277        let pack_path = self.dir.join(format!("pack-{pack_hash}.pack"));
278        let idx_path = self.dir.join(format!("pack-{pack_hash}.idx"));
279        if pack_path.exists() && idx_path.exists() {
280            return Ok(());
281        }
282        let pack_tmp = pack_path.with_extension("pack.tmp");
283        let idx_tmp = idx_path.with_extension("idx.tmp");
284        let mut offsets: BTreeMap<OpId, u64> = BTreeMap::new();
285        {
286            let mut f = fs::File::create(&pack_tmp)?;
287            let mut cursor: u64 = 0;
288            for (op_id, bytes) in &ops {
289                offsets.insert(op_id.clone(), cursor);
290                let len = bytes.len() as u64;
291                f.write_all(&len.to_be_bytes())?;
292                f.write_all(bytes)?;
293                cursor += 8 + len;
294            }
295            f.sync_all()?;
296        }
297        let idx = PackIndex { version: 1, ops: offsets };
298        idx.save(&idx_tmp)?;
299        fs::rename(&pack_tmp, &pack_path)?;
300        fs::rename(&idx_tmp, &idx_path)?;
301        Ok(())
302    }
303
304    /// Enumerate every loose `<op_id>.json` in the ops directory.
305    /// Used by [`Self::repack`] and [`Self::list_all`].
306    fn list_loose_files(&self) -> io::Result<Vec<(OpId, PathBuf)>> {
307        let mut out = Vec::new();
308        for entry in fs::read_dir(&self.dir)? {
309            let entry = entry?;
310            let name = match entry.file_name().into_string() {
311                Ok(s) => s,
312                Err(_) => continue,
313            };
314            if let Some(id) = name.strip_suffix(".json") {
315                if !id.starts_with("pack-") {
316                    out.push((id.to_string(), entry.path()));
317                }
318            }
319        }
320        Ok(out)
321    }
322
323    /// Remove a record from the log. Used by [`crate::migrate`] to
324    /// delete the old `<op_id>.json` files after a format migration
325    /// has written their replacements. Idempotent on missing files.
326    ///
327    /// **Not** part of the day-to-day op-log API — the log is
328    /// append-only by design (#129). The only legitimate caller is
329    /// the migration tool, which is supervising a destructive,
330    /// `--confirm`-gated batch.
331    pub fn delete(&self, op_id: &OpId) -> io::Result<()> {
332        let path = self.path(op_id);
333        match fs::remove_file(&path) {
334            Ok(()) => Ok(()),
335            Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
336            Err(e) => Err(e),
337        }
338    }
339
340    /// Walk parents transitively. Newest-first, BFS, dedup'd by op_id.
341    /// Stops at parentless ops or after `limit` records.
342    pub fn walk_back(
343        &self,
344        head: &OpId,
345        limit: Option<usize>,
346    ) -> io::Result<Vec<OperationRecord>> {
347        let mut out = Vec::new();
348        let mut seen = BTreeSet::new();
349        let mut frontier: VecDeque<OpId> = VecDeque::from([head.clone()]);
350        while let Some(id) = frontier.pop_back() {
351            if !seen.insert(id.clone()) {
352                continue;
353            }
354            if let Some(rec) = self.get(&id)? {
355                // Push parents before recording so traversal order is
356                // a stable BFS-by-discovery: children-first, then their
357                // parents, parents of those, etc.
358                for p in &rec.op.parents {
359                    if !seen.contains(p) {
360                        frontier.push_front(p.clone());
361                    }
362                }
363                out.push(rec);
364                if let Some(n) = limit {
365                    if out.len() >= n {
366                        break;
367                    }
368                }
369            }
370        }
371        Ok(out)
372    }
373
374    /// Same set as walk_back but oldest-first, and **topological**: every op
375    /// comes after all of its parents. Used by branch_head (and every other
376    /// consumer) for left-to-right transition replay, so this is the one
377    /// definition of "the order a head is replayed in".
378    ///
379    /// `walk_back` is a breadth-first walk, and its reverse is *not* a
380    /// topological order once a merge joins lines of different lengths: an op
381    /// reachable by a short path is emitted before its own descendant on the
382    /// long one, so it lands after that descendant here. Replaying it then
383    /// overwrote a later change with an earlier one (#1062). The order is now
384    /// repaired by [`Self::linearize`], which keeps the BFS order wherever it
385    /// was already topological — every linear history, and any merge of
386    /// equal-length lines — so those replay exactly as they always did.
387    pub fn walk_forward(
388        &self,
389        head: &OpId,
390        limit: Option<usize>,
391    ) -> io::Result<Vec<OperationRecord>> {
392        let mut all = self.walk_back(head, None)?;
393        all.reverse();
394        let mut all = Self::linearize(all);
395        if let Some(n) = limit {
396            all.truncate(n);
397        }
398        Ok(all)
399    }
400
401    /// Reorder `records` into a topological order: every record comes after
402    /// all of its parents that are in the set (a parent missing from the set
403    /// — an incomplete local log — imposes no constraint, the same leniency
404    /// the walks have). Among the records ready at any point, the one listed
405    /// first in `records` goes first, so an input that is already
406    /// topological comes back unchanged, and the result is a pure function
407    /// of the input order.
408    pub fn linearize(records: Vec<OperationRecord>) -> Vec<OperationRecord> {
409        use std::cmp::Reverse;
410        use std::collections::{BinaryHeap, HashMap};
411        let index: HashMap<&str, usize> =
412            records.iter().enumerate().map(|(i, r)| (r.op_id.as_str(), i)).collect();
413        // Fast path (every linear history): nothing to reorder.
414        let already = records.iter().enumerate().all(|(i, r)| {
415            r.op.parents.iter().all(|p| index.get(p.as_str()).is_none_or(|&j| j < i))
416        });
417        if already {
418            return records;
419        }
420        let n = records.len();
421        let mut indegree = vec![0usize; n];
422        let mut children: Vec<Vec<usize>> = vec![Vec::new(); n];
423        for (i, r) in records.iter().enumerate() {
424            let mut seen: Vec<usize> = Vec::new();
425            for p in &r.op.parents {
426                if let Some(&j) = index.get(p.as_str()) {
427                    if j != i && !seen.contains(&j) {
428                        seen.push(j);
429                        indegree[i] += 1;
430                        children[j].push(i);
431                    }
432                }
433            }
434        }
435        let mut ready: BinaryHeap<Reverse<usize>> =
436            (0..n).filter(|&i| indegree[i] == 0).map(Reverse).collect();
437        let mut order: Vec<usize> = Vec::with_capacity(n);
438        while let Some(Reverse(i)) = ready.pop() {
439            order.push(i);
440            for &c in &children[i] {
441                indegree[c] -= 1;
442                if indegree[c] == 0 {
443                    ready.push(Reverse(c));
444                }
445            }
446        }
447        // A cycle is impossible for content-addressed ops; if a corrupt log
448        // ever produced one, keep the leftovers rather than dropping them.
449        if order.len() < n {
450            let placed: BTreeSet<usize> = order.iter().copied().collect();
451            order.extend((0..n).filter(|i| !placed.contains(i)));
452        }
453        let mut slots: Vec<Option<OperationRecord>> = records.into_iter().map(Some).collect();
454        order.into_iter().filter_map(|i| slots[i].take()).collect()
455    }
456
457    /// Whether `records` is exactly a continuation of `since`: every record
458    /// descends from `since` through parents that are themselves in
459    /// `records` (or are `since`), so none of them is an ancestor of `since`
460    /// and none has an ancestor outside the set. This is the shape
461    /// [`Self::walk_forward_since`] returns for a plain fast-forward, and
462    /// NOT the shape it returns for a merge: there it also walks the second
463    /// parent's history back to genesis, which is history `since` already
464    /// contains (#1062). Replaying such a set on top of a state computed for
465    /// `since` re-applies old changes over newer ones.
466    pub fn continues_from(records: &[OperationRecord], since: &OpId) -> bool {
467        use std::collections::HashMap;
468        let index: HashMap<&str, usize> =
469            records.iter().enumerate().map(|(i, r)| (r.op_id.as_str(), i)).collect();
470        // 0 = unvisited, 1 = in progress, 2 = descends from `since`, 3 = does not.
471        let mut state = vec![0u8; records.len()];
472        for start in 0..records.len() {
473            let mut stack = vec![start];
474            while let Some(&i) = stack.last() {
475                match state[i] {
476                    2 | 3 => {
477                        stack.pop();
478                        continue;
479                    }
480                    _ => {}
481                }
482                state[i] = 1;
483                let mut pending = false;
484                let mut reaches = false;
485                let parents = &records[i].op.parents;
486                if parents.is_empty() {
487                    state[i] = 3;
488                    stack.pop();
489                    continue;
490                }
491                for p in parents {
492                    if p == since {
493                        reaches = true;
494                    } else if let Some(&j) = index.get(p.as_str()) {
495                        match state[j] {
496                            2 => reaches = true,
497                            0 => {
498                                stack.push(j);
499                                pending = true;
500                            }
501                            // 3: an ancestor of `since` (or unrooted): the
502                            // set is not a pure continuation.
503                            _ => return false,
504                        }
505                    } else {
506                        // A parent that is neither `since` nor in the set.
507                        return false;
508                    }
509                }
510                if pending {
511                    continue;
512                }
513                state[i] = if reaches { 2 } else { 3 };
514                if state[i] == 3 {
515                    return false;
516                }
517                stack.pop();
518            }
519        }
520        true
521    }
522
523    /// Like [`Self::walk_forward`], but bounded: walk from `head` back
524    /// toward genesis and stop as soon as `since` is reached, without
525    /// visiting `since`'s own parents or including `since` itself in the
526    /// result. Returns oldest-first, suitable for incrementally
527    /// extending a transition map already computed as of `since`.
528    ///
529    /// Returns `Ok(None)` if `since` is never reached (not an ancestor
530    /// of `head` — e.g. after a branch reset or a merge that reordered
531    /// history): callers should fall back to a full `walk_forward` in
532    /// that case, since there is nothing valid to incrementally extend.
533    ///
534    /// This is the piece `walk_forward` itself doesn't provide: its own
535    /// `limit` truncates the *result* after a full walk_back to genesis
536    /// has already completed (see its body above), so it can't turn an
537    /// O(N)-in-total-history walk into an O(ops since a checkpoint) one.
538    /// `head == since` returns `Ok(Some(vec![]))` without touching the
539    /// op log at all.
540    pub fn walk_forward_since(
541        &self,
542        head: &OpId,
543        since: &OpId,
544    ) -> io::Result<Option<Vec<OperationRecord>>> {
545        if head == since {
546            return Ok(Some(Vec::new()));
547        }
548        let mut out = Vec::new();
549        let mut seen = BTreeSet::new();
550        let mut frontier: VecDeque<OpId> = VecDeque::from([head.clone()]);
551        let mut found = false;
552        while let Some(id) = frontier.pop_back() {
553            if !seen.insert(id.clone()) {
554                continue;
555            }
556            if id == *since {
557                found = true;
558                continue; // boundary: don't include it, don't descend into its parents
559            }
560            if let Some(rec) = self.get(&id)? {
561                for p in &rec.op.parents {
562                    if !seen.contains(p) {
563                        frontier.push_front(p.clone());
564                    }
565                }
566                out.push(rec);
567            }
568        }
569        if !found {
570            return Ok(None);
571        }
572        out.reverse();
573        Ok(Some(out))
574    }
575
576    /// Common ancestor of two op_ids in the DAG.
577    ///
578    /// On tree-shaped histories and chain merges this is the
579    /// **lowest** common ancestor — the closest shared op. On
580    /// criss-cross merges (two ops each with two parents from
581    /// independent histories) there can be multiple
582    /// incomparable common ancestors; this picks one
583    /// deterministically (the first hit when traversing `b`'s
584    /// ancestors newest-first), but not via a recursive merge.
585    /// `None` if no shared ancestor exists.
586    ///
587    /// Tier-1 merge in #129 covers linear and tree-shaped
588    /// histories; criss-cross resolution is deferred to a
589    /// future tier (Git's `recursive` strategy is the reference).
590    pub fn lca(&self, a: &OpId, b: &OpId) -> io::Result<Option<OpId>> {
591        let a_anc: BTreeSet<OpId> = self
592            .walk_back(a, None)?
593            .into_iter()
594            .map(|r| r.op_id)
595            .collect();
596        // Walk b's ancestors newest-first; first hit is the deepest
597        // common ancestor on tree-shaped histories. In criss-cross
598        // DAGs this picks deterministically but not via recursive
599        // resolution — see the doc comment above.
600        for rec in self.walk_back(b, None)? {
601            if a_anc.contains(&rec.op_id) {
602                return Ok(Some(rec.op_id));
603            }
604        }
605        Ok(None)
606    }
607
608    /// Every record in the log. Order is whatever the directory
609    /// listing produces — undefined and not stable. Used by the
610    /// [`crate::predicate`] evaluator when no narrower candidate
611    /// set is available.
612    pub fn list_all(&self) -> io::Result<Vec<OperationRecord>> {
613        let mut out = Vec::new();
614        let mut seen: BTreeSet<OpId> = BTreeSet::new();
615        // Loose first so dedup wins for them on collision (loose
616        // and pack should never both exist for the same op_id post-
617        // repack, but during an interrupted repack both can be
618        // present transiently).
619        for (id, _) in self.list_loose_files()? {
620            if let Some(rec) = self.get(&id)? {
621                if seen.insert(rec.op_id.clone()) {
622                    out.push(rec);
623                }
624            }
625        }
626        for pack_idx in self.list_pack_indices()? {
627            let idx = PackIndex::load(&pack_idx)?;
628            let pack_path = pack_idx.with_extension("pack");
629            for (op_id, &offset) in &idx.ops {
630                if seen.insert(op_id.clone()) {
631                    out.push(read_packed_op(&pack_path, offset)?);
632                }
633            }
634        }
635        Ok(out)
636    }
637
638    /// Ops in `head`'s history that are not in `base`'s history.
639    /// `base = None` means "include all of head's history" (used for
640    /// independent-histories case where the LCA is None).
641    pub fn ops_since(
642        &self,
643        head: &OpId,
644        base: Option<&OpId>,
645    ) -> io::Result<Vec<OperationRecord>> {
646        let exclude: BTreeSet<OpId> = match base {
647            Some(b) => self
648                .walk_back(b, None)?
649                .into_iter()
650                .map(|r| r.op_id)
651                .collect(),
652            None => BTreeSet::new(),
653        };
654        Ok(self
655            .walk_back(head, None)?
656            .into_iter()
657            .filter(|r| !exclude.contains(&r.op_id))
658            .collect())
659    }
660}
661
662/// Sidecar index for a packfile. Maps `op_id` to the byte offset
663/// of the record's length header inside the `.pack`. JSON for
664/// inspectability and forward-compat.
665#[derive(serde::Serialize, serde::Deserialize)]
666struct PackIndex {
667    version: u32,
668    ops: BTreeMap<OpId, u64>,
669}
670
671impl PackIndex {
672    fn load(path: &Path) -> io::Result<Self> {
673        let bytes = fs::read(path)?;
674        serde_json::from_slice(&bytes)
675            .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))
676    }
677
678    fn save(&self, path: &Path) -> io::Result<()> {
679        let bytes = serde_json::to_vec(self)
680            .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
681        let mut f = fs::File::create(path)?;
682        f.write_all(&bytes)?;
683        f.sync_all()?;
684        Ok(())
685    }
686}
687
688/// Read one record from a packfile at `offset`. The record is
689/// framed as `[8-byte BE length][canonical JSON]`.
690fn read_packed_op(pack_path: &Path, offset: u64) -> io::Result<OperationRecord> {
691    let mut f = fs::File::open(pack_path)?;
692    f.seek(SeekFrom::Start(offset))?;
693    let mut len_buf = [0u8; 8];
694    f.read_exact(&mut len_buf)?;
695    let len = u64::from_be_bytes(len_buf) as usize;
696    let mut buf = vec![0u8; len];
697    f.read_exact(&mut buf)?;
698    serde_json::from_slice(&buf)
699        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))
700}
701
702#[cfg(test)]
703mod tests {
704    use super::*;
705    use crate::operation::{Operation, OperationKind, StageTransition};
706    use std::collections::{BTreeMap, BTreeSet};
707
708    fn add_op() -> OperationRecord {
709        let op = Operation::new(
710            OperationKind::AddFunction {
711                sig_id: "fac::Int->Int".into(),
712                stage_id: "abc123".into(),
713                effects: BTreeSet::new(),
714                budget_cost: None,
715                in_file: None,
716            },
717            [],
718        );
719        OperationRecord::new(
720            op,
721            StageTransition::Create {
722                sig_id: "fac::Int->Int".into(),
723                stage_id: "abc123".into(),
724            },
725        )
726    }
727
728    fn modify_op(parent: &OpId, sig: &str, from: &str, to: &str) -> OperationRecord {
729        let op = Operation::new(
730            OperationKind::ModifyBody {
731                sig_id: sig.into(),
732                from_stage_id: from.into(),
733                to_stage_id: to.into(),
734                from_budget: None,
735                to_budget: None,
736                to_sig_id: None,
737            },
738            [parent.clone()],
739        );
740        OperationRecord::new(
741            op,
742            StageTransition::Replace {
743                sig_id: sig.into(),
744                from: from.into(),
745                to: to.into(),
746            },
747        )
748    }
749
750    #[test]
751    fn put_then_get_round_trips() {
752        let tmp = tempfile::tempdir().unwrap();
753        let log = OpLog::open(tmp.path()).unwrap();
754        let rec = add_op();
755        log.put(&rec).unwrap();
756        let back = log.get(&rec.op_id).unwrap().unwrap();
757        assert_eq!(back, rec);
758    }
759
760    #[test]
761    fn put_is_idempotent() {
762        let tmp = tempfile::tempdir().unwrap();
763        let log = OpLog::open(tmp.path()).unwrap();
764        let rec = add_op();
765        log.put(&rec).unwrap();
766        log.put(&rec).unwrap(); // second write is a no-op
767        assert!(log.get(&rec.op_id).unwrap().is_some());
768    }
769
770    #[test]
771    fn get_missing_returns_none() {
772        let tmp = tempfile::tempdir().unwrap();
773        let log = OpLog::open(tmp.path()).unwrap();
774        assert!(log.get(&"deadbeef".to_string()).unwrap().is_none());
775    }
776
777    #[test]
778    fn walk_back_returns_newest_first() {
779        let tmp = tempfile::tempdir().unwrap();
780        let log = OpLog::open(tmp.path()).unwrap();
781        let a = add_op();
782        log.put(&a).unwrap();
783        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
784        log.put(&b).unwrap();
785        let c = modify_op(&b.op_id, "fac::Int->Int", "def456", "789aaa");
786        log.put(&c).unwrap();
787
788        let walked = log.walk_back(&c.op_id, None).unwrap();
789        let ids: Vec<_> = walked.iter().map(|r| r.op_id.as_str()).collect();
790        assert_eq!(
791            ids,
792            vec![c.op_id.as_str(), b.op_id.as_str(), a.op_id.as_str()]
793        );
794    }
795
796    #[test]
797    fn walk_forward_returns_oldest_first() {
798        let tmp = tempfile::tempdir().unwrap();
799        let log = OpLog::open(tmp.path()).unwrap();
800        let a = add_op();
801        log.put(&a).unwrap();
802        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
803        log.put(&b).unwrap();
804
805        let walked = log.walk_forward(&b.op_id, None).unwrap();
806        let ids: Vec<_> = walked.iter().map(|r| r.op_id.as_str()).collect();
807        assert_eq!(ids, vec![a.op_id.as_str(), b.op_id.as_str()]);
808    }
809
810    #[test]
811    fn walk_forward_since_returns_only_ops_after_the_boundary() {
812        let tmp = tempfile::tempdir().unwrap();
813        let log = OpLog::open(tmp.path()).unwrap();
814        let a = add_op();
815        log.put(&a).unwrap();
816        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
817        log.put(&b).unwrap();
818        let c = modify_op(&b.op_id, "fac::Int->Int", "def456", "789aaa");
819        log.put(&c).unwrap();
820
821        // Everything strictly after `a`: b, c, oldest-first.
822        let since_a = log.walk_forward_since(&c.op_id, &a.op_id).unwrap().unwrap();
823        let ids: Vec<_> = since_a.iter().map(|r| r.op_id.as_str()).collect();
824        assert_eq!(ids, vec![b.op_id.as_str(), c.op_id.as_str()]);
825
826        // Everything strictly after `b`: just c.
827        let since_b = log.walk_forward_since(&c.op_id, &b.op_id).unwrap().unwrap();
828        let ids: Vec<_> = since_b.iter().map(|r| r.op_id.as_str()).collect();
829        assert_eq!(ids, vec![c.op_id.as_str()]);
830    }
831
832    #[test]
833    fn walk_forward_since_head_equals_since_returns_empty_without_touching_the_log() {
834        let tmp = tempfile::tempdir().unwrap();
835        let log = OpLog::open(tmp.path()).unwrap();
836        let a = add_op();
837        log.put(&a).unwrap();
838
839        // Deliberately pass an op_id that was never `put` -- if this
840        // took the "walk and look for it" path it would return `None`
841        // (not found). The `head == since` fast path must short-circuit
842        // before ever touching the log.
843        let ghost = "never-written-anywhere".to_string();
844        let result = log.walk_forward_since(&ghost, &ghost).unwrap();
845        assert_eq!(result, Some(Vec::new()));
846    }
847
848    #[test]
849    fn walk_forward_since_returns_none_when_boundary_is_not_an_ancestor() {
850        let tmp = tempfile::tempdir().unwrap();
851        let log = OpLog::open(tmp.path()).unwrap();
852        let a = add_op();
853        log.put(&a).unwrap();
854        // A second root with different content (distinct sig_id), so it
855        // gets a different content-addressed op_id and shares no
856        // history with `a` -- `add_op()` alone is parameterless and
857        // would collide with itself.
858        let op = Operation::new(
859            OperationKind::AddFunction {
860                sig_id: "unrelated::Str->Str".into(),
861                stage_id: "zzz999".into(),
862                effects: BTreeSet::new(),
863                budget_cost: None,
864                in_file: None,
865            },
866            [],
867        );
868        let unrelated = OperationRecord::new(
869            op,
870            StageTransition::Create {
871                sig_id: "unrelated::Str->Str".into(),
872                stage_id: "zzz999".into(),
873            },
874        );
875        log.put(&unrelated).unwrap();
876        assert_ne!(a.op_id, unrelated.op_id, "test setup must produce two distinct ops");
877
878        let result = log.walk_forward_since(&a.op_id, &unrelated.op_id).unwrap();
879        assert_eq!(
880            result, None,
881            "unrelated op_id is not an ancestor of `a` -- callers must fall back to a full walk"
882        );
883    }
884
885    #[test]
886    fn lca_finds_common_ancestor() {
887        let tmp = tempfile::tempdir().unwrap();
888        let log = OpLog::open(tmp.path()).unwrap();
889        let root = add_op();
890        log.put(&root).unwrap();
891        let left = modify_op(&root.op_id, "fac::Int->Int", "abc123", "left1");
892        log.put(&left).unwrap();
893        let right = modify_op(&root.op_id, "fac::Int->Int", "abc123", "right1");
894        log.put(&right).unwrap();
895
896        let lca = log.lca(&left.op_id, &right.op_id).unwrap();
897        assert_eq!(lca, Some(root.op_id));
898    }
899
900    #[test]
901    fn lca_none_for_independent_histories() {
902        let tmp = tempfile::tempdir().unwrap();
903        let log = OpLog::open(tmp.path()).unwrap();
904        let a = add_op();
905        log.put(&a).unwrap();
906        // A second parentless op (different sig, so different op_id).
907        let b = OperationRecord::new(
908            Operation::new(
909                OperationKind::AddFunction {
910                    sig_id: "double::Int->Int".into(),
911                    stage_id: "ddd111".into(),
912                    effects: BTreeSet::new(),
913                    budget_cost: None,
914                    in_file: None,
915                },
916                [],
917            ),
918            StageTransition::Create {
919                sig_id: "double::Int->Int".into(),
920                stage_id: "ddd111".into(),
921            },
922        );
923        log.put(&b).unwrap();
924
925        assert_eq!(log.lca(&a.op_id, &b.op_id).unwrap(), None);
926    }
927
928    #[test]
929    fn ops_since_excludes_base_history() {
930        let tmp = tempfile::tempdir().unwrap();
931        let log = OpLog::open(tmp.path()).unwrap();
932        let a = add_op();
933        log.put(&a).unwrap();
934        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
935        log.put(&b).unwrap();
936        let c = modify_op(&b.op_id, "fac::Int->Int", "def456", "789aaa");
937        log.put(&c).unwrap();
938
939        let since: Vec<_> = log
940            .ops_since(&c.op_id, Some(&a.op_id))
941            .unwrap()
942            .into_iter()
943            .map(|r| r.op_id)
944            .collect();
945        assert_eq!(since.len(), 2);
946        assert!(since.contains(&b.op_id));
947        assert!(since.contains(&c.op_id));
948        assert!(!since.contains(&a.op_id));
949    }
950
951    #[test]
952    fn repack_consolidates_loose_files_into_a_pack() {
953        let tmp = tempfile::tempdir().unwrap();
954        let log = OpLog::open(tmp.path()).unwrap();
955        let a = add_op();
956        log.put(&a).unwrap();
957        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
958        log.put(&b).unwrap();
959
960        let n = log.repack(0).unwrap();  // threshold 0 = always
961        assert_eq!(n, 2);
962        let ops_dir = tmp.path().join("ops");
963        let loose: Vec<_> = fs::read_dir(&ops_dir).unwrap()
964            .filter_map(|e| e.ok())
965            .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
966            .filter(|e| !e.file_name().to_string_lossy().starts_with("pack-"))
967            .collect();
968        assert!(loose.is_empty(), "loose .json files should be deleted");
969        let packs: Vec<_> = fs::read_dir(&ops_dir).unwrap()
970            .filter_map(|e| e.ok())
971            .filter(|e| e.path().extension().is_some_and(|x| x == "pack"))
972            .collect();
973        assert_eq!(packs.len(), 1);
974
975        // After repack, get() must still return both ops via the pack.
976        assert_eq!(log.get(&a.op_id).unwrap().unwrap(), a);
977        assert_eq!(log.get(&b.op_id).unwrap().unwrap(), b);
978    }
979
980    #[test]
981    fn repack_below_threshold_is_a_noop() {
982        let tmp = tempfile::tempdir().unwrap();
983        let log = OpLog::open(tmp.path()).unwrap();
984        log.put(&add_op()).unwrap();
985        let n = log.repack(10).unwrap();
986        assert_eq!(n, 0);
987    }
988
989    #[test]
990    fn repack_is_deterministic_on_same_input() {
991        // Two stores with the same loose ops repack to the same
992        // pack hash — content addressing all the way down.
993        let make_log = || {
994            let tmp = tempfile::tempdir().unwrap();
995            let log = OpLog::open(tmp.path()).unwrap();
996            let a = add_op();
997            log.put(&a).unwrap();
998            let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "def456");
999            log.put(&b).unwrap();
1000            log.repack(0).unwrap();
1001            (tmp, log)
1002        };
1003        let (tmp1, _log1) = make_log();
1004        let (tmp2, _log2) = make_log();
1005        let pack_name = |dir: &std::path::Path| -> String {
1006            fs::read_dir(dir.join("ops")).unwrap()
1007                .filter_map(|e| e.ok())
1008                .find(|e| e.path().extension().is_some_and(|x| x == "pack"))
1009                .unwrap()
1010                .file_name().into_string().unwrap()
1011        };
1012        assert_eq!(pack_name(tmp1.path()), pack_name(tmp2.path()));
1013    }
1014
1015    #[test]
1016    fn walk_back_works_across_loose_and_packed_ops() {
1017        // Pack the older history, leave newer ops loose. walk_back
1018        // must traverse seamlessly.
1019        let tmp = tempfile::tempdir().unwrap();
1020        let log = OpLog::open(tmp.path()).unwrap();
1021        let a = add_op();
1022        log.put(&a).unwrap();
1023        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "b1");
1024        log.put(&b).unwrap();
1025        log.repack(0).unwrap();
1026        // Now add a newer op as a loose file.
1027        let c = modify_op(&b.op_id, "fac::Int->Int", "b1", "c1");
1028        log.put(&c).unwrap();
1029
1030        let walked = log.walk_back(&c.op_id, None).unwrap();
1031        let ids: Vec<_> = walked.iter().map(|r| r.op_id.as_str()).collect();
1032        assert_eq!(ids, vec![c.op_id.as_str(), b.op_id.as_str(), a.op_id.as_str()]);
1033    }
1034
1035    #[test]
1036    fn list_all_dedups_across_loose_and_pack() {
1037        let tmp = tempfile::tempdir().unwrap();
1038        let log = OpLog::open(tmp.path()).unwrap();
1039        let a = add_op();
1040        log.put(&a).unwrap();
1041        log.repack(0).unwrap();
1042        // Re-put the same op as a loose file (simulate an
1043        // interrupted repack). list_all should still report
1044        // exactly one record per op_id.
1045        log.put(&a).unwrap();
1046
1047        let all = log.list_all().unwrap();
1048        assert_eq!(all.len(), 1);
1049        assert_eq!(all[0].op_id, a.op_id);
1050    }
1051
1052    #[test]
1053    fn walk_back_orders_ancestors_after_descendants() {
1054        // Build a small DAG with a merge:
1055        //
1056        //     a
1057        //    / \
1058        //   b   c
1059        //    \ /
1060        //     m  (merge with parents [b, c])
1061        //
1062        // The merge engine relies on the property that any ancestor of
1063        // X appears strictly after X in the walk_back output. Pin it.
1064        let tmp = tempfile::tempdir().unwrap();
1065        let log = OpLog::open(tmp.path()).unwrap();
1066        let a = add_op();
1067        log.put(&a).unwrap();
1068        let b = modify_op(&a.op_id, "fac::Int->Int", "abc123", "b1");
1069        log.put(&b).unwrap();
1070        let c = OperationRecord::new(
1071            Operation::new(
1072                OperationKind::ModifyBody {
1073                    sig_id: "double::Int->Int".into(),
1074                    from_stage_id: "ddd000".into(),
1075                    to_stage_id: "c1".into(),
1076                    from_budget: None,
1077                    to_budget: None,
1078                    to_sig_id: None,
1079                },
1080                [a.op_id.clone()],
1081            ),
1082            StageTransition::Replace {
1083                sig_id: "double::Int->Int".into(),
1084                from: "ddd000".into(),
1085                to: "c1".into(),
1086            },
1087        );
1088        log.put(&c).unwrap();
1089        let m = OperationRecord::new(
1090            Operation::new(
1091                OperationKind::Merge { resolved: 0 },
1092                [b.op_id.clone(), c.op_id.clone()],
1093            ),
1094            StageTransition::Merge { entries: BTreeMap::new() },
1095        );
1096        log.put(&m).unwrap();
1097
1098        let walked = log.walk_back(&m.op_id, None).unwrap();
1099        let pos = |id: &str| walked.iter().position(|r| r.op_id == id).unwrap();
1100        let (m_pos, b_pos, c_pos, a_pos) =
1101            (pos(&m.op_id), pos(&b.op_id), pos(&c.op_id), pos(&a.op_id));
1102        // Each ancestor must appear strictly after its descendants.
1103        assert!(m_pos < b_pos, "merge before its parent b");
1104        assert!(m_pos < c_pos, "merge before its parent c");
1105        assert!(b_pos < a_pos, "b before its parent a");
1106        assert!(c_pos < a_pos, "c before its parent a");
1107    }
1108}