Skip to main content

vole_document/store/
pack.rs

1//! Packed procedural seed store (`fieldpack`) — Phase 15.3.
2//!
3//! The reference [`FsSeedStore`](super::FsSeedStore) pays one file (+ two
4//! dentries) per content-addressed node, and a `openat`/`newfstatat`/`close`
5//! per read. This backend replaces the `seed/` namespace with a small number of
6//! append-only **segments** and one immutable hash index per segment, so a
7//! `NodeId -> bytes` lookup is a couple of bounded `pread`s plus the *unchanged*
8//! BLAKE3 re-hash gate.
9//!
10//! ## On-disk layout (`<root>/fieldpack/`)
11//!
12//! ```text
13//! seg-NNNNNNNN.pack   append-only payload; header + `u32_len || body` records
14//! seg-NNNNNNNN.idx    immutable hash index written at seal (Git-pack convention)
15//! ```
16//!
17//! A `.pack` **with** a sibling `.idx` is *sealed* (never written again); the
18//! single `.pack` **without** an `.idx` is the open append segment. No mutable
19//! manifest file is ever written, so every metadata file is immutable once
20//! written.
21//!
22//! ### `seg-N.pack`
23//!
24//! ```text
25//! header (24 B): magic b"VOLFPAK1" | version u8 | reserved [u8;7] | seg_id u32 | reserved2 u32
26//! records from offset 24: len u32 (1..=u32::MAX) | body [u8; len]
27//! ```
28//!
29//! The node id is **not** stored; it is recomputed as `NodeId::of_node(body)` on
30//! read, exactly the reference verification discipline.
31//!
32//! ### `seg-N.idx`
33//!
34//! ```text
35//! header (32 B): magic b"VOLFPIDX" | version u8 | log2_buckets u8 | reserved [u8;6]
36//!                | seg_id u32 | entry_count u64 | reserved_tail u32
37//! bucket_off: (nbuckets + 1) * u64   // absolute offset of each bucket's record run
38//! records:    entry_count * 48 B, sorted by (bucket, id)
39//!             id [u8;32] | off u64 (absolute offset of the record BODY) | len u32 | flags u32
40//! ```
41//!
42//! `bucket(id) = u16_le(id[0..2]) & (nbuckets - 1)`; the ids are already
43//! uniformly random BLAKE3 digests, so no extra hashing is needed. `nbuckets` is
44//! the smallest power of two (`log2_buckets` in `8..=24`) holding the load factor
45//! at or below 0.7.
46//!
47//! ## Exactness
48//!
49//! Node identity, canonical bytes, the whole-node re-hash gate, the strict
50//! no-clip range semantics, and the `(id, len)` enumeration set are byte-for-byte
51//! identical to [`FsSeedStore`](super::FsSeedStore). The segment framing lives
52//! *outside* the payload, so it can never change a node's bytes or id. Nothing
53//! here is on the `materialize --exact` path.
54//!
55//! ## Durability policy and crash consistency
56//!
57//! Appended records live in the open segment until it is **sealed** (an `.idx`
58//! is published atomically and a new segment is opened) or explicitly
59//! **flushed** (a durability barrier with no visibility change). [`SyncPolicy`]
60//! decides when an appended record is forced to stable storage:
61//!
62//! * [`SyncPolicy::Batch`] (default) — one `fsync` per segment, at seal and at
63//!   [`PackedSeedStore::flush`], never per record. This is the sound choice for
64//!   this format: the segment is append-only and self-describing, so crash
65//!   recovery is a **prefix** recovery (see below), and a published field
66//!   manifest is written only after [`PackedSeedStore::flush`] has made every
67//!   node it references durable (`FieldStore::put_field`).
68//! * [`SyncPolicy::Each`] — a `fdatasync` after every record (the pre-18.5
69//!   behavior). Stronger per-`put_node` durability, at one sync round-trip per
70//!   seed node.
71//!
72//! **What a crash guarantees.** A record is a `u32` length prefix followed by
73//! its body. On reopen the open segment is scanned from its header: scanning
74//! stops at the first prefix that is absent, zero, oversized, or whose body does
75//! not fit inside the file. Every record before that point is complete and is
76//! recovered; the torn tail is discarded (and truncated on a read-write reopen;
77//! a read-only reopen simply ignores it). Therefore:
78//!
79//! * no **partial** node is ever observable — a record is recovered whole or not
80//!   at all, and every fetched node is re-hashed against its id (the unchanged
81//!   [`SeedStore::get_node`] gate);
82//! * the recovered set is exactly a **prefix** of the appended record sequence;
83//! * only records not yet forced to stable storage at the crash are at risk
84//!   (with [`SyncPolicy::Batch`], at most the records since the last seal/flush);
85//! * a **sealed** segment is never written again, so its `.pack`/`.idx` pair is
86//!   immutable and its entries cannot be lost or reordered;
87//! * an unsealed segment that lost its tail can be **left incomplete**: recovery
88//!   drops the torn tail, and the caller may append from the recovered end. Any
89//!   node referenced by a durable field manifest is durable by construction
90//!   (the flush-before-publish ordering above), so a dangling manifest cannot
91//!   survive a crash.
92//!
93//! Reads use `std::os::unix::fs::FileExt::read_exact_at` (a *safe* `pread`);
94//! mmap is deliberately out of scope because the crate `forbid`s `unsafe_code`
95//! (see the Phase 15.3 design, §2.5).
96
97use core::cmp::Ordering;
98use std::cell::RefCell;
99use std::collections::HashMap;
100use std::fs;
101use std::io::{Read, Write};
102use std::path::{Path, PathBuf};
103use std::rc::Rc;
104
105use crate::error::{Error, Result};
106
107use super::{IoCounters, NodeId, SeedStore};
108
109/// Segment header magic.
110pub const PACK_MAGIC: &[u8; 8] = b"VOLFPAK1";
111/// Segment header version.
112pub const PACK_VERSION: u8 = 1;
113/// Fixed segment header length (bytes).
114pub const PACK_HEADER_LEN: u64 = 24;
115/// Index header magic.
116pub const IDX_MAGIC: &[u8; 8] = b"VOLFPIDX";
117/// Index header version.
118pub const IDX_VERSION: u8 = 1;
119/// Fixed index header length (bytes).
120pub const IDX_HEADER_LEN: usize = 32;
121/// Fixed index record length (bytes).
122pub const IDX_RECORD_LEN: usize = 48;
123/// The `fieldpack/` directory name under a store root.
124pub const PACK_DIR: &str = "fieldpack";
125/// Default append-segment size threshold: a new record that would cross this is
126/// written to a freshly-opened segment instead.
127pub const MAX_SEGMENT_BYTES: u64 = 64 * 1024 * 1024;
128/// A single canonical node may not exceed this (the record length field is `u32`).
129pub const MAX_PACK_NODE_BYTES: u64 = u32::MAX as u64;
130
131/// When the packed writer forces appended records to stable storage.
132///
133/// The default is [`SyncPolicy::Batch`]; see the module docs for the crash
134/// consistency it provides. [`SyncPolicy::Each`] restores the pre-18.5 per-record
135/// barrier for callers that want the strongest per-`put_node` durability at the
136/// cost of one sync round-trip per node.
137#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
138pub enum SyncPolicy {
139    /// One `fsync` per segment (at seal and at [`PackedSeedStore::flush`]).
140    #[default]
141    Batch,
142    /// A `fdatasync` after every appended record.
143    Each,
144}
145
146/// The 4-byte little-endian record length prefix.
147const RECORD_PREFIX: u64 = 4;
148const MIN_LOG2_BUCKETS: u8 = 8;
149const MAX_LOG2_BUCKETS: u8 = 24;
150/// Target load factor ≤ 7/10 keeps each bucket run near one record.
151const LOAD_NUM: u64 = 7;
152const LOAD_DEN: u64 = 10;
153
154fn seg_name(seg_id: u32, ext: &str) -> String {
155    format!("seg-{seg_id:08}.{ext}")
156}
157
158fn pack_path(dir: &Path, seg_id: u32) -> PathBuf {
159    dir.join(seg_name(seg_id, "pack"))
160}
161
162fn idx_path(dir: &Path, seg_id: u32) -> PathBuf {
163    dir.join(seg_name(seg_id, "idx"))
164}
165
166fn encode_pack_header(seg_id: u32) -> [u8; PACK_HEADER_LEN as usize] {
167    let mut h = [0u8; PACK_HEADER_LEN as usize];
168    h[0..8].copy_from_slice(PACK_MAGIC);
169    h[8] = PACK_VERSION;
170    // reserved 9..16 stays zero
171    h[16..20].copy_from_slice(&seg_id.to_le_bytes());
172    // reserved2 20..24 stays zero
173    h
174}
175
176fn parse_pack_header(head: &[u8; PACK_HEADER_LEN as usize], path: &Path) -> Result<u32> {
177    if &head[0..8] != PACK_MAGIC {
178        return Err(Error::integrity_mismatch(format!(
179            "packed segment {} has an unrecognised magic",
180            path.display()
181        )));
182    }
183    if head[8] != PACK_VERSION {
184        return Err(Error::unsupported_version(format!(
185            "packed segment {} has version {}, expected {PACK_VERSION}",
186            path.display(),
187            head[8]
188        )));
189    }
190    Ok(u32::from_le_bytes(head[16..20].try_into().unwrap()))
191}
192
193fn bucket_of(id: &NodeId, nbuckets: usize) -> usize {
194    let b = id.as_bytes();
195    usize::from(u16::from_le_bytes([b[0], b[1]])) & (nbuckets - 1)
196}
197
198fn choose_log2(entry_count: u64) -> u8 {
199    let mut l = MIN_LOG2_BUCKETS;
200    while l < MAX_LOG2_BUCKETS {
201        let nb = 1u64 << l;
202        if entry_count <= nb * LOAD_NUM / LOAD_DEN {
203            break;
204        }
205        l += 1;
206    }
207    l
208}
209
210/// A safe `pread`-style exact read. On unix this is
211/// [`std::os::unix::fs::FileExt::read_exact_at`]; elsewhere it degrades to a
212/// `seek`+`read_exact` (same bounded-range semantics, no `unsafe`).
213fn pread_exact(f: &fs::File, buf: &mut [u8], offset: u64) -> Result<()> {
214    #[cfg(unix)]
215    {
216        use std::os::unix::fs::FileExt;
217        f.read_exact_at(buf, offset)
218            .map_err(|e| Error::io(format!("reading packed segment at {offset}: {e}")))
219    }
220    #[cfg(not(unix))]
221    {
222        use std::io::{Seek, SeekFrom};
223        let mut f = f;
224        f.seek(SeekFrom::Start(offset))
225            .and_then(|_| f.read_exact(buf))
226            .map_err(|e| Error::io(format!("reading packed segment at {offset}: {e}")))
227    }
228}
229
230/// Build the immutable `.idx` bytes for one sealed segment.
231fn build_idx(seg_id: u32, entries: &[(NodeId, u64, u32)]) -> Vec<u8> {
232    let n = entries.len();
233    let log2 = choose_log2(n as u64);
234    let nbuckets = 1usize << log2;
235
236    // Sort by (bucket, id) so each bucket is a contiguous, id-sorted run.
237    let mut order: Vec<usize> = (0..n).collect();
238    order.sort_unstable_by(|&a, &b| {
239        let ba = bucket_of(&entries[a].0, nbuckets);
240        let bb = bucket_of(&entries[b].0, nbuckets);
241        ba.cmp(&bb)
242            .then_with(|| entries[a].0.as_bytes().cmp(entries[b].0.as_bytes()))
243    });
244
245    let rec_base = IDX_HEADER_LEN as u64 + (nbuckets as u64 + 1) * 8;
246    let mut bucket_off = vec![0u64; nbuckets + 1];
247    let mut pos = 0usize;
248    for (b, slot) in bucket_off.iter_mut().enumerate().take(nbuckets) {
249        *slot = rec_base + (pos as u64) * IDX_RECORD_LEN as u64;
250        while pos < n && bucket_of(&entries[order[pos]].0, nbuckets) == b {
251            pos += 1;
252        }
253    }
254    bucket_off[nbuckets] = rec_base + (n as u64) * IDX_RECORD_LEN as u64;
255
256    let mut out = Vec::with_capacity(rec_base as usize + n * IDX_RECORD_LEN);
257    out.extend_from_slice(IDX_MAGIC);
258    out.push(IDX_VERSION);
259    out.push(log2);
260    out.extend_from_slice(&[0u8; 6]);
261    out.extend_from_slice(&seg_id.to_le_bytes());
262    out.extend_from_slice(&(n as u64).to_le_bytes());
263    out.extend_from_slice(&[0u8; 4]);
264    debug_assert_eq!(out.len(), IDX_HEADER_LEN);
265    for off in &bucket_off {
266        out.extend_from_slice(&off.to_le_bytes());
267    }
268    for &oi in &order {
269        let (id, off, len) = entries[oi];
270        out.extend_from_slice(id.as_bytes());
271        out.extend_from_slice(&off.to_le_bytes());
272        out.extend_from_slice(&len.to_le_bytes());
273        out.extend_from_slice(&0u32.to_le_bytes());
274    }
275    out
276}
277
278/// A sealed segment: its `.idx` bucket directory is resident; the pack/index
279/// files are opened lazily.
280#[derive(Debug)]
281struct SealedSeg {
282    pack: PathBuf,
283    log2_buckets: u8,
284    bucket_off: Vec<u64>,
285    entry_count: u64,
286    idx_file: fs::File,
287    pack_file: Option<fs::File>,
288}
289
290impl SealedSeg {
291    fn open(pack: PathBuf, idx: PathBuf, seg_id: u32) -> Result<Self> {
292        let mut idx_file = fs::File::open(&idx)
293            .map_err(|e| Error::io(format!("opening packed index {}: {e}", idx.display())))?;
294        let file_len = idx_file.metadata()?.len();
295        if file_len < IDX_HEADER_LEN as u64 {
296            return Err(Error::integrity_mismatch(format!(
297                "packed index {} is shorter than its header",
298                idx.display()
299            )));
300        }
301        let mut head = [0u8; IDX_HEADER_LEN];
302        idx_file.read_exact(&mut head)?;
303        if &head[0..8] != IDX_MAGIC {
304            return Err(Error::integrity_mismatch(format!(
305                "packed index {} has an unrecognised magic",
306                idx.display()
307            )));
308        }
309        if head[8] != IDX_VERSION {
310            return Err(Error::unsupported_version(format!(
311                "packed index {} has version {}, expected {IDX_VERSION}",
312                idx.display(),
313                head[8]
314            )));
315        }
316        let log2_buckets = head[9];
317        let idx_seg_id = u32::from_le_bytes(head[16..20].try_into().unwrap());
318        let entry_count = u64::from_le_bytes(head[20..28].try_into().unwrap());
319        if idx_seg_id != seg_id {
320            return Err(Error::integrity_mismatch(format!(
321                "packed index {} names segment {idx_seg_id}, expected {seg_id}",
322                idx.display()
323            )));
324        }
325        if !(MIN_LOG2_BUCKETS..=MAX_LOG2_BUCKETS).contains(&log2_buckets) {
326            return Err(Error::integrity_mismatch(format!(
327                "packed index {} has log2_buckets {log2_buckets} outside \
328                 {MIN_LOG2_BUCKETS}..={MAX_LOG2_BUCKETS}",
329                idx.display()
330            )));
331        }
332        let nbuckets = 1usize << log2_buckets;
333        let dir_len = (nbuckets + 1) * 8;
334        if IDX_HEADER_LEN as u64 + dir_len as u64 > file_len {
335            return Err(Error::integrity_mismatch(format!(
336                "packed index {} is too short for its bucket directory",
337                idx.display()
338            )));
339        }
340        let mut dir = vec![0u8; dir_len];
341        idx_file.read_exact(&mut dir)?;
342        let mut bucket_off = Vec::with_capacity(nbuckets + 1);
343        for i in 0..=nbuckets {
344            bucket_off.push(u64::from_le_bytes(
345                dir[i * 8..i * 8 + 8].try_into().unwrap(),
346            ));
347        }
348        // The directory must be non-decreasing and describe exactly the record
349        // run, which must fit inside the file.
350        let rec_base = IDX_HEADER_LEN as u64 + dir_len as u64;
351        let rec_end = rec_base + entry_count * IDX_RECORD_LEN as u64;
352        let sane = bucket_off[0] == rec_base
353            && bucket_off[nbuckets] == rec_end
354            && rec_end <= file_len
355            && bucket_off.windows(2).all(|w| w[0] <= w[1]);
356        if !sane {
357            return Err(Error::integrity_mismatch(format!(
358                "packed index {} has an inconsistent bucket directory",
359                idx.display()
360            )));
361        }
362        Ok(SealedSeg {
363            pack,
364            log2_buckets,
365            bucket_off,
366            entry_count,
367            idx_file,
368            pack_file: None,
369        })
370    }
371
372    fn pack_file(&mut self) -> Result<&fs::File> {
373        if self.pack_file.is_none() {
374            self.pack_file = Some(fs::File::open(&self.pack).map_err(|e| {
375                Error::io(format!(
376                    "opening packed segment {}: {e}",
377                    self.pack.display()
378                ))
379            })?);
380        }
381        Ok(self.pack_file.as_ref().unwrap())
382    }
383
384    /// Locate `id` in this segment, returning `(body_offset, len)`.
385    fn lookup(&mut self, id: &NodeId) -> Result<Option<(u64, u32)>> {
386        let nbuckets = 1usize << self.log2_buckets;
387        let b = bucket_of(id, nbuckets);
388        let start = self.bucket_off[b];
389        let end = self.bucket_off[b + 1];
390        if end <= start {
391            return Ok(None);
392        }
393        let mut run = vec![0u8; (end - start) as usize];
394        pread_exact(&self.idx_file, &mut run, start)?;
395        let key = id.as_bytes();
396        let n = run.len() / IDX_RECORD_LEN;
397        let (mut lo, mut hi) = (0usize, n);
398        while lo < hi {
399            let mid = (lo + hi) / 2;
400            let base = mid * IDX_RECORD_LEN;
401            match run[base..base + 32].cmp(key.as_slice()) {
402                Ordering::Less => lo = mid + 1,
403                Ordering::Greater => hi = mid,
404                Ordering::Equal => {
405                    let off = u64::from_le_bytes(run[base + 32..base + 40].try_into().unwrap());
406                    let len = u32::from_le_bytes(run[base + 40..base + 44].try_into().unwrap());
407                    return Ok(Some((off, len)));
408                }
409            }
410        }
411        Ok(None)
412    }
413
414    fn collect_all(&mut self, out: &mut Vec<(NodeId, u64)>) -> Result<()> {
415        if self.entry_count == 0 {
416            return Ok(());
417        }
418        let nbuckets = 1usize << self.log2_buckets;
419        let rec_base = IDX_HEADER_LEN as u64 + (nbuckets as u64 + 1) * 8;
420        let total = usize::try_from(self.entry_count)
421            .ok()
422            .and_then(|n| n.checked_mul(IDX_RECORD_LEN))
423            .ok_or_else(|| Error::resource_limit("packed index record count overflows"))?;
424        let mut buf = vec![0u8; total];
425        pread_exact(&self.idx_file, &mut buf, rec_base)?;
426        for chunk in buf.as_chunks::<IDX_RECORD_LEN>().0 {
427            let id = NodeId::from_bytes(chunk[0..32].try_into().unwrap());
428            let len = u32::from_le_bytes(chunk[40..44].try_into().unwrap());
429            out.push((id, u64::from(len)));
430        }
431        Ok(())
432    }
433}
434
435/// The single open (unsealed) segment for a read-only open.
436#[derive(Debug)]
437struct OpenSeg {
438    path: PathBuf,
439    entries: HashMap<NodeId, (u64, u32)>,
440    file: Option<fs::File>,
441}
442
443impl OpenSeg {
444    fn file(&mut self) -> Result<&fs::File> {
445        if self.file.is_none() {
446            self.file = Some(fs::File::open(&self.path).map_err(|e| {
447                Error::io(format!(
448                    "opening packed segment {}: {e}",
449                    self.path.display()
450                ))
451            })?);
452        }
453        Ok(self.file.as_ref().unwrap())
454    }
455}
456
457#[derive(Debug, Default)]
458struct PackReader {
459    sealed: Vec<SealedSeg>,
460    open: Option<OpenSeg>,
461}
462
463/// The append-only writer for the open segment.
464#[derive(Debug)]
465struct PackWriter {
466    dir: PathBuf,
467    current: Option<fs::File>,
468    current_len: u64,
469    next_seg_id: u32,
470    max_segment_bytes: u64,
471    sync_policy: SyncPolicy,
472    pending: HashMap<NodeId, (u64, u32)>,
473}
474
475impl PackWriter {
476    fn ensure_open(&mut self) -> Result<()> {
477        if self.current.is_none() {
478            let path = pack_path(&self.dir, self.next_seg_id);
479            let mut f = fs::OpenOptions::new()
480                .create(true)
481                .truncate(true)
482                .read(true)
483                .write(true)
484                .open(&path)
485                .map_err(|e| {
486                    Error::io(format!("creating packed segment {}: {e}", path.display()))
487                })?;
488            f.write_all(&encode_pack_header(self.next_seg_id))?;
489            self.current = Some(f);
490            self.current_len = PACK_HEADER_LEN;
491        }
492        Ok(())
493    }
494}
495
496#[derive(Debug)]
497struct PackState {
498    reader: RefCell<PackReader>,
499    writer: Option<RefCell<PackWriter>>,
500}
501
502/// Where a located record's bytes live.
503#[derive(Clone, Copy, Debug)]
504enum Src {
505    /// The live writer's current segment (in-process, unsealed).
506    Writer,
507    /// The recovered open segment of a read-only store.
508    ReaderOpen,
509    /// Sealed segment `n`.
510    Sealed(usize),
511}
512
513#[derive(Clone, Copy, Debug)]
514struct Located {
515    src: Src,
516    off: u64,
517    len: u32,
518}
519
520/// A packed content-addressed store of canonical procedural seed nodes.
521///
522/// Implements [`SeedStore`] with the same identity, bytes, verification, range,
523/// and enumeration semantics as [`super::FsSeedStore`]; only the physical layout
524/// differs (`fieldpack/` segments instead of one file per node).
525pub struct PackedSeedStore {
526    root: PathBuf,
527    io: IoCounters,
528    state: Rc<PackState>,
529}
530
531impl std::fmt::Debug for PackedSeedStore {
532    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
533        f.debug_struct("PackedSeedStore")
534            .field("root", &self.root)
535            .finish_non_exhaustive()
536    }
537}
538
539/// Discover the segment ids under `dir`: sealed ids (with an `.idx`) ascending,
540/// and the single open id (a `.pack` with no `.idx`), if any.
541fn discover(dir: &Path) -> Result<(Vec<u32>, Option<u32>)> {
542    let mut sealed = Vec::new();
543    let mut open: Option<u32> = None;
544    let entries = match fs::read_dir(dir) {
545        Ok(e) => e,
546        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok((sealed, open)),
547        Err(e) => {
548            return Err(Error::io(format!(
549                "reading packed store {}: {e}",
550                dir.display()
551            )));
552        }
553    };
554    for entry in entries.flatten() {
555        let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
556            continue;
557        };
558        let Some(idpart) = name
559            .strip_prefix("seg-")
560            .and_then(|r| r.strip_suffix(".pack"))
561        else {
562            continue;
563        };
564        let Ok(seg_id) = idpart.parse::<u32>() else {
565            continue;
566        };
567        if idx_path(dir, seg_id).exists() {
568            sealed.push(seg_id);
569        } else {
570            open = Some(open.map_or(seg_id, |o| o.max(seg_id)));
571        }
572    }
573    sealed.sort_unstable();
574    Ok((sealed, open))
575}
576
577/// `(id, body_offset, len)` for every complete record in an unsealed segment,
578/// paired with the length of the last complete record boundary.
579type ScannedOpenSegment = (Vec<(NodeId, u64, u32)>, u64);
580
581/// Scan an unsealed segment from its header, returning every complete record's
582/// `(id, body_offset, len)` and the length of the last complete record boundary
583/// (a torn tail is truncated to this on reopen).
584fn scan_open_segment(path: &Path, expected: u32) -> Result<ScannedOpenSegment> {
585    let f = fs::File::open(path)
586        .map_err(|e| Error::io(format!("opening packed segment {}: {e}", path.display())))?;
587    let file_len = f.metadata()?.len();
588    if file_len < PACK_HEADER_LEN {
589        return Err(Error::integrity_mismatch(format!(
590            "packed segment {} has a truncated header",
591            path.display()
592        )));
593    }
594    let mut head = [0u8; PACK_HEADER_LEN as usize];
595    pread_exact(&f, &mut head, 0)?;
596    if parse_pack_header(&head, path)? != expected {
597        return Err(Error::integrity_mismatch(format!(
598            "packed segment {} names a different segment id",
599            path.display()
600        )));
601    }
602    let mut entries = Vec::new();
603    let mut p = PACK_HEADER_LEN;
604    loop {
605        if p + RECORD_PREFIX > file_len {
606            break;
607        }
608        let mut pre = [0u8; 4];
609        pread_exact(&f, &mut pre, p)?;
610        let len = u64::from(u32::from_le_bytes(pre));
611        if len == 0 || len > MAX_PACK_NODE_BYTES {
612            break;
613        }
614        if p + RECORD_PREFIX + len > file_len {
615            break;
616        }
617        let mut body = vec![0u8; len as usize];
618        pread_exact(&f, &mut body, p + RECORD_PREFIX)?;
619        let id = NodeId::of_node(&body);
620        entries.push((
621            id,
622            p + RECORD_PREFIX,
623            u32::try_from(len).unwrap_or(u32::MAX),
624        ));
625        p += RECORD_PREFIX + len;
626    }
627    Ok((entries, p))
628}
629
630impl PackedSeedStore {
631    fn pack_dir(&self) -> PathBuf {
632        self.root.join(PACK_DIR)
633    }
634
635    /// Open the store read-only: sealed segments via their `.idx`, plus the open
636    /// segment (if any) recovered by a framing scan.
637    pub fn open_read(root: impl AsRef<Path>, io: IoCounters) -> Result<Self> {
638        let root = root.as_ref().to_path_buf();
639        let dir = root.join(PACK_DIR);
640        let (sealed_ids, open_id) = discover(&dir)?;
641        let mut sealed = Vec::with_capacity(sealed_ids.len());
642        for seg_id in sealed_ids {
643            sealed.push(SealedSeg::open(
644                pack_path(&dir, seg_id),
645                idx_path(&dir, seg_id),
646                seg_id,
647            )?);
648        }
649        let open = match open_id {
650            Some(seg_id) => {
651                let path = pack_path(&dir, seg_id);
652                let (entries, _valid_len) = scan_open_segment(&path, seg_id)?;
653                let mut map = HashMap::with_capacity(entries.len());
654                for (id, off, len) in entries {
655                    map.insert(id, (off, len));
656                }
657                Some(OpenSeg {
658                    path,
659                    entries: map,
660                    file: None,
661                })
662            }
663            None => None,
664        };
665        Ok(PackedSeedStore {
666            root,
667            io,
668            state: Rc::new(PackState {
669                reader: RefCell::new(PackReader { sealed, open }),
670                writer: None,
671            }),
672        })
673    }
674
675    /// Open the store read-write, recovering (and truncating the torn tail of) an
676    /// unsealed segment so appends continue from the last complete record.
677    ///
678    /// Uses the default [`SyncPolicy::Batch`].
679    pub fn open_write(root: impl AsRef<Path>, io: IoCounters) -> Result<Self> {
680        Self::open_write_internal(root.as_ref(), io, MAX_SEGMENT_BYTES, SyncPolicy::default())
681    }
682
683    /// Open the store read-write with an explicit durability [`SyncPolicy`].
684    pub fn open_write_with_policy(
685        root: impl AsRef<Path>,
686        io: IoCounters,
687        policy: SyncPolicy,
688    ) -> Result<Self> {
689        Self::open_write_internal(root.as_ref(), io, MAX_SEGMENT_BYTES, policy)
690    }
691
692    fn open_write_internal(
693        root: &Path,
694        io: IoCounters,
695        max_segment_bytes: u64,
696        policy: SyncPolicy,
697    ) -> Result<Self> {
698        let root = root.to_path_buf();
699        let dir = root.join(PACK_DIR);
700        fs::create_dir_all(&dir)?;
701        let (sealed_ids, open_id) = discover(&dir)?;
702        let next_seg_for_fresh = sealed_ids.last().copied().map_or(0, |m| m + 1);
703        let mut sealed = Vec::with_capacity(sealed_ids.len());
704        for seg_id in sealed_ids {
705            sealed.push(SealedSeg::open(
706                pack_path(&dir, seg_id),
707                idx_path(&dir, seg_id),
708                seg_id,
709            )?);
710        }
711        let reader = PackReader { sealed, open: None };
712        let writer = match open_id {
713            Some(seg_id) => {
714                let path = pack_path(&dir, seg_id);
715                let (entries, valid_len) = scan_open_segment(&path, seg_id)?;
716                if valid_len < fs::metadata(&path)?.len() {
717                    fs::OpenOptions::new()
718                        .read(true)
719                        .write(true)
720                        .open(&path)?
721                        .set_len(valid_len)?;
722                }
723                let f = fs::OpenOptions::new()
724                    .read(true)
725                    .append(true)
726                    .open(&path)
727                    .map_err(|e| {
728                        Error::io(format!("opening packed segment {}: {e}", path.display()))
729                    })?;
730                let mut pending = HashMap::with_capacity(entries.len());
731                for (id, off, len) in entries {
732                    pending.insert(id, (off, len));
733                }
734                PackWriter {
735                    dir: dir.clone(),
736                    current: Some(f),
737                    current_len: valid_len,
738                    next_seg_id: seg_id,
739                    max_segment_bytes,
740                    sync_policy: policy,
741                    pending,
742                }
743            }
744            None => PackWriter {
745                dir: dir.clone(),
746                current: None,
747                current_len: 0,
748                next_seg_id: next_seg_for_fresh,
749                max_segment_bytes,
750                sync_policy: policy,
751                pending: HashMap::new(),
752            },
753        };
754        Ok(PackedSeedStore {
755            root,
756            io,
757            state: Rc::new(PackState {
758                reader: RefCell::new(reader),
759                writer: Some(RefCell::new(writer)),
760            }),
761        })
762    }
763
764    /// The store root.
765    pub fn root(&self) -> &Path {
766        &self.root
767    }
768
769    /// Seal the open segment: fsync it, publish `seg-N.idx` atomically, and make
770    /// it readable. Called by [`crate::field::FieldStore::sync`].
771    pub fn seal(&self) -> Result<()> {
772        let Some(w_cell) = self.state.writer.as_ref() else {
773            return Ok(());
774        };
775        let (seg_id, pending) = {
776            let mut w = w_cell.borrow_mut();
777            if w.pending.is_empty() {
778                return Ok(());
779            }
780            if let Some(f) = w.current.as_ref() {
781                f.sync_all()?;
782            }
783            let seg_id = w.next_seg_id;
784            let pending = std::mem::take(&mut w.pending);
785            w.current = None;
786            w.current_len = 0;
787            w.next_seg_id = w
788                .next_seg_id
789                .checked_add(1)
790                .ok_or_else(|| Error::resource_limit("packed segment id space exhausted"))?;
791            (seg_id, pending)
792        };
793        let entries: Vec<(NodeId, u64, u32)> = pending
794            .into_iter()
795            .map(|(id, (off, len))| (id, off, len))
796            .collect();
797        let idx = build_idx(seg_id, &entries);
798        let dir = self.pack_dir();
799        let idx_file = idx_path(&dir, seg_id);
800        #[cfg(feature = "fault-inject")]
801        crate::fault::hit("seal.before_idx");
802        crate::field::write_atomic(&idx_file, &idx)?;
803        #[cfg(feature = "fault-inject")]
804        crate::fault::hit("seal.after_idx");
805        let seg = SealedSeg::open(pack_path(&dir, seg_id), idx_file, seg_id)?;
806        self.state.reader.borrow_mut().sealed.push(seg);
807        Ok(())
808    }
809
810    fn locate(&self, id: &NodeId) -> Result<Option<Located>> {
811        if let Some(w_cell) = self.state.writer.as_ref()
812            && let Some(&(off, len)) = w_cell.borrow().pending.get(id)
813        {
814            return Ok(Some(Located {
815                src: Src::Writer,
816                off,
817                len,
818            }));
819        }
820        let mut r = self.state.reader.borrow_mut();
821        if let Some(o) = r.open.as_ref()
822            && let Some(&(off, len)) = o.entries.get(id)
823        {
824            return Ok(Some(Located {
825                src: Src::ReaderOpen,
826                off,
827                len,
828            }));
829        }
830        for (i, s) in r.sealed.iter_mut().enumerate() {
831            if let Some((off, len)) = s.lookup(id)? {
832                return Ok(Some(Located {
833                    src: Src::Sealed(i),
834                    off,
835                    len,
836                }));
837            }
838        }
839        Ok(None)
840    }
841
842    fn read_at(&self, loc: &Located, delta: u64, buf: &mut [u8]) -> Result<()> {
843        let off = loc.off + delta;
844        match loc.src {
845            Src::Writer => {
846                let w = self
847                    .state
848                    .writer
849                    .as_ref()
850                    .ok_or_else(|| Error::internal_invariant("packed writer vanished"))?
851                    .borrow();
852                let f = w.current.as_ref().ok_or_else(|| {
853                    Error::internal_invariant("packed writer has no open segment")
854                })?;
855                pread_exact(f, buf, off)
856            }
857            Src::ReaderOpen => {
858                let mut r = self.state.reader.borrow_mut();
859                let o = r
860                    .open
861                    .as_mut()
862                    .ok_or_else(|| Error::internal_invariant("packed open segment vanished"))?;
863                let f = o.file()?;
864                pread_exact(f, buf, off)
865            }
866            Src::Sealed(i) => {
867                let mut r = self.state.reader.borrow_mut();
868                let s = r
869                    .sealed
870                    .get_mut(i)
871                    .ok_or_else(|| Error::internal_invariant("packed sealed segment vanished"))?;
872                let f = s.pack_file()?;
873                pread_exact(f, buf, off)
874            }
875        }
876    }
877
878    /// Store the canonical bytes of one node (idempotent).
879    ///
880    /// With [`SyncPolicy::Each`] the record is `fdatasync`'d before returning;
881    /// with the default [`SyncPolicy::Batch`] it is visible immediately (through
882    /// the pending map) but becomes crash-durable at the next seal or
883    /// [`PackedSeedStore::flush`]. See the module docs for the exact semantics.
884    pub fn insert(&self, canonical: &[u8]) -> Result<NodeId> {
885        let id = NodeId::of_node(canonical);
886        if self.has(&id)? {
887            return Ok(id);
888        }
889        if canonical.len() as u64 > MAX_PACK_NODE_BYTES {
890            return Err(Error::resource_limit(format!(
891                "seed node is {} bytes, exceeding the packed maximum {MAX_PACK_NODE_BYTES}",
892                canonical.len()
893            )));
894        }
895        let w_cell =
896            self.state.writer.as_ref().ok_or_else(|| {
897                Error::unsupported_feature("packed seed store is opened read-only")
898            })?;
899        let need_seal = {
900            let w = w_cell.borrow();
901            w.current.is_some()
902                && w.current_len + RECORD_PREFIX + canonical.len() as u64 > w.max_segment_bytes
903        };
904        if need_seal {
905            self.seal()?;
906        }
907        let mut w = w_cell.borrow_mut();
908        w.ensure_open()?;
909        let body_off = w.current_len + RECORD_PREFIX;
910        let sync_each = w.sync_policy == SyncPolicy::Each;
911        let mut rec = Vec::with_capacity(RECORD_PREFIX as usize + canonical.len());
912        rec.extend_from_slice(&(canonical.len() as u32).to_le_bytes());
913        rec.extend_from_slice(canonical);
914        let f = w
915            .current
916            .as_mut()
917            .ok_or_else(|| Error::internal_invariant("packed writer failed to open a segment"))?;
918        // The `fault-inject` build splits the record write at the framing
919        // boundary so the court can abort between prefix and body; the shipped
920        // build writes the whole record in one call (unchanged).
921        #[cfg(feature = "fault-inject")]
922        {
923            crate::fault::hit("record.before_prefix");
924            f.write_all(&rec[..RECORD_PREFIX as usize])?;
925            crate::fault::hit("record.after_prefix");
926            f.write_all(&rec[RECORD_PREFIX as usize..])?;
927            crate::fault::hit("record.after_body");
928        }
929        #[cfg(not(feature = "fault-inject"))]
930        f.write_all(&rec)?;
931        if sync_each {
932            f.sync_data()?;
933        }
934        w.current_len = body_off + canonical.len() as u64;
935        w.pending.insert(id, (body_off, canonical.len() as u32));
936        Ok(id)
937    }
938
939    /// Force every appended record in the open segment to stable storage without
940    /// sealing it. A durability barrier, not a visibility change: the pending map
941    /// already serves every appended node. Called by
942    /// [`FieldStore::put_field`](crate::field::FieldStore) so a published manifest
943    /// only ever references nodes that are already durable.
944    pub fn flush(&self) -> Result<()> {
945        let Some(w_cell) = self.state.writer.as_ref() else {
946            return Ok(());
947        };
948        let w = w_cell.borrow();
949        if let Some(f) = w.current.as_ref() {
950            #[cfg(feature = "fault-inject")]
951            crate::fault::hit("flush.before_sync");
952            f.sync_all()?;
953            #[cfg(feature = "fault-inject")]
954            crate::fault::hit("flush.after_sync");
955        }
956        Ok(())
957    }
958
959    /// Fetch a node, verifying `NodeId::of_node(bytes) == id`.
960    pub fn fetch(&self, id: &NodeId) -> Result<Vec<u8>> {
961        let loc = self.locate(id)?.ok_or_else(|| {
962            Error::missing_external_object(format!("seed node {id} is not present"))
963        })?;
964        let mut buf = vec![0u8; loc.len as usize];
965        self.read_at(&loc, 0, &mut buf)?;
966        let actual = NodeId::of_node(&buf);
967        if actual != *id {
968            return Err(Error::integrity_mismatch(format!(
969                "seed node {id} content hashes to {actual}"
970            )));
971        }
972        self.io.add_seed(buf.len() as u64);
973        Ok(buf)
974    }
975
976    /// Fetch `len` bytes at `offset` (strict; no EOF clip, no whole-node gate).
977    pub fn fetch_range(&self, id: &NodeId, offset: u64, len: u64) -> Result<Vec<u8>> {
978        let loc = self.locate(id)?.ok_or_else(|| {
979            Error::missing_external_object(format!("seed node {id} is not present"))
980        })?;
981        let stored = u64::from(loc.len);
982        if offset.checked_add(len).is_none_or(|end| end > stored) {
983            return Err(Error::integrity_mismatch(format!(
984                "seed node {id} range [{offset}, {}) exceeds stored length {stored}",
985                offset.saturating_add(len)
986            )));
987        }
988        let mut buf = vec![0u8; usize::try_from(len).unwrap_or(usize::MAX)];
989        self.read_at(&loc, offset, &mut buf)?;
990        self.io.add_seed(buf.len() as u64);
991        Ok(buf)
992    }
993
994    /// Whether `id` is present.
995    pub fn has(&self, id: &NodeId) -> Result<bool> {
996        Ok(self.locate(id)?.is_some())
997    }
998
999    /// Every stored `(id, canonical_len)`, sorted ascending.
1000    pub fn entries(&self) -> Result<Vec<(NodeId, u64)>> {
1001        let mut out: Vec<(NodeId, u64)> = Vec::new();
1002        {
1003            let mut r = self.state.reader.borrow_mut();
1004            for s in r.sealed.iter_mut() {
1005                s.collect_all(&mut out)?;
1006            }
1007            if let Some(o) = r.open.as_ref() {
1008                for (id, (_, len)) in &o.entries {
1009                    out.push((*id, u64::from(*len)));
1010                }
1011            }
1012        }
1013        if let Some(w_cell) = self.state.writer.as_ref() {
1014            for (id, (_, len)) in &w_cell.borrow().pending {
1015                out.push((*id, u64::from(*len)));
1016            }
1017        }
1018        out.sort_unstable();
1019        Ok(out)
1020    }
1021}
1022
1023impl SeedStore for PackedSeedStore {
1024    fn put_node(&mut self, canonical: &[u8]) -> Result<NodeId> {
1025        self.insert(canonical)
1026    }
1027
1028    fn get_node(&self, id: &NodeId) -> Result<Vec<u8>> {
1029        self.fetch(id)
1030    }
1031
1032    fn get_node_range(&self, id: &NodeId, offset: u64, len: u64) -> Result<Vec<u8>> {
1033        self.fetch_range(id, offset, len)
1034    }
1035
1036    fn contains_node(&self, id: &NodeId) -> Result<bool> {
1037        self.has(id)
1038    }
1039
1040    fn list_nodes(&self) -> Result<Vec<(NodeId, u64)>> {
1041        self.entries()
1042    }
1043}
1044
1045#[cfg(test)]
1046mod tests {
1047    use super::*;
1048    use crate::store::{FsSeedStore, SeedStore};
1049
1050    fn temp_root(label: &str) -> PathBuf {
1051        let mut p = std::env::temp_dir();
1052        p.push(format!(
1053            "vole-pack-{label}-{}-{}",
1054            std::process::id(),
1055            std::time::SystemTime::now()
1056                .duration_since(std::time::UNIX_EPOCH)
1057                .unwrap()
1058                .as_nanos()
1059        ));
1060        p
1061    }
1062
1063    fn sample_nodes() -> Vec<Vec<u8>> {
1064        vec![
1065            b"alpha".to_vec(),
1066            b"beta-node".to_vec(),
1067            b"gamma-delta-epsilon".to_vec(),
1068            vec![0u8; 97],
1069            (0..255u8).collect(),
1070        ]
1071    }
1072
1073    fn pack_file(root: &Path, seg_id: u32) -> PathBuf {
1074        root.join(PACK_DIR).join(seg_name(seg_id, "pack"))
1075    }
1076
1077    #[test]
1078    fn put_get_contains_list_match_fs_semantics() {
1079        let root = temp_root("rt");
1080        let nodes = sample_nodes();
1081
1082        // The reference backend: one file per node.
1083        let mut fs_store = FsSeedStore::open(&root).unwrap();
1084        // The packed backend: segments under the same root.
1085        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1086
1087        let fs_ids: Vec<NodeId> = nodes
1088            .iter()
1089            .map(|n| fs_store.put_node(n).unwrap())
1090            .collect();
1091        let packed_ids: Vec<NodeId> = nodes.iter().map(|n| packed.insert(n).unwrap()).collect();
1092        assert_eq!(fs_ids, packed_ids, "node ids are content-derived and equal");
1093
1094        // Idempotent: re-putting does not change ids or the enumeration set.
1095        let again: Vec<NodeId> = nodes.iter().map(|n| packed.insert(n).unwrap()).collect();
1096        assert_eq!(again, packed_ids);
1097
1098        for (id, node) in packed_ids.iter().zip(&nodes) {
1099            assert!(packed.has(id).unwrap());
1100            assert!(packed.contains_node(id).unwrap());
1101            assert_eq!(&packed.fetch(id).unwrap(), node);
1102            assert_eq!(&packed.get_node(id).unwrap(), node);
1103            assert_eq!(&fs_store.get_node(id).unwrap(), node);
1104        }
1105
1106        // Same (id, len) set and ordering as the reference backend.
1107        assert_eq!(packed.entries().unwrap(), fs_store.list_nodes().unwrap());
1108        assert_eq!(packed.list_nodes().unwrap(), fs_store.list_nodes().unwrap());
1109
1110        // Seal, then a fresh read-only open sees exactly the same set.
1111        packed.seal().unwrap();
1112        assert!(pack_file(&root, 0).exists(), "segment 0 is written");
1113        assert!(
1114            idx_path(&root.join(PACK_DIR), 0).exists(),
1115            "segment 0 is sealed"
1116        );
1117        let ro = PackedSeedStore::open_read(&root, IoCounters::new()).unwrap();
1118        assert_eq!(ro.entries().unwrap(), fs_store.list_nodes().unwrap());
1119        for (id, node) in packed_ids.iter().zip(&nodes) {
1120            assert!(ro.has(id).unwrap());
1121            assert_eq!(&ro.fetch(id).unwrap(), node);
1122        }
1123
1124        fs::remove_dir_all(&root).ok();
1125    }
1126
1127    #[test]
1128    fn missing_node_is_a_miss() {
1129        let root = temp_root("miss");
1130        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1131        packed.insert(b"only one").unwrap();
1132
1133        let absent = NodeId::from_bytes([0x5A; 32]);
1134        assert!(!packed.has(&absent).unwrap());
1135        let e = packed.fetch(&absent).unwrap_err();
1136        assert_eq!(e.class(), crate::ErrorClass::MissingExternalObject);
1137        assert_eq!(
1138            packed.fetch_range(&absent, 0, 1).unwrap_err().class(),
1139            crate::ErrorClass::MissingExternalObject
1140        );
1141
1142        fs::remove_dir_all(&root).ok();
1143    }
1144
1145    #[test]
1146    fn corrupted_body_fails_the_hash_gate() {
1147        let root = temp_root("corrupt");
1148        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1149        let id = packed.insert(b"canonical node bytes").unwrap();
1150        packed.seal().unwrap();
1151
1152        // Flip one body byte in the sealed segment (body starts after the 24-byte
1153        // header and the 4-byte length prefix).
1154        let path = pack_file(&root, 0);
1155        let mut bytes = fs::read(&path).unwrap();
1156        let body_at = PACK_HEADER_LEN as usize + RECORD_PREFIX as usize;
1157        bytes[body_at] ^= 0xFF;
1158        fs::write(&path, &bytes).unwrap();
1159
1160        let ro = PackedSeedStore::open_read(&root, IoCounters::new()).unwrap();
1161        assert_eq!(
1162            ro.fetch(&id).unwrap_err().class(),
1163            crate::ErrorClass::IntegrityMismatch
1164        );
1165
1166        fs::remove_dir_all(&root).ok();
1167    }
1168
1169    #[test]
1170    fn range_reads_are_strict_and_ungated() {
1171        let root = temp_root("range");
1172        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1173        let id = packed.insert(b"0123456789").unwrap();
1174        assert_eq!(packed.fetch_range(&id, 2, 3).unwrap(), b"234");
1175        assert_eq!(packed.fetch_range(&id, 0, 10).unwrap(), b"0123456789");
1176        assert_eq!(
1177            packed.fetch_range(&id, 8, 5).unwrap_err().class(),
1178            crate::ErrorClass::IntegrityMismatch
1179        );
1180        fs::remove_dir_all(&root).ok();
1181    }
1182
1183    #[test]
1184    fn sealing_spans_multiple_segments() {
1185        let root = temp_root("segments");
1186        // A tiny limit forces a new segment every couple of records.
1187        let packed =
1188            PackedSeedStore::open_write_internal(&root, IoCounters::new(), 64, SyncPolicy::Batch)
1189                .unwrap();
1190        let nodes = sample_nodes();
1191        let ids: Vec<NodeId> = nodes.iter().map(|n| packed.insert(n).unwrap()).collect();
1192        packed.seal().unwrap();
1193
1194        let dir = root.join(PACK_DIR);
1195        let segs: Vec<u32> = discover(&dir).unwrap().0;
1196        assert!(
1197            segs.len() >= 2,
1198            "the tiny limit must roll segments, got {segs:?}"
1199        );
1200
1201        // A read-only reopen still locates every node across segments.
1202        let ro = PackedSeedStore::open_read(&root, IoCounters::new()).unwrap();
1203        assert_eq!(ro.entries().unwrap().len(), nodes.len());
1204        for (id, node) in ids.iter().zip(&nodes) {
1205            assert_eq!(&ro.fetch(id).unwrap(), node);
1206        }
1207        fs::remove_dir_all(&root).ok();
1208    }
1209
1210    #[test]
1211    fn torn_tail_is_truncated_on_reopen() {
1212        let root = temp_root("torn");
1213        {
1214            let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1215            packed.insert(b"first record").unwrap();
1216            packed.insert(b"second record").unwrap();
1217            // Deliberately do NOT seal: the segment stays open with no `.idx`.
1218        }
1219        let path = pack_file(&root, 0);
1220        // Append a torn record: a length prefix promising a body that is not there.
1221        {
1222            let mut bytes = fs::read(&path).unwrap();
1223            bytes.extend_from_slice(&9u32.to_le_bytes());
1224            bytes.extend_from_slice(b"partial");
1225            fs::write(&path, &bytes).unwrap();
1226        }
1227        let before = fs::metadata(&path).unwrap().len();
1228
1229        // Reopening read-write truncates the torn tail and keeps the complete records.
1230        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1231        let after = fs::metadata(&path).unwrap().len();
1232        assert!(after < before, "torn tail must be truncated");
1233        let entries = packed.entries().unwrap();
1234        assert_eq!(entries.len(), 2, "both complete records recovered");
1235        assert_eq!(packed.fetch(&entries[0].0).unwrap(), b"first record");
1236        assert_eq!(packed.fetch(&entries[1].0).unwrap(), b"second record");
1237
1238        // Appending after recovery keeps the file consistent.
1239        packed.insert(b"third record").unwrap();
1240        packed.seal().unwrap();
1241        let ro = PackedSeedStore::open_read(&root, IoCounters::new()).unwrap();
1242        assert_eq!(ro.entries().unwrap().len(), 3);
1243        fs::remove_dir_all(&root).ok();
1244    }
1245
1246    #[test]
1247    fn same_ingest_produces_identical_segments() {
1248        let a = temp_root("det-a");
1249        let b = temp_root("det-b");
1250        let nodes = sample_nodes();
1251        for root in [&a, &b] {
1252            let packed = PackedSeedStore::open_write(root, IoCounters::new()).unwrap();
1253            for n in &nodes {
1254                packed.insert(n).unwrap();
1255            }
1256            packed.seal().unwrap();
1257        }
1258        assert_eq!(
1259            fs::read(pack_file(&a, 0)).unwrap(),
1260            fs::read(pack_file(&b, 0)).unwrap(),
1261            "pack segments are byte-identical"
1262        );
1263        assert_eq!(
1264            fs::read(idx_path(&a.join(PACK_DIR), 0)).unwrap(),
1265            fs::read(idx_path(&b.join(PACK_DIR), 0)).unwrap(),
1266            "index segments are byte-identical"
1267        );
1268        fs::remove_dir_all(&a).ok();
1269        fs::remove_dir_all(&b).ok();
1270    }
1271
1272    /// The default [`SyncPolicy::Batch`] makes appended records visible and
1273    /// recoverable without a per-record sync: a crash that tears the trailing
1274    /// record recovers exactly the complete prefix, and no partial node is ever
1275    /// returned (the per-fetch re-hash gate still applies).
1276    #[test]
1277    fn batch_policy_recovers_a_complete_prefix_after_a_torn_tail() {
1278        let root = temp_root("batch-torn");
1279        let nodes = sample_nodes();
1280        let ids: Vec<NodeId> = {
1281            let packed = PackedSeedStore::open_write_with_policy(
1282                &root,
1283                IoCounters::new(),
1284                SyncPolicy::Batch,
1285            )
1286            .unwrap();
1287            nodes.iter().map(|n| packed.insert(n).unwrap()).collect()
1288            // Drop without sealing: the tail is durable only to the page cache.
1289        };
1290        let path = pack_file(&root, 0);
1291        // Simulate a power-loss torn tail: a length prefix whose body is short.
1292        {
1293            let mut bytes = fs::read(&path).unwrap();
1294            bytes.extend_from_slice(&11u32.to_le_bytes());
1295            bytes.extend_from_slice(b"partial");
1296            fs::write(&path, &bytes).unwrap();
1297        }
1298        let before = fs::metadata(&path).unwrap().len();
1299
1300        let packed = PackedSeedStore::open_write(&root, IoCounters::new()).unwrap();
1301        assert!(
1302            fs::metadata(&path).unwrap().len() < before,
1303            "the torn tail is truncated to the last complete record boundary"
1304        );
1305        assert_eq!(
1306            packed.entries().unwrap().len(),
1307            nodes.len(),
1308            "exactly the complete prefix is recovered"
1309        );
1310        for (id, node) in ids.iter().zip(&nodes) {
1311            assert_eq!(&packed.fetch(id).unwrap(), node);
1312        }
1313        fs::remove_dir_all(&root).ok();
1314    }
1315
1316    /// The sync policy is a durability distinction only: it must not change any
1317    /// stored byte, node id, or the `(id, len)` enumeration set. An explicit
1318    /// [`PackedSeedStore::flush`] is a no-op when nothing has been appended.
1319    #[test]
1320    fn sync_policy_does_not_change_stored_bytes() {
1321        let a = temp_root("policy-batch");
1322        let b = temp_root("policy-each");
1323        let nodes = sample_nodes();
1324        for (root, policy) in [(&a, SyncPolicy::Batch), (&b, SyncPolicy::Each)] {
1325            let packed =
1326                PackedSeedStore::open_write_with_policy(root, IoCounters::new(), policy).unwrap();
1327            packed.flush().unwrap();
1328            for n in &nodes {
1329                packed.insert(n).unwrap();
1330            }
1331            packed.flush().unwrap();
1332            packed.seal().unwrap();
1333        }
1334        let ro = PackedSeedStore::open_read(&a, IoCounters::new()).unwrap();
1335        for node in &nodes {
1336            let id = NodeId::of_node(node);
1337            assert!(ro.has(&id).unwrap());
1338            assert_eq!(&ro.fetch(&id).unwrap(), node);
1339        }
1340        assert_eq!(
1341            fs::read(pack_file(&a, 0)).unwrap(),
1342            fs::read(pack_file(&b, 0)).unwrap(),
1343            "the policy must not change the pack bytes"
1344        );
1345        assert_eq!(
1346            fs::read(idx_path(&a.join(PACK_DIR), 0)).unwrap(),
1347            fs::read(idx_path(&b.join(PACK_DIR), 0)).unwrap(),
1348            "the policy must not change the index bytes"
1349        );
1350        fs::remove_dir_all(&a).ok();
1351        fs::remove_dir_all(&b).ok();
1352    }
1353}