1use 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
109pub const PACK_MAGIC: &[u8; 8] = b"VOLFPAK1";
111pub const PACK_VERSION: u8 = 1;
113pub const PACK_HEADER_LEN: u64 = 24;
115pub const IDX_MAGIC: &[u8; 8] = b"VOLFPIDX";
117pub const IDX_VERSION: u8 = 1;
119pub const IDX_HEADER_LEN: usize = 32;
121pub const IDX_RECORD_LEN: usize = 48;
123pub const PACK_DIR: &str = "fieldpack";
125pub const MAX_SEGMENT_BYTES: u64 = 64 * 1024 * 1024;
128pub const MAX_PACK_NODE_BYTES: u64 = u32::MAX as u64;
130
131#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
138pub enum SyncPolicy {
139 #[default]
141 Batch,
142 Each,
144}
145
146const RECORD_PREFIX: u64 = 4;
148const MIN_LOG2_BUCKETS: u8 = 8;
149const MAX_LOG2_BUCKETS: u8 = 24;
150const 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 h[16..20].copy_from_slice(&seg_id.to_le_bytes());
172 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
210fn 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
230fn 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 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#[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 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 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#[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#[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#[derive(Clone, Copy, Debug)]
504enum Src {
505 Writer,
507 ReaderOpen,
509 Sealed(usize),
511}
512
513#[derive(Clone, Copy, Debug)]
514struct Located {
515 src: Src,
516 off: u64,
517 len: u32,
518}
519
520pub 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
539fn 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
577type ScannedOpenSegment = (Vec<(NodeId, u64, u32)>, u64);
580
581fn 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 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 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 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 pub fn root(&self) -> &Path {
766 &self.root
767 }
768
769 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 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 #[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 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 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 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 pub fn has(&self, id: &NodeId) -> Result<bool> {
996 Ok(self.locate(id)?.is_some())
997 }
998
999 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 let mut fs_store = FsSeedStore::open(&root).unwrap();
1084 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 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 assert_eq!(packed.entries().unwrap(), fs_store.list_nodes().unwrap());
1108 assert_eq!(packed.list_nodes().unwrap(), fs_store.list_nodes().unwrap());
1109
1110 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 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 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 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 }
1219 let path = pack_file(&root, 0);
1220 {
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 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 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 #[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 };
1290 let path = pack_file(&root, 0);
1291 {
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 #[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}