Skip to main content

refs/refs/
refs_manager.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Reference manager: threads, markers, HEAD, and packed refs.
3
4use std::{
5    path::{Path, PathBuf},
6    sync::{
7        Arc, Mutex,
8        atomic::{AtomicU64, Ordering},
9    },
10    time::SystemTime,
11};
12
13use heddle_object_model::op_record::OpRecord;
14use objects::{
15    error::{HeddleError, Result},
16    object::{MarkerName, StateId, ThreadName},
17};
18
19use super::{
20    Head, RefExpectation, RefUpdate,
21    backend::CoreRefBackend,
22    format_state_id_text,
23    packed_refs::PackedRefs,
24    reconcile::{LoadRequest, Loaded, RefClass, RefCommitter, RefReconciler},
25    ref_backend::RefBackend,
26    refs_storage::RefsLock,
27    resolve_refspec,
28};
29use crate::fs_atomic::{create_dir_all_durable, sync_directory};
30
31/// Sentinel meaning "no batch has been reconciled yet" — distinct from a real
32/// `head_id` so the first read after a reconciler is injected always reconciles.
33const WATERMARK_UNSET: u64 = u64::MAX;
34
35fn watermark_covers(cached: u64, tip: u64) -> bool {
36    cached != WATERMARK_UNSET && cached >= tip
37}
38
39/// Per-worktree persisted **local**-class watermark (HEAD + undo-recovery),
40/// stored beside the per-checkout `HEAD`. Local refs are worktree-private, so
41/// each checkout tracks its own.
42const RECONCILE_WATERMARK_LOCAL: &str = "RECONCILE_WATERMARK_LOCAL";
43
44/// Persisted **shared**-class watermark (thread / marker / remote-thread),
45/// stored in the SHARED Heddle dir so every sibling worktree advances and seeds
46/// from the SAME value (heddle#354 r6, cid 3329711893). A per-worktree shared
47/// watermark let a checkout opened with a lagging file re-fold a shared create a
48/// sibling had already processed/published — resurrecting it cross-worktree.
49const RECONCILE_WATERMARK_SHARED: &str = "RECONCILE_WATERMARK_SHARED";
50
51/// Reconstructible, request-scoped proof that the canonical HEAD written by a
52/// committed snapshot already matches the current oplog tip. Unlike a class
53/// watermark this is trusted only when both its tip and raw ref value match.
54const SNAPSHOT_WITNESS_LOCAL: &str = "SNAPSHOT_WITNESS_LOCAL";
55
56/// Shared-thread counterpart to [`SNAPSHOT_WITNESS_LOCAL`].
57const SNAPSHOT_WITNESS_SHARED: &str = "SNAPSHOT_WITNESS_SHARED";
58
59/// Well-known refspec that resolves the heddle-internal pre-undo recovery
60/// pointer used by `heddle undo --recover`. It is UNSHADOWABLE by any user
61/// marker or thread, in BOTH directions (heddle#305 r3):
62///
63/// - **Write side:** the leading `.` is rejected by [`validate_ref_name`], so a
64///   user can never `marker create` / `thread` a ref with this name. The
65///   recovery state therefore lives in a reserved namespace no user ref can
66///   occupy.
67/// - **Resolve side:** [`resolve_refspec`] routes this handle to the internal
68///   recovery pointer BEFORE consulting user threads/markers, so no user ref
69///   can intercept the advertised handle.
70///
71/// Invariant: an advertised handle for an internal ref must use a reserved form
72/// that user-namespace names cannot take — never a bare user-namespace name.
73///
74/// [`validate_ref_name`]: super::name::validate_ref_name
75pub const UNDO_RECOVERY_HANDLE: &str = ".undo-recovery";
76
77/// Process-local packed-refs snapshot keyed by on-disk identity.
78///
79/// Avoids re-`read_to_string` + parse on every cold `get_thread` /
80/// `get_marker` when the file has not changed. Invalidated on write and
81/// revalidated via `(mtime, len)` when another process rewrites the file.
82struct CachedPackedRefs {
83    stamp: Option<(SystemTime, u64)>,
84    packed: Arc<PackedRefs>,
85}
86
87/// Manager for references (threads, markers, HEAD).
88pub struct RefManager {
89    pub(crate) root: PathBuf,
90    pub(crate) local_head: Option<PathBuf>,
91    /// Oplog-backed reconciler (heddle#330 read chokepoint). `None` for the
92    /// bootstrap/test path — then `reconciled_load` returns the plain cache,
93    /// behaviourally identical to the pre-chokepoint code.
94    reconciler: Option<Arc<dyn RefReconciler>>,
95    /// Oplog-backed committer (heddle#330 write chokepoint). `None` for the
96    /// bootstrap/test path — then `commit_and_publish` publishes without a
97    /// record, like the pre-chokepoint code.
98    committer: Option<Arc<dyn RefCommitter>>,
99    /// Watermark of fully-materialized **local**-class batches (HEAD,
100    /// undo-recovery) — `op_scope`-scoped. `WATERMARK_UNSET` until first reconcile.
101    cached_local_generation: AtomicU64,
102    /// Watermark of fully-materialized **shared**-class batches (thread, marker,
103    /// remote-thread) — global across lanes.
104    cached_shared_generation: AtomicU64,
105    /// In-process packed-refs cache (see [`CachedPackedRefs`]).
106    packed_refs_cache: Mutex<Option<CachedPackedRefs>>,
107}
108
109impl RefManager {
110    pub fn new(heddle_dir: impl AsRef<Path>) -> Self {
111        Self {
112            root: heddle_dir.as_ref().to_path_buf(),
113            local_head: None,
114            reconciler: None,
115            committer: None,
116            cached_local_generation: AtomicU64::new(WATERMARK_UNSET),
117            cached_shared_generation: AtomicU64::new(WATERMARK_UNSET),
118            packed_refs_cache: Mutex::new(None),
119        }
120    }
121
122    pub fn with_local_head(mut self, path: PathBuf) -> Self {
123        self.local_head = Some(path);
124        self
125    }
126
127    /// Inject the oplog-backed reconciler (heddle#330 §2.2). Once set, every
128    /// logical read funnels through [`RefManager::reconciled_load`] and
129    /// reconciles against the committed oplog tail. Mirrors the
130    /// [`with_local_head`](Self::with_local_head) builder shape.
131    ///
132    /// The class watermarks are seeded to the **current** generation: this
133    /// handle trusts the already-published canonical cache as of open and
134    /// reconciles only commits made *after* it — the load-bearing long-held
135    /// handle cell (the daemon's `Arc<Repository>`, cid 3328112197) an
136    /// open-time-only pass cannot reach. (Catching a *pre-open* crash lag is the
137    /// job of the optional `Repository::open` eager pass, deferred here as the
138    /// spike's stated optimization, not the guarantee.) Seeding to the current
139    /// generation also keeps reconciliation from re-deriving long-since-deleted
140    /// refs from old records in the un-migrated tree, where deletes do not all
141    /// record yet.
142    pub fn with_reconciler(mut self, reconciler: Arc<dyn RefReconciler>) -> Self {
143        // Seed both watermarks to the current generation. On a header read error
144        // (cid 3329631081) seed `WATERMARK_UNSET` rather than swallowing it as a
145        // generation — the next [`reconciled_load`] re-reads `generation()` and
146        // propagates the error loudly instead of trusting a fabricated value.
147        let generation = reconciler.generation().unwrap_or(WATERMARK_UNSET);
148        self.cached_local_generation
149            .store(generation, Ordering::Release);
150        self.cached_shared_generation
151            .store(generation, Ordering::Release);
152        self.reconciler = Some(reconciler);
153        self
154    }
155
156    /// Inject the oplog-backed committer (heddle#330 §2.2 write chokepoint).
157    /// Once set, [`commit_and_publish`](Self::commit_and_publish) appends the
158    /// caller's ref-carrying records before publishing the ref batch.
159    pub fn with_committer(mut self, committer: Arc<dyn RefCommitter>) -> Self {
160        self.committer = Some(committer);
161        self
162    }
163
164    /// The atomic-write entry of THE write chokepoint (heddle#330 §2.2): commit
165    /// the caller-supplied ref-carrying record batch (phase 4) **before**
166    /// publishing the atomic ref batch (phase 5), record-before-publish, the
167    /// whole batch published as one unit. The bare publish (temp→rename, via
168    /// `update_refs_with_lock`) is reachable through this seam, never with a
169    /// ref published ahead of its record. With no committer it degrades to a
170    /// plain publish (bootstrap).
171    ///
172    /// **Invariant (cid 3329490978 / 3329490984): the oplog record and the ref
173    /// publish commit together under the refs lock; a record exists iff its
174    /// publish succeeded, and concurrent publishes to the same ref serialize
175    /// record-and-publish as a unit.** Routed through
176    /// [`write_chokepoint`](Self::write_chokepoint), which takes the refs lock
177    /// FIRST and materializes the committed-but-unpublished tail of every class
178    /// BEFORE the body runs; the ref expectations are then validated (phase 3)
179    /// against that reconciled state, BEFORE the record is appended (phase 4),
180    /// and the publish (phase 5) follows under the same lock — so a failed
181    /// expectation never leaks a record, and two concurrent callers can never
182    /// append in one order and publish in another. (For `PgRefBackend` the single
183    /// `pool.begin()…commit()` gives the same atomicity natively.)
184    pub fn commit_and_publish(
185        &self,
186        records: &[OpRecord],
187        ref_updates: &[RefUpdate],
188        scope: Option<&str>,
189    ) -> Result<()> {
190        self.write_chokepoint(|lock| {
191            self.validate_commit_publish(ref_updates, lock, || {
192                // Phase 4 — the commit point: append the ref-carrying records
193                // only after phase-3 validation has passed, under the held lock.
194                let committed_for_reconcile = self.committer.is_some() && !records.is_empty();
195                if let Some(committer) = self.committer.as_ref() {
196                    committer.commit_records(records, scope)?;
197                } else if !records.is_empty() {
198                    // Fail closed (heddle#354 r9, cid 3330304656): no committer
199                    // is wired but records were handed in. Publishing the refs
200                    // here would silently drop them — committed data must never
201                    // be lost. The bootstrap/no-committer path only legitimately
202                    // runs with an empty record batch.
203                    return Err(HeddleError::Config(format!(
204                        "commit_and_publish was handed {} record(s) but this RefManager has no \
205                         committer; refusing to publish and silently drop committed data",
206                        records.len()
207                    )));
208                }
209                Ok(committed_for_reconcile)
210            })
211        })
212    }
213
214    /// THE write chokepoint (heddle#354 r7): the SOLE path by which any ref
215    /// write reaches the backend. Under ONE held publish lock it
216    ///
217    /// 1. reconciles AND materializes the committed-but-unpublished tail of
218    ///    BOTH ref classes ([`materialize_committed_tail`](Self::materialize_committed_tail)),
219    ///    advancing+persisting each class watermark to the current oplog tip, and
220    /// 2. runs the caller's `body` (validate → commit → publish) against the
221    ///    now-materialized canonical state.
222    ///
223    /// Materializing FIRST is what closes the lost-/clobbered-record class (cid
224    /// 3329765073): a non-atomic [`update_refs`](Self::update_refs) — or any
225    /// remote-thread / undo-recovery setter — used to fold the committed tail
226    /// only for *validation* and discard it, so the canonical never caught up
227    /// and a later read could re-fold an ancient record over the just-written
228    /// value. Now every write first persists the committed tail and advances the
229    /// watermark to the tip, so the write lands on a fully-reconciled canonical
230    /// and no later read re-folds across it.
231    ///
232    /// **No-bypass invariant.** The raw backend writers (`publish_ref_plans`,
233    /// `set_remote_thread_locked`, `delete_remote_thread_locked`,
234    /// `set_undo_recovery_locked`) are private and reached ONLY from a
235    /// chokepoint `body` or from [`materialize`](Self::materialize) (which itself
236    /// runs only here and on the read chokepoint). A source-level conformance
237    /// check (`write_read_conformance` in the refs tests) fails CI if any other
238    /// path calls a raw writer, so a path-by-path allowlist cannot silently
239    /// re-open the class.
240    pub(super) fn write_chokepoint<T>(
241        &self,
242        body: impl FnOnce(&RefsLock) -> Result<T>,
243    ) -> Result<T> {
244        let lock = self.lock_refs()?;
245        self.materialize_committed_tail(&lock)?;
246        body(&lock)
247    }
248
249    /// Reconcile + materialize the committed-but-unpublished tail of BOTH ref
250    /// classes under the held lock — the first step of every write
251    /// ([`write_chokepoint`](Self::write_chokepoint)). The O(1) gate
252    /// (watermark == tip) skips a class with no lag, so on the hot path this is
253    /// two `generation()` header reads and no fold.
254    fn materialize_committed_tail(&self, lock: &RefsLock) -> Result<()> {
255        self.materialize_class(RefClass::Local, lock)?;
256        self.materialize_class(RefClass::Shared, lock)?;
257        Ok(())
258    }
259
260    /// Atomically publish a just-committed attached snapshot as a
261    /// reconstructible materialized view. The caller already owns the
262    /// authoritative state and oplog tip, so replaying the same oplog tail
263    /// through a fresh reconciler here only adds file/index I/O to every
264    /// capture. The persisted watermarks deliberately remain at their prior
265    /// durable floor for fresh-process recovery.
266    pub fn materialize_snapshot_thread_after_commit(
267        &self,
268        thread: &ThreadName,
269        state: StateId,
270        tip: u64,
271    ) -> Result<()> {
272        let lock = self.lock_refs()?;
273        let outcome = super::reconcile::ReconcileOutcome {
274            loaded: Loaded::Point(Some(state)),
275            republish: vec![RefUpdate::Thread {
276                name: thread.clone(),
277                expected: RefExpectation::Any,
278                new: Some(state),
279            }],
280            remote_updates: Vec::new(),
281            undo_recovery: None,
282        };
283        self.materialize_with_ref_durability(&outcome, &lock, false)?;
284        self.cached_shared_generation.store(tip, Ordering::Release);
285        // Attached snapshots do not change HEAD's identity; this commit has no
286        // other local-class effect, so the local view is complete at `tip` too.
287        self.cached_local_generation.store(tip, Ordering::Release);
288        let _ = std::fs::write(
289            self.snapshot_witness_local_path(),
290            format!("{tip}\nattached\n{thread}\n"),
291        );
292        let _ = std::fs::write(
293            self.snapshot_witness_shared_path(),
294            format!("{tip}\nthread\n{thread}\n{}\n", state.to_string_full()),
295        );
296        Ok(())
297    }
298
299    /// Detached-HEAD counterpart to
300    /// [`materialize_snapshot_thread_after_commit`](Self::materialize_snapshot_thread_after_commit).
301    pub fn materialize_snapshot_head_after_commit(&self, state: StateId, tip: u64) -> Result<()> {
302        let lock = self.lock_refs()?;
303        let outcome = super::reconcile::ReconcileOutcome {
304            loaded: Loaded::Head(Head::Detached { state }),
305            republish: vec![RefUpdate::Head {
306                expected: RefExpectation::Any,
307                new: Head::Detached { state },
308            }],
309            remote_updates: Vec::new(),
310            undo_recovery: None,
311        };
312        self.materialize_with_ref_durability(&outcome, &lock, false)?;
313        self.cached_local_generation.store(tip, Ordering::Release);
314        // Detached snapshots do not touch a shared ref.
315        self.cached_shared_generation.store(tip, Ordering::Release);
316        let _ = std::fs::write(
317            self.snapshot_witness_local_path(),
318            format!("{tip}\ndetached\n{}\n", state.to_string_full()),
319        );
320        Ok(())
321    }
322
323    /// Materialize one class's committed tail under the held lock — the
324    /// class-scoped form of [`reconciled_load`](Self::reconciled_load)'s lag
325    /// branch, driven by a lightweight probe request of the class (the
326    /// materialization set covers EVERY ref the lagged batches touched and does
327    /// not depend on the specific request — only on the class + fold). No-op
328    /// when the class is not lagging or no reconciler is injected.
329    fn materialize_class(&self, class: RefClass, lock: &RefsLock) -> Result<()> {
330        let Some(reconciler) = self.reconciler.as_ref() else {
331            return Ok(());
332        };
333        let watermark = self.class_watermark(class);
334        let tip = reconciler.generation()?;
335        if watermark_covers(watermark.load(Ordering::Acquire), tip) {
336            return Ok(());
337        }
338        // Re-read the persisted (possibly sibling-advanced) last-clean point
339        // before folding, so a long-lived handle never re-derives a record a
340        // sibling already materialized past (cid 3329765075).
341        self.refresh_persisted_watermark(class, lock)?;
342        let cached = watermark.load(Ordering::Acquire);
343        if watermark_covers(cached, tip) {
344            return Ok(());
345        }
346        let req = Self::class_probe(class);
347        let raw = self.raw_load(&req)?;
348        let since = if cached == WATERMARK_UNSET { 0 } else { cached };
349        let outcome = reconciler.reconcile(&req, raw, since)?;
350        self.materialize(&outcome, lock)?;
351        watermark.store(tip, Ordering::Release);
352        let _ = self.persist_reconcile_watermark(lock);
353        Ok(())
354    }
355
356    /// The per-read class watermark atomic.
357    fn class_watermark(&self, class: RefClass) -> &AtomicU64 {
358        match class {
359            RefClass::Local => &self.cached_local_generation,
360            RefClass::Shared => &self.cached_shared_generation,
361        }
362    }
363
364    /// A lightweight probe request of `class` to drive a class-wide reconcile.
365    /// The materialization set (`republish` / `remote_updates` / `undo_recovery`)
366    /// is class-derived, not request-derived, so any request of the class yields
367    /// the full set; the projected `loaded` value is discarded by the caller.
368    fn class_probe(class: RefClass) -> LoadRequest {
369        match class {
370            RefClass::Local => LoadRequest::Head,
371            RefClass::Shared => LoadRequest::MarkerList,
372        }
373    }
374
375    /// Advance the in-memory class watermark to the persisted last-clean point
376    /// when a sibling worktree (or a prior read) has advanced it past our
377    /// in-memory value (heddle#354 r7, cid 3329765075). The persisted watermark
378    /// is a known-materialized point, so adopting it is always safe and
379    /// ADVANCE-ONLY — it never regresses what this handle already materialized.
380    ///
381    /// This is what makes a long-lived handle behave like a fresh open: rather
382    /// than re-folding from a value frozen at open, each reconcile re-reads the
383    /// CURRENT shared (or local) last-clean point and folds only what genuinely
384    /// lags above it. Called under the held lock so the load→store is serialized
385    /// with every other watermark writer.
386    fn refresh_persisted_watermark(&self, class: RefClass, _lock: &RefsLock) -> Result<()> {
387        let path = match class {
388            RefClass::Local => self.reconcile_watermark_local_path(),
389            RefClass::Shared => self.reconcile_watermark_shared_path(),
390        };
391        let Some(persisted) = self.read_single_watermark(&path)? else {
392            return Ok(());
393        };
394        let watermark = self.class_watermark(class);
395        let cached = watermark.load(Ordering::Acquire);
396        // The `UNSET` sentinel is `u64::MAX`, so a plain `max` would keep it;
397        // adopt the persisted value outright in that case.
398        let next = if cached == WATERMARK_UNSET {
399            persisted
400        } else {
401            cached.max(persisted)
402        };
403        if next != cached {
404            watermark.store(next, Ordering::Release);
405        }
406        Ok(())
407    }
408
409    /// THE read chokepoint (heddle#330 §2.2): the sole path for a **logical
410    /// read** to obtain ref data. The raw loaders
411    /// (`read_state_id_at`/`read_head_state`/`try_read_ref_summary_index`/
412    /// `*_from_storage`/`PackedRefs::load`) are reached from a logical read only
413    /// from inside here — the maintenance path `pack_refs` is the one allowlisted
414    /// non-logical caller. With no reconciler this is the plain raw load.
415    fn reconciled_load(&self, req: LoadRequest) -> Result<Loaded> {
416        heddle_perf_contract::record_ref_read();
417        let Some(reconciler) = self.reconciler.as_ref() else {
418            return self.raw_load(&req);
419        };
420
421        let watermark = match req.ref_class() {
422            RefClass::Local => &self.cached_local_generation,
423            RefClass::Shared => &self.cached_shared_generation,
424        };
425
426        // Cheap O(1) gate: when this class's watermark already equals the oplog
427        // tip, every committed record of the class is materialized into
428        // canonical ⇒ the raw read is authoritative ⇒ no lock, no tail scan.
429        // A `generation()` error propagates (cid 3329631081) — never silently
430        // treated as generation 0.
431        let tip = reconciler.generation()?;
432        if watermark_covers(watermark.load(Ordering::Acquire), tip) {
433            return self.raw_load(&req);
434        }
435        if let Some(loaded) = self.snapshot_witnessed_load(&req, tip)? {
436            return Ok(loaded);
437        }
438
439        // Lag: the fold AND the lazy re-publish must be atomic w.r.t. a
440        // concurrent `commit_and_publish` (cid 3329631077). Take the publish
441        // lock FIRST, then re-read tip + raw and fold UNDER the lock — so a
442        // concurrent publish that lands a newer value cannot interpose between
443        // the fold and the materialize. The fold sees the newest committed
444        // record (highest id wins), so materialization never republishes a stale
445        // value over a freshly-published newer one.
446        let lock = self.lock_refs()?;
447        let tip = reconciler.generation()?;
448        // Re-read the CURRENT persisted watermark (heddle#354 r7, cid 3329765075):
449        // a long-lived handle's in-memory watermark is frozen at open, so a
450        // sibling worktree that advanced the shared last-clean point past it
451        // would otherwise be re-folded from the stale frozen value. Refreshing
452        // here makes a long-lived handle fold from the same floor a fresh open
453        // would — re-deriving only what genuinely lags above the current point.
454        self.refresh_persisted_watermark(req.ref_class(), &lock)?;
455        let cached = watermark.load(Ordering::Acquire);
456        let raw = self.raw_load(&req)?;
457        if watermark_covers(cached, tip) {
458            // A concurrent reconcile materialized the lag while we waited for the
459            // lock; the freshly-read canonical is now authoritative.
460            return Ok(raw);
461        }
462
463        // The reconcile is batch-atomic — it returns the re-materialization set
464        // for every ref the lagged batches touched, which we publish (under the
465        // held lock) so the watermark can advance without leaving a batch sibling
466        // stale.
467        let since = if cached == WATERMARK_UNSET { 0 } else { cached };
468        let outcome = reconciler.reconcile(&req, raw, since)?;
469        self.materialize(&outcome, &lock)?;
470        watermark.store(tip, Ordering::Release);
471        // Persist the advanced watermark so a future process seeds from this
472        // last-clean point and folds only the genuine crash tail above it, never
473        // re-deriving long-since-deleted refs from ancient records (cid
474        // 3329631074). Best-effort: a write failure only costs extra folding next
475        // open, never correctness.
476        let _ = self.persist_reconcile_watermark(&lock);
477        Ok(outcome.loaded)
478    }
479
480    /// Fold the committed-but-unpublished oplog tail over the raw value WITHOUT
481    /// taking the refs lock — the caller already holds it (phase-3 validation in
482    /// [`plan_ref_updates`](Self::plan_ref_updates)). Closes the stale-validation
483    /// gap (cid 3329631079): a `Missing`/CAS expectation, and the publish base,
484    /// are computed from the reconciled state — never a pre-lock raw read that a
485    /// crash-left committed-but-unpublished record has made stale. The same O(1)
486    /// gate applies: with the watermark current, the raw value is authoritative.
487    pub(super) fn reconciled_value_under_lock(&self, req: &LoadRequest) -> Result<Loaded> {
488        let raw = self.raw_load(req)?;
489        let Some(reconciler) = self.reconciler.as_ref() else {
490            return Ok(raw);
491        };
492        let tip = reconciler.generation()?;
493        let watermark = match req.ref_class() {
494            RefClass::Local => &self.cached_local_generation,
495            RefClass::Shared => &self.cached_shared_generation,
496        };
497        let cached = watermark.load(Ordering::Acquire);
498        if watermark_covers(cached, tip) {
499            return Ok(raw);
500        }
501        let since = if cached == WATERMARK_UNSET { 0 } else { cached };
502        Ok(reconciler.reconcile(req, raw, since)?.loaded)
503    }
504
505    /// Seed the per-read watermarks from the persisted last-clean point
506    /// (heddle#354 r5, cid 3329631074), so a fresh handle recovers a prior
507    /// process's committed-but-unpublished crash tail.
508    ///
509    /// A `RefManager` seeds its in-memory watermarks at the current generation
510    /// ([`with_reconciler`](Self::with_reconciler)) — so the per-read gate, on a
511    /// fresh process, would never fold a record committed *before* this handle
512    /// opened, and a cross-process crash (phase-4 committed, phase-5 publish
513    /// never ran) would be silently lost. The fix is NOT an eager open-time fold
514    /// (that would re-derive long-since-deleted refs from ancient records, since
515    /// the un-migrated delete paths do not all record yet): it is a **persisted
516    /// watermark**. Reads advance and persist it past every materialized record,
517    /// so on open the seed sits at the last point canonical was known-consistent;
518    /// the per-read reconcile then folds only `(seed, tip]` — the genuine crash
519    /// tail — and never the ancient records below the seed.
520    ///
521    /// When no watermark has been persisted yet (a fresh repo, or a repo from
522    /// before this version), seed conservatively at the current generation and
523    /// write the file, so the next process has a real last-clean point.
524    ///
525    /// The two classes seed from SEPARATE files: the local watermark from the
526    /// per-worktree file, the shared watermark from the shared-dir file (cid
527    /// 3329711893). A sibling worktree that already advanced the shared
528    /// watermark publishes it to the shared file, so this checkout seeds at that
529    /// shared last-clean point and never re-folds a shared create the sibling
530    /// already processed.
531    pub fn init_reconcile_watermark(&self) -> Result<()> {
532        if self.reconciler.is_none() {
533            return Ok(());
534        }
535        let (local, shared) = self.read_persisted_reconcile_watermark()?;
536        if let Some(local) = local {
537            self.cached_local_generation.store(local, Ordering::Release);
538        }
539        if let Some(shared) = shared {
540            self.cached_shared_generation
541                .store(shared, Ordering::Release);
542        }
543        // Any class with no persisted last-clean point yet (fresh repo, or a
544        // repo from before this version) keeps the current-generation seed from
545        // `with_reconciler` and gets written, so the next process has a real
546        // last-clean point.
547        if local.is_none() || shared.is_none() {
548            let lock = self.lock_refs()?;
549            self.persist_reconcile_watermark(&lock)?;
550        }
551        Ok(())
552    }
553
554    /// Per-worktree LOCAL watermark file (HEAD + undo-recovery), beside the
555    /// per-checkout `HEAD` and `UNDO_RECOVERY` — local refs are worktree-private.
556    fn reconcile_watermark_local_path(&self) -> PathBuf {
557        self.head_path()
558            .parent()
559            .map(|dir| dir.join(RECONCILE_WATERMARK_LOCAL))
560            .unwrap_or_else(|| self.root.join(RECONCILE_WATERMARK_LOCAL))
561    }
562
563    /// SHARED watermark file (thread / marker / remote-thread), in the SHARED
564    /// Heddle dir (`self.root`, objectstore-pointed). Every sibling worktree
565    /// resolves the SAME path, so a shared create one worktree advances past is
566    /// never re-folded by another (cid 3329711893). Mirrors `refs/`, which lives
567    /// under the same shared root and whose `LOCK` already serializes writers.
568    fn reconcile_watermark_shared_path(&self) -> PathBuf {
569        self.root.join(RECONCILE_WATERMARK_SHARED)
570    }
571
572    fn snapshot_witness_local_path(&self) -> PathBuf {
573        self.head_path()
574            .parent()
575            .map(|dir| dir.join(SNAPSHOT_WITNESS_LOCAL))
576            .unwrap_or_else(|| self.root.join(SNAPSHOT_WITNESS_LOCAL))
577    }
578
579    fn snapshot_witness_shared_path(&self) -> PathBuf {
580        self.root.join(SNAPSHOT_WITNESS_SHARED)
581    }
582
583    /// Return the raw point value when a snapshot witness proves this exact
584    /// read is already materialized at `tip`. Missing, torn, stale, and
585    /// mismatched witnesses are cache misses and fall through to reconciliation.
586    fn snapshot_witnessed_load(&self, request: &LoadRequest, tip: u64) -> Result<Option<Loaded>> {
587        let parse_tip = |lines: &mut std::str::Lines<'_>| {
588            lines
589                .next()
590                .and_then(|value| value.parse::<u64>().ok())
591                .filter(|value| *value == tip)
592        };
593        match request {
594            LoadRequest::Head => {
595                let Some(contents) =
596                    self.read_optional_string(&self.snapshot_witness_local_path())?
597                else {
598                    return Ok(None);
599                };
600                let mut lines = contents.lines();
601                if parse_tip(&mut lines).is_none() {
602                    return Ok(None);
603                }
604                let witnessed = match lines.next() {
605                    Some("attached") => lines.next().map(|thread| Head::Attached {
606                        thread: ThreadName::new(thread),
607                    }),
608                    Some("detached") => lines
609                        .next()
610                        .and_then(|state| StateId::parse(state).ok())
611                        .map(|state| Head::Detached { state }),
612                    _ => None,
613                };
614                let Some(witnessed) = witnessed else {
615                    return Ok(None);
616                };
617                let raw = self.read_head_state()?.head;
618                Ok((raw == witnessed).then_some(Loaded::Head(raw)))
619            }
620            LoadRequest::Thread(requested) => {
621                let Some(contents) =
622                    self.read_optional_string(&self.snapshot_witness_shared_path())?
623                else {
624                    return Ok(None);
625                };
626                let mut lines = contents.lines();
627                if parse_tip(&mut lines).is_none() || lines.next() != Some("thread") {
628                    return Ok(None);
629                }
630                let Some(thread) = lines.next() else {
631                    return Ok(None);
632                };
633                let Some(state) = lines.next().and_then(|value| StateId::parse(value).ok()) else {
634                    return Ok(None);
635                };
636                if thread != requested.as_str() {
637                    return Ok(None);
638                }
639                let raw = self.raw_get_thread(requested)?;
640                Ok((raw == Some(state)).then_some(Loaded::Point(raw)))
641            }
642            _ => Ok(None),
643        }
644    }
645
646    /// Read the persisted `(local, shared)` watermark from their two scope
647    /// files; each component is `None` when absent / unparseable ("no last-clean
648    /// point yet" for that class).
649    fn read_persisted_reconcile_watermark(&self) -> Result<(Option<u64>, Option<u64>)> {
650        let local = self.read_single_watermark(&self.reconcile_watermark_local_path())?;
651        let shared = self.read_single_watermark(&self.reconcile_watermark_shared_path())?;
652        Ok((local, shared))
653    }
654
655    /// Read a single `u64` watermark from `path`, or `None` when absent /
656    /// unparseable.
657    fn read_single_watermark(&self, path: &Path) -> Result<Option<u64>> {
658        let Some(contents) = self.read_optional_string(path)? else {
659            return Ok(None);
660        };
661        Ok(contents
662            .split_whitespace()
663            .next()
664            .and_then(|s| s.parse::<u64>().ok()))
665    }
666
667    /// Persist the current in-memory local + shared watermarks, each to its own
668    /// scope file, under the held refs lock. Called after a read advances a
669    /// watermark, and at open when a file does not exist yet.
670    fn persist_reconcile_watermark(&self, _lock: &RefsLock) -> Result<()> {
671        let local = self.cached_local_generation.load(Ordering::Acquire);
672        let shared = self.cached_shared_generation.load(Ordering::Acquire);
673        self.persist_watermark_file(&self.reconcile_watermark_local_path(), local)?;
674        self.persist_watermark_file(&self.reconcile_watermark_shared_path(), shared)?;
675        Ok(())
676    }
677
678    /// Write `value` to a watermark file ADVANCE-ONLY: never below what is
679    /// already on disk. The shared file is written by concurrent sibling
680    /// worktrees (serialized by the shared refs `LOCK`); a checkout whose
681    /// in-memory shared watermark lags a sibling's published value must not
682    /// regress the file when it persists after a local-only read (cid
683    /// 3329711893). A still-`UNSET` watermark (no reconciler, or seeded UNSET on
684    /// a header error) is not a meaningful last-clean point — skip it.
685    fn persist_watermark_file(&self, path: &Path, value: u64) -> Result<()> {
686        if value == WATERMARK_UNSET {
687            return Ok(());
688        }
689        let on_disk = self.read_single_watermark(path)?.unwrap_or(0);
690        let next = value.max(on_disk);
691        self.write_string(path, &format!("{next}\n"))
692    }
693
694    /// Lazily re-publish (phase-5 materialization) the refs a reconcile found
695    /// committed-but-unpublished — the records already exist, so this writes
696    /// canonical only (never the oplog).
697    ///
698    /// **Authoritative-apply (cid 3329490981):** a committed record past the
699    /// class watermark is authoritative over the live canonical, so a folded
700    /// value is materialized when it CREATES a missing ref *or* UPDATES a stale
701    /// present one (the crash-replayed update-to-existing case) — not
702    /// fill-if-absent, which silently dropped a committed update to an
703    /// already-existing ref. The folded set only ever holds refs touched by
704    /// commits newer than the watermark, so applying it respects the
705    /// two-watermark scoping and a ref with no recent committed record is never
706    /// rewritten; a write equal to the canonical is skipped as a no-op. (The
707    /// rare un-migrated case where an unrecorded direct write raced in *after*
708    /// the commit is the residual the writers' record-first migration closes.)
709    fn materialize(
710        &self,
711        outcome: &super::reconcile::ReconcileOutcome,
712        lock: &RefsLock,
713    ) -> Result<()> {
714        self.materialize_with_ref_durability(outcome, lock, true)
715    }
716
717    fn materialize_with_ref_durability(
718        &self,
719        outcome: &super::reconcile::ReconcileOutcome,
720        lock: &RefsLock,
721        durable_refs: bool,
722    ) -> Result<()> {
723        // The whole materialization runs under the caller's single held lock so
724        // the fold that produced `outcome` and these re-publishes are one atomic
725        // unit vs a concurrent publish (cid 3329631077). The publish values are
726        // the authoritative folded values; the no-op skip is against the current
727        // canonical, computed inside `plan_materialization`.
728        let plans = self.plan_materialization(&outcome.republish)?;
729        if !plans.is_empty() {
730            if durable_refs {
731                self.publish_ref_plans(plans, lock)?;
732            } else {
733                self.publish_ref_plans_reconstructible(plans, lock)?;
734            }
735        }
736        for (remote, thread, value) in &outcome.remote_updates {
737            if self.raw_get_remote_thread(remote, thread)? != *value {
738                match value {
739                    Some(state) => self.set_remote_thread_locked(remote, thread, state, lock)?,
740                    None => {
741                        self.delete_remote_thread_locked(remote, thread, lock)?;
742                    }
743                }
744            }
745        }
746        if let Some(state) = &outcome.undo_recovery {
747            let current = self.read_state_id_at(
748                &self.undo_recovery_path(),
749                "undo recovery",
750                UNDO_RECOVERY_HANDLE,
751            )?;
752            if current.as_ref() != Some(state) {
753                self.set_undo_recovery_locked(state, lock)?;
754            }
755        }
756        Ok(())
757    }
758
759    /// Request-scoped raw read — the private sub-step `reconciled_load` calls.
760    /// Each arm touches exactly one raw loader for a point read (no whole-set
761    /// scan on the hot path).
762    fn raw_load(&self, req: &LoadRequest) -> Result<Loaded> {
763        Ok(match req {
764            LoadRequest::Head => Loaded::Head(self.read_head_state()?.head),
765            LoadRequest::Thread(name) => Loaded::Point(self.raw_get_thread(name)?),
766            LoadRequest::Marker(name) => Loaded::Point(self.raw_get_marker(name)?),
767            LoadRequest::UndoRecovery => Loaded::Point(self.read_state_id_at(
768                &self.undo_recovery_path(),
769                "undo recovery",
770                UNDO_RECOVERY_HANDLE,
771            )?),
772            LoadRequest::RemoteThread { remote, thread } => {
773                Loaded::Point(self.raw_get_remote_thread(remote, thread)?)
774            }
775            LoadRequest::ThreadList => Loaded::ThreadList(self.raw_list_threads()?),
776            LoadRequest::MarkerList => Loaded::MarkerList(self.raw_list_markers()?),
777            LoadRequest::RemoteList => Loaded::RemoteList(self.raw_list_remotes()?),
778            LoadRequest::RemoteThreadList { remote } => {
779                Loaded::RemoteThreadList(self.raw_list_remote_threads(remote)?)
780            }
781        })
782    }
783
784    fn raw_get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
785        let path = self.thread_path(name)?;
786        if let Some(id) = self.read_state_id_at(&path, "thread", name)? {
787            return Ok(Some(id));
788        }
789        Ok(self.load_packed_refs_cached()?.get_thread(name))
790    }
791
792    fn raw_get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
793        let path = self.marker_path(name)?;
794        if let Some(id) = self.read_state_id_at(&path, "marker", name)? {
795            return Ok(Some(id));
796        }
797        Ok(self.load_packed_refs_cached()?.get_marker(name))
798    }
799
800    /// On-disk identity for the packed-refs file: `(mtime, len)`, or `None`
801    /// when the file is absent. Used to detect external rewrites without
802    /// re-reading the body on every lookup.
803    fn packed_refs_stamp(path: &Path) -> Option<(SystemTime, u64)> {
804        let meta = std::fs::metadata(path).ok()?;
805        let modified = meta.modified().ok()?;
806        Some((modified, meta.len()))
807    }
808
809    /// Load packed-refs with a process-local cache. Safe under concurrent
810    /// readers in this process; writers call [`invalidate_packed_refs_cache`]
811    /// after mutating the file.
812    pub(super) fn load_packed_refs_cached(&self) -> Result<Arc<PackedRefs>> {
813        let path = self.packed_refs_path();
814        let stamp = Self::packed_refs_stamp(&path);
815        let mut guard = self.packed_refs_cache.lock().map_err(|_| {
816            HeddleError::Config("Failed to acquire packed-refs cache lock".to_string())
817        })?;
818        if let Some(cached) = guard.as_ref()
819            && cached.stamp == stamp
820        {
821            return Ok(cached.packed.clone());
822        }
823        let packed = Arc::new(PackedRefs::load(&path)?);
824        *guard = Some(CachedPackedRefs {
825            stamp,
826            packed: packed.clone(),
827        });
828        Ok(packed)
829    }
830
831    /// Drop the process-local packed-refs cache after a write so the next
832    /// read reloads from disk.
833    pub(super) fn invalidate_packed_refs_cache(&self) {
834        if let Ok(mut guard) = self.packed_refs_cache.lock() {
835            *guard = None;
836        }
837    }
838
839    pub(super) fn raw_get_remote_thread(
840        &self,
841        remote: &str,
842        thread: &ThreadName,
843    ) -> Result<Option<StateId>> {
844        let path = self.remote_thread_path(remote, thread)?;
845        self.read_state_id_at(&path, "remote thread", &format!("{}/{}", remote, thread))
846    }
847
848    fn raw_list_threads(&self) -> Result<Vec<ThreadName>> {
849        if let Some(summary) = self.try_read_ref_summary_index() {
850            return Ok(summary.thread_names());
851        }
852        self.list_threads_from_storage()
853    }
854
855    fn raw_list_markers(&self) -> Result<Vec<MarkerName>> {
856        if let Some(summary) = self.try_read_ref_summary_index() {
857            return Ok(summary.marker_names());
858        }
859        self.list_markers_from_storage()
860    }
861
862    fn raw_list_remotes(&self) -> Result<Vec<String>> {
863        if let Some(summary) = self.try_read_ref_summary_index() {
864            return Ok(summary.remote_names());
865        }
866        self.list_remotes_from_storage()
867    }
868
869    fn raw_list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
870        if let Some(summary) = self.try_read_ref_summary_index() {
871            return Ok(summary.remote_thread_names(remote));
872        }
873        self.list_remote_threads_from_storage(remote)
874    }
875
876    pub fn init(&self) -> Result<()> {
877        create_dir_all_durable(&self.threads_dir())?;
878        create_dir_all_durable(&self.markers_dir())?;
879        create_dir_all_durable(&self.remotes_dir())?;
880        Ok(())
881    }
882
883    pub fn cleanup_stale_temps(&self) {
884        let refs_dir = self.refs_dir();
885        if let Ok(entries) = std::fs::read_dir(&refs_dir) {
886            for entry in entries.flatten() {
887                let path = entry.path();
888                if path
889                    .extension()
890                    .and_then(|e| e.to_str())
891                    .map(|e| e.starts_with("tmp-"))
892                    .unwrap_or(false)
893                {
894                    let _ = std::fs::remove_file(&path);
895                }
896            }
897        }
898    }
899
900    pub fn read_head(&self) -> Result<Head> {
901        match self.reconciled_load(LoadRequest::Head)? {
902            Loaded::Head(head) => Ok(head),
903            _ => unreachable!("Head request yields Head"),
904        }
905    }
906
907    pub fn write_head(&self, head: &Head) -> Result<()> {
908        self.write_head_cas(RefExpectation::Any, head)
909    }
910
911    pub fn write_head_cas(&self, expected: RefExpectation<Head>, head: &Head) -> Result<()> {
912        self.update_refs(&[RefUpdate::Head {
913            expected,
914            new: head.clone(),
915        }])
916    }
917
918    /// Resolve a point-valued ref request (thread / marker / remote-thread /
919    /// undo-recovery) through reconciliation. All four share the same
920    /// `Loaded::Point` shape; the catch-all is unreachable because the load
921    /// request and the returned variant are paired by construction.
922    fn reconciled_point(&self, request: LoadRequest) -> Result<Option<StateId>> {
923        match self.reconciled_load(request)? {
924            Loaded::Point(id) => Ok(id),
925            _ => unreachable!("point request yields Point"),
926        }
927    }
928
929    pub fn get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
930        self.reconciled_point(LoadRequest::Thread(name.clone()))
931    }
932
933    pub fn set_thread(&self, name: &ThreadName, state: &StateId) -> Result<()> {
934        self.set_thread_cas(name, RefExpectation::Any, state)
935    }
936
937    pub fn set_thread_cas(
938        &self,
939        name: &ThreadName,
940        expected: RefExpectation<StateId>,
941        state: &StateId,
942    ) -> Result<()> {
943        self.update_refs(&[RefUpdate::Thread {
944            name: name.clone(),
945            expected,
946            new: Some(*state),
947        }])
948    }
949
950    pub fn delete_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
951        let state = self.get_thread(name)?;
952        if state.is_some() {
953            self.update_refs(&[RefUpdate::Thread {
954                name: name.clone(),
955                expected: RefExpectation::Any,
956                new: None,
957            }])?;
958        }
959        Ok(state)
960    }
961
962    pub fn delete_thread_cas(
963        &self,
964        name: &ThreadName,
965        expected: RefExpectation<StateId>,
966    ) -> Result<()> {
967        self.update_refs(&[RefUpdate::Thread {
968            name: name.clone(),
969            expected,
970            new: None,
971        }])
972    }
973
974    pub fn list_threads(&self) -> Result<Vec<ThreadName>> {
975        match self.reconciled_load(LoadRequest::ThreadList)? {
976            Loaded::ThreadList(names) => Ok(names),
977            _ => unreachable!("ThreadList request yields ThreadList"),
978        }
979    }
980
981    /// List thread names and targets through one reconciled list read. The
982    /// summary index is updated by reconciliation before it is read here, so
983    /// callers avoid one logical point read per listed thread.
984    pub fn list_threads_with_states(&self) -> Result<Vec<(ThreadName, StateId)>> {
985        let names = self.list_threads()?;
986        if let Some(summary) = self.try_read_ref_summary_index() {
987            let states = summary.thread_states();
988            if states.len() == names.len()
989                && states
990                    .iter()
991                    .zip(&names)
992                    .all(|((state_name, _), name)| state_name == name)
993            {
994                return Ok(states);
995            }
996        }
997        names
998            .into_iter()
999            .filter_map(|name| match self.raw_get_thread(&name) {
1000                Ok(Some(state)) => Some(Ok((name, state))),
1001                Ok(None) => None,
1002                Err(error) => Some(Err(error)),
1003            })
1004            .collect()
1005    }
1006
1007    pub fn get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1008        self.reconciled_point(LoadRequest::Marker(name.clone()))
1009    }
1010
1011    pub fn create_marker(&self, name: &MarkerName, state: &StateId) -> Result<()> {
1012        self.set_marker_cas(name, RefExpectation::Missing, state)
1013    }
1014
1015    pub fn set_marker_cas(
1016        &self,
1017        name: &MarkerName,
1018        expected: RefExpectation<StateId>,
1019        state: &StateId,
1020    ) -> Result<()> {
1021        self.update_refs(&[RefUpdate::Marker {
1022            name: name.clone(),
1023            expected,
1024            new: Some(*state),
1025        }])
1026    }
1027
1028    pub fn delete_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1029        let state = self.get_marker(name)?;
1030        if state.is_some() {
1031            self.delete_marker_cas(name, RefExpectation::Any)?;
1032        }
1033        Ok(state)
1034    }
1035
1036    pub fn delete_marker_cas(
1037        &self,
1038        name: &MarkerName,
1039        expected: RefExpectation<StateId>,
1040    ) -> Result<()> {
1041        self.update_refs(&[RefUpdate::Marker {
1042            name: name.clone(),
1043            expected,
1044            new: None,
1045        }])
1046    }
1047
1048    pub fn list_markers(&self) -> Result<Vec<MarkerName>> {
1049        match self.reconciled_load(LoadRequest::MarkerList)? {
1050            Loaded::MarkerList(names) => Ok(names),
1051            _ => unreachable!("MarkerList request yields MarkerList"),
1052        }
1053    }
1054
1055    /// Record the heddle-internal pre-undo recovery pointer (ORIG_HEAD-style:
1056    /// a single rolling ref each undo overwrites). Stored OUTSIDE the
1057    /// user-writable marker namespace so `marker create/delete` — and their
1058    /// undo inverses — can never collide with it. See
1059    /// [`UNDO_RECOVERY_HANDLE`] for the resolution handle.
1060    pub fn set_undo_recovery(&self, state: &StateId) -> Result<()> {
1061        self.set_undo_recovery_raw(state)
1062    }
1063
1064    /// The undo-recovery write entry of the write chokepoint: materialize the
1065    /// committed tail FIRST, then write canonical under the held lock. Has no
1066    /// oplog append of its own (undo-recovery is recorded via the atomic
1067    /// `commit_and_publish` path); routing it through
1068    /// [`write_chokepoint`](Self::write_chokepoint) keeps it from bypassing
1069    /// reconciliation (heddle#354 r7).
1070    fn set_undo_recovery_raw(&self, state: &StateId) -> Result<()> {
1071        self.write_chokepoint(|lock| self.set_undo_recovery_locked(state, lock))
1072    }
1073
1074    /// The lock-free core of [`set_undo_recovery_raw`](Self::set_undo_recovery_raw):
1075    /// the caller already holds the refs lock (e.g. the reconciler's
1076    /// materialization runs the whole fold + re-publish under one lock).
1077    fn set_undo_recovery_locked(&self, state: &StateId, _lock: &RefsLock) -> Result<()> {
1078        self.write_string(
1079            &self.undo_recovery_path(),
1080            &super::format_state_id_text(state),
1081        )
1082    }
1083
1084    /// Remove the heddle-internal pre-undo recovery pointer, returning the repo
1085    /// to the "no undo has run" state. Routes through the same
1086    /// [`write_chokepoint`](Self::write_chokepoint) as the setter so it cannot
1087    /// bypass reconciliation, and is a no-op when no pointer exists. Used as the
1088    /// inverse of [`set_undo_recovery`](Self::set_undo_recovery) when the atomic
1089    /// `undo` transaction rewinds and the pointer had no prior value to restore
1090    /// (the first-ever undo): the pointer is written with no oplog record of its
1091    /// own, so deleting the canonical file is a complete clear.
1092    pub fn clear_undo_recovery(&self) -> Result<()> {
1093        self.write_chokepoint(|lock| self.clear_undo_recovery_locked(lock))
1094    }
1095
1096    fn clear_undo_recovery_locked(&self, _lock: &RefsLock) -> Result<()> {
1097        let path = self.undo_recovery_path();
1098        match std::fs::remove_file(&path) {
1099            Ok(()) => Ok(()),
1100            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1101            Err(e) => Err(HeddleError::Io(e)),
1102        }
1103    }
1104
1105    /// Read the heddle-internal pre-undo recovery pointer, if one has been
1106    /// recorded. Returns `None` when no undo has run in this repo.
1107    pub fn get_undo_recovery(&self) -> Result<Option<StateId>> {
1108        self.reconciled_point(LoadRequest::UndoRecovery)
1109    }
1110
1111    pub fn get_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1112        self.reconciled_point(LoadRequest::RemoteThread {
1113            remote: remote.to_string(),
1114            thread: thread.clone(),
1115        })
1116    }
1117
1118    pub fn set_remote_thread(
1119        &self,
1120        remote: &str,
1121        thread: &ThreadName,
1122        state: &StateId,
1123    ) -> Result<()> {
1124        self.set_remote_thread_raw(remote, thread, state)
1125    }
1126
1127    /// The remote-thread write entry of the write chokepoint: materialize the
1128    /// committed tail FIRST, then write canonical under the held lock
1129    /// (heddle#354 r7), so a remote-thread setter cannot bypass reconciliation.
1130    fn set_remote_thread_raw(
1131        &self,
1132        remote: &str,
1133        thread: &ThreadName,
1134        state: &StateId,
1135    ) -> Result<()> {
1136        self.write_chokepoint(|lock| self.set_remote_thread_locked(remote, thread, state, lock))
1137    }
1138
1139    /// The lock-free core of [`set_remote_thread_raw`](Self::set_remote_thread_raw):
1140    /// the caller already holds the refs lock.
1141    fn set_remote_thread_locked(
1142        &self,
1143        remote: &str,
1144        thread: &ThreadName,
1145        state: &StateId,
1146        lock: &RefsLock,
1147    ) -> Result<()> {
1148        let path = self.remote_thread_path(remote, thread)?;
1149        let content = format_state_id_text(state);
1150        let parent = path.parent().ok_or_else(|| {
1151            HeddleError::Config(format!(
1152                "invalid remote thread path for {}/{}",
1153                remote, thread
1154            ))
1155        })?;
1156        create_dir_all_durable(parent)?;
1157        self.write_string(&path, &content)?;
1158        if self.rebuild_ref_summary_index_with_lock(lock).is_err() {
1159            self.invalidate_ref_summary_index();
1160        }
1161        Ok(())
1162    }
1163
1164    pub fn delete_remote_thread(
1165        &self,
1166        remote: &str,
1167        thread: &ThreadName,
1168    ) -> Result<Option<StateId>> {
1169        self.delete_remote_thread_raw(remote, thread)
1170    }
1171
1172    /// The remote-thread delete entry of the write chokepoint: materialize the
1173    /// committed tail FIRST, then delete canonical under the held lock
1174    /// (heddle#354 r7), so a remote-thread delete cannot bypass reconciliation.
1175    fn delete_remote_thread_raw(
1176        &self,
1177        remote: &str,
1178        thread: &ThreadName,
1179    ) -> Result<Option<StateId>> {
1180        self.write_chokepoint(|lock| self.delete_remote_thread_locked(remote, thread, lock))
1181    }
1182
1183    /// The lock-free core of [`delete_remote_thread_raw`](Self::delete_remote_thread_raw):
1184    /// the caller already holds the refs lock.
1185    fn delete_remote_thread_locked(
1186        &self,
1187        remote: &str,
1188        thread: &ThreadName,
1189        lock: &RefsLock,
1190    ) -> Result<Option<StateId>> {
1191        let state = self.raw_get_remote_thread(remote, thread)?;
1192        if state.is_some() {
1193            let path = self.remote_thread_path(remote, thread)?;
1194            match std::fs::remove_file(&path) {
1195                Ok(()) => {}
1196                Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
1197                Err(e) => return Err(HeddleError::from(e)),
1198            }
1199        }
1200        if self.rebuild_ref_summary_index_with_lock(lock).is_err() {
1201            self.invalidate_ref_summary_index();
1202        }
1203        Ok(state)
1204    }
1205
1206    pub fn list_remotes(&self) -> Result<Vec<String>> {
1207        match self.reconciled_load(LoadRequest::RemoteList)? {
1208            Loaded::RemoteList(names) => Ok(names),
1209            _ => unreachable!("RemoteList request yields RemoteList"),
1210        }
1211    }
1212
1213    pub fn list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
1214        match self.reconciled_load(LoadRequest::RemoteThreadList {
1215            remote: remote.to_string(),
1216        })? {
1217            Loaded::RemoteThreadList(names) => Ok(names),
1218            _ => unreachable!("RemoteThreadList request yields RemoteThreadList"),
1219        }
1220    }
1221
1222    pub fn update_refs(&self, updates: &[RefUpdate]) -> Result<()> {
1223        if updates.is_empty() {
1224            return Ok(());
1225        }
1226        // The non-atomic write path funnels through the same chokepoint as the
1227        // atomic one: materialize the committed tail FIRST, then validate +
1228        // publish under the held lock (heddle#354 r7). Validating + writing
1229        // against the reconciled-and-materialized canonical is what stops a
1230        // non-atomic write from losing a committed record or being re-folded
1231        // over by an ancient one (cid 3329765073).
1232        self.write_chokepoint(|lock| self.update_refs_with_lock(updates, lock))
1233    }
1234
1235    pub fn resolve(&self, refspec: &str) -> Result<Option<StateId>> {
1236        resolve_refspec(
1237            refspec,
1238            || self.read_head(),
1239            |name| self.get_thread(&ThreadName::new(name)),
1240            |name| self.get_marker(&MarkerName::new(name)),
1241            || self.get_undo_recovery(),
1242        )
1243    }
1244
1245    pub fn pack_refs(&self) -> Result<()> {
1246        let lock = self.lock_refs()?;
1247        let packed_path = self.packed_refs_path();
1248        let mut packed = (*self.load_packed_refs_cached()?).clone();
1249
1250        let threads = self.list_threads_from_storage()?;
1251        for name in &threads {
1252            let path = self.thread_path(name)?;
1253            if let Some(id) = self.read_state_id_at(&path, "thread", name)? {
1254                packed.set_thread(name, id);
1255            }
1256        }
1257        let markers = self.list_markers_from_storage()?;
1258        for name in &markers {
1259            let path = self.marker_path(name)?;
1260            if let Some(id) = self.read_state_id_at(&path, "marker", name)? {
1261                packed.set_marker(name, id);
1262            }
1263        }
1264        if !packed.is_empty() {
1265            packed.save(&packed_path)?;
1266            self.invalidate_packed_refs_cache();
1267            let packed_parent = packed_path
1268                .parent()
1269                .ok_or_else(|| HeddleError::Config("invalid packed-refs path".to_string()))?;
1270            sync_directory(packed_parent)?;
1271            for name in &threads {
1272                let path = self.thread_path(name)?;
1273                if path.exists() {
1274                    std::fs::remove_file(&path)?;
1275                }
1276            }
1277            for name in &markers {
1278                let path = self.marker_path(name)?;
1279                if path.exists() {
1280                    std::fs::remove_file(&path)?;
1281                }
1282            }
1283        }
1284        if self.rebuild_ref_summary_index_with_lock(&lock).is_err() {
1285            self.invalidate_ref_summary_index();
1286        }
1287        drop(lock);
1288        Ok(())
1289    }
1290}
1291
1292impl CoreRefBackend for RefManager {
1293    type Error = HeddleError;
1294
1295    fn read_head(&self) -> Result<Head> {
1296        RefManager::read_head(self)
1297    }
1298    fn write_head(&self, head: &Head) -> Result<()> {
1299        RefManager::write_head(self, head)
1300    }
1301    fn write_head_cas(&self, expected: RefExpectation<Head>, head: &Head) -> Result<()> {
1302        RefManager::write_head_cas(self, expected, head)
1303    }
1304    async fn get_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
1305        RefManager::get_thread(self, name)
1306    }
1307    fn set_thread(&self, name: &ThreadName, state: &StateId) -> Result<()> {
1308        RefManager::set_thread(self, name, state)
1309    }
1310    fn set_thread_cas(
1311        &self,
1312        name: &ThreadName,
1313        expected: RefExpectation<StateId>,
1314        state: &StateId,
1315    ) -> Result<()> {
1316        RefManager::set_thread_cas(self, name, expected, state)
1317    }
1318    fn delete_thread(&self, name: &ThreadName) -> Result<Option<StateId>> {
1319        RefManager::delete_thread(self, name)
1320    }
1321    fn delete_thread_cas(
1322        &self,
1323        name: &ThreadName,
1324        expected: RefExpectation<StateId>,
1325    ) -> Result<()> {
1326        RefManager::delete_thread_cas(self, name, expected)
1327    }
1328    fn list_threads(&self) -> Result<Vec<ThreadName>> {
1329        RefManager::list_threads(self)
1330    }
1331    async fn get_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1332        RefManager::get_marker(self, name)
1333    }
1334    async fn create_marker(&self, name: &MarkerName, state: &StateId) -> Result<()> {
1335        RefManager::create_marker(self, name, state)
1336    }
1337    fn set_marker_cas(
1338        &self,
1339        name: &MarkerName,
1340        expected: RefExpectation<StateId>,
1341        state: &StateId,
1342    ) -> Result<()> {
1343        RefManager::set_marker_cas(self, name, expected, state)
1344    }
1345    fn delete_marker(&self, name: &MarkerName) -> Result<Option<StateId>> {
1346        RefManager::delete_marker(self, name)
1347    }
1348    fn delete_marker_cas(
1349        &self,
1350        name: &MarkerName,
1351        expected: RefExpectation<StateId>,
1352    ) -> Result<()> {
1353        RefManager::delete_marker_cas(self, name, expected)
1354    }
1355    fn list_markers(&self) -> Result<Vec<MarkerName>> {
1356        RefManager::list_markers(self)
1357    }
1358    fn update_refs(&self, updates: &[RefUpdate]) -> Result<()> {
1359        RefManager::update_refs(self, updates)
1360    }
1361    async fn resolve(&self, refspec: &str) -> Result<Option<StateId>> {
1362        RefManager::resolve(self, refspec)
1363    }
1364}
1365
1366impl RefBackend for RefManager {
1367    fn can_commit_records(&self) -> bool {
1368        self.committer.is_some()
1369    }
1370
1371    fn get_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1372        RefManager::get_remote_thread(self, remote, thread)
1373    }
1374    fn set_remote_thread(&self, remote: &str, thread: &ThreadName, state: &StateId) -> Result<()> {
1375        RefManager::set_remote_thread(self, remote, thread, state)
1376    }
1377    fn delete_remote_thread(&self, remote: &str, thread: &ThreadName) -> Result<Option<StateId>> {
1378        RefManager::delete_remote_thread(self, remote, thread)
1379    }
1380    fn list_remotes(&self) -> Result<Vec<String>> {
1381        RefManager::list_remotes(self)
1382    }
1383    fn list_remote_threads(&self, remote: &str) -> Result<Vec<ThreadName>> {
1384        RefManager::list_remote_threads(self, remote)
1385    }
1386    fn commit_and_publish(
1387        &self,
1388        records: &[OpRecord],
1389        ref_updates: &[RefUpdate],
1390        scope: Option<&str>,
1391    ) -> Result<()> {
1392        RefManager::commit_and_publish(self, records, ref_updates, scope)
1393    }
1394    fn inspect_ref_summary_index(&self) -> Result<super::RefSummaryIndexInspection> {
1395        RefManager::inspect_ref_summary_index(self)
1396    }
1397    fn rebuild_ref_summary_index(&self) -> Result<super::RefSummaryIndexInspection> {
1398        RefManager::rebuild_ref_summary_index(self)
1399    }
1400    fn pack_refs(&self) -> Result<()> {
1401        RefManager::pack_refs(self)
1402    }
1403    fn cleanup_stale_temps(&self) {
1404        RefManager::cleanup_stale_temps(self)
1405    }
1406}