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