1use 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
71pub const PACK_MAGIC: &[u8; 8] = b"VOLFPAK1";
73pub const PACK_VERSION: u8 = 1;
75pub const PACK_HEADER_LEN: u64 = 24;
77pub const IDX_MAGIC: &[u8; 8] = b"VOLFPIDX";
79pub const IDX_VERSION: u8 = 1;
81pub const IDX_HEADER_LEN: usize = 32;
83pub const IDX_RECORD_LEN: usize = 48;
85pub const PACK_DIR: &str = "fieldpack";
87pub const MAX_SEGMENT_BYTES: u64 = 64 * 1024 * 1024;
90pub const MAX_PACK_NODE_BYTES: u64 = u32::MAX as u64;
92
93const RECORD_PREFIX: u64 = 4;
95const MIN_LOG2_BUCKETS: u8 = 8;
96const MAX_LOG2_BUCKETS: u8 = 24;
97const 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 h[16..20].copy_from_slice(&seg_id.to_le_bytes());
119 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
157fn 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
177fn 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 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#[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 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 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#[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#[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#[derive(Clone, Copy, Debug)]
450enum Src {
451 Writer,
453 ReaderOpen,
455 Sealed(usize),
457}
458
459#[derive(Clone, Copy, Debug)]
460struct Located {
461 src: Src,
462 off: u64,
463 len: u32,
464}
465
466pub 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
485fn 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
523type ScannedOpenSegment = (Vec<(NodeId, u64, u32)>, u64);
526
527fn 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 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 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 pub fn root(&self) -> &Path {
694 &self.root
695 }
696
697 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 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 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 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 pub fn has(&self, id: &NodeId) -> Result<bool> {
880 Ok(self.locate(id)?.is_some())
881 }
882
883 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 let mut fs_store = FsSeedStore::open(&root).unwrap();
968 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 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 assert_eq!(packed.entries().unwrap(), fs_store.list_nodes().unwrap());
992 assert_eq!(packed.list_nodes().unwrap(), fs_store.list_nodes().unwrap());
993
994 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 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 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 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 }
1101 let path = pack_file(&root, 0);
1102 {
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 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 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}