1#![deny(unsafe_code)]
21#![allow(warnings)]
22
23pub mod chunker;
24pub mod classifier;
25pub mod compaction;
26pub mod config;
27pub mod delta_builder;
28pub mod dictionary;
29pub mod file_categorizer;
30pub mod flatten;
31#[cfg(feature = "pipeline-parallelism")]
32pub mod pipeline;
33pub mod rw;
34#[cfg(feature = "sparse-index")]
35pub mod sparse_index;
36pub mod turnover;
37
38pub use config::{
39 profile, CategorizerConfig, ChunkingConfig, CodecRegistry, CodecTunables, Defaults,
40 DictionaryConfig, EncryptionConfig, TournamentConfig, WriteConfig,
41};
42
43use std::collections::{HashMap, HashSet};
44use std::path::{Path, PathBuf};
45
46use crate::chunker::FastCDC;
47use limnifs_core::codec::CODEC_REFERENCED;
48use limnifs_core::slab_store::SlabStore;
49use limnifs_core::{
50 compute_merkle_root, hash_empty_section, hash_section, parse_manifest_header, parse_slab_index,
51 ManifestCursor, ManifestHeader, SectionHashes, FEATURE_FLAGS_SECTION_VERSION,
52 HISTORY_SECTION_VERSION, INODE_FLAG_INLINE_DATA, INODE_FLAG_SHARED_INLINE,
53 METADATA_REFERENCE_SECTION_VERSION_2, SLAB_INDEX_SECTION_VERSION,
54};
55use limnifs_format::{ManifestRoot, SlabId};
56
57pub const INLINE_THRESHOLD: usize = 4096;
60
61pub const MMAP_READ_THRESHOLD: usize = 1024 * 1024;
67
68pub const WHOLE_FILE_MAX_SIZE: usize = 64 * 1024 * 1024;
73
74pub const MAX_SLAB_TOTAL_BYTES: usize = 60 * 1024 * 1024;
79
80const SLAB_HEADER_LEN: usize = 56;
84
85pub const METADATA_EXTERNALIZE_THRESHOLD: usize =
93 limnifs_core::metadata_reference::DEFAULT_INLINE_METADATA_MAX_BYTES as usize - 24 * 1024;
94
95pub const METADATA_LARGE_BLOB_THRESHOLD: usize = 256 * 1024;
100
101pub const METADATA_SMALL_BLOB_QUALITY: i32 = 5;
104
105pub const METADATA_LARGE_BLOB_QUALITY: i32 = 2;
110
111#[derive(Clone, Debug)]
114pub struct SlabArtifact {
115 pub id: SlabId,
116 pub bytes: Vec<u8>,
117 pub locator: String,
118 pub drop_ids: Vec<[u8; 32]>,
122}
123
124#[derive(Clone, Debug)]
128pub struct MetadataSidecar {
129 pub bytes: Vec<u8>,
130 pub locator: String,
131}
132
133#[derive(Clone, Debug)]
135pub struct WriteArtifact {
136 pub bytes: Vec<u8>,
137 pub merkle_root: ManifestRoot,
138 pub slabs: Vec<SlabArtifact>,
141 pub metadata_sidecar: Option<MetadataSidecar>,
145 pub inode_count: usize,
146 pub file_count: usize,
147 pub dir_count: usize,
148 pub drop_count: usize,
149 pub root_inode_number: u64,
154}
155
156impl WriteArtifact {
157 #[must_use]
161 pub fn slab_bytes(&self) -> Option<&[u8]> {
162 if self.slabs.len() == 1 {
163 Some(&self.slabs[0].bytes)
164 } else {
165 None
166 }
167 }
168
169 #[must_use]
171 pub fn slab_locator(&self) -> Option<&str> {
172 if self.slabs.len() == 1 {
173 Some(&self.slabs[0].locator)
174 } else {
175 None
176 }
177 }
178}
179
180#[derive(Debug)]
182pub enum WriteError {
183 Io(std::io::Error),
184 UnsupportedFileType {
189 path: PathBuf,
190 kind: String,
191 },
192}
193
194impl std::fmt::Display for WriteError {
195 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
196 match self {
197 Self::Io(e) => write!(f, "I/O error: {e}"),
198 Self::UnsupportedFileType { path, kind } => write!(
199 f,
200 "unsupported file type ({kind}): {} — limnifs stores files, \
201 directories, and symlinks; remove the entry or file an issue \
202 if you need it carried",
203 path.display()
204 ),
205 }
206 }
207}
208
209impl std::error::Error for WriteError {}
210
211impl From<std::io::Error> for WriteError {
212 fn from(e: std::io::Error) -> Self {
213 Self::Io(e)
214 }
215}
216
217pub fn write_directory(root: &Path) -> Result<WriteArtifact, WriteError> {
231 write_directory_with_config(root, &WriteConfig::default_v0_1())
232}
233
234pub fn write_stream<R: std::io::Read>(
250 name: &str,
251 reader: R,
252 config: &WriteConfig,
253) -> Result<WriteArtifact, WriteError> {
254 let mut ctx = WriteContext::new();
255 ctx.categorizers_disabled = config.categorizers.is_empty();
256 ctx.rw_mode = matches!(config.mode, crate::config::ImageMode::ReadWrite(_));
257 ctx.auto_turnover = config.turnover_threshold > 0;
258 ctx.collect_dict_samples = config.dictionaries.enabled;
259
260 let drop_id_root = [0u8; 32]; let pending = PendingFile {
265 path: std::path::PathBuf::from(name),
266 inode_number: 1,
267 file_len: 0, mtime_ns: 0,
269 };
270 ctx.pending_files.push(pending);
271 ctx.root_inode_number = 1;
272
273 let chunker = ctx.chunker.clone();
275 let chunks = chunker.chunk_reader(reader)?;
276
277 let total_len: u64 = chunks.iter().map(|c| c.len() as u64).sum();
279
280 let text_codec = config.text_codec_id().unwrap_or(0x04);
284 let binary_codec = config.binary_codec_id().unwrap_or(0x01);
285 let tunables = config.to_core_tunables();
286 let classifier = ctx.classifier;
287 let registry = config
288 .codec_registry()
289 .map_err(|e| WriteError::Io(std::io::Error::other(format!("codec registry: {e}"))))?;
290 let tournament_codec_ids: Vec<u8> = config
291 .tournament
292 .codecs
293 .iter()
294 .filter_map(|n| registry.lookup_by_name(n))
295 .collect();
296 let tournament = TournamentSpec {
297 codec_ids: tournament_codec_ids,
298 min_size: config.tournament.min_size_threshold as usize,
299 skip_for_binary: config.tournament.skip_for_binary,
300 short_circuit_permille: config.tournament.short_circuit_threshold,
301 };
302
303 let mut drops: Vec<RawDrop> = Vec::with_capacity(chunks.len());
304 let mut slices: Vec<PendingSlice> = Vec::with_capacity(chunks.len());
305 let mut offset: u64 = 0;
306 for chunk in &chunks {
307 let drop_id = hash_section(chunk);
308 slices.push(PendingSlice {
309 drop_id,
310 file_byte_start: offset,
311 file_byte_end: offset + chunk.len() as u64,
312 });
313 offset += chunk.len() as u64;
314 let class = classifier.classify(chunk);
315 let (codec_id, compressed) = compress_chunk_with_tournament(
316 chunk,
317 class,
318 text_codec,
319 binary_codec,
320 &tunables,
321 &tournament,
322 );
323 drops.push((drop_id, chunk.clone(), compressed, codec_id));
324 }
325 let _ = drop_id_root;
326
327 let result = ChunkedFileResult { drops, slices };
329 let pf = ctx.pending_files[0].clone();
330 ctx.merge_chunked_file(&pf, result);
331 ctx.pending_files[0].file_len = total_len;
333 if let Some(inode) = ctx.inodes.last_mut() {
336 if let PendingContent::DropBacked { file_len, .. } = &mut inode.content {
337 *file_len = total_len;
338 }
339 }
340
341 ctx.train_and_apply_dictionary(&config.dictionaries);
342 let artifact = ctx.assemble();
343 Ok(artifact)
344}
345
346pub fn write_layer(
387 base_image: &Path,
388 root: &Path,
389 config: &WriteConfig,
390) -> Result<WriteArtifact, WriteError> {
391 let (base_drop_index, base_root) = load_base_drop_index(base_image)?;
393
394 let mut ctx = WriteContext::new();
395 ctx.categorizers_disabled = config.categorizers.is_empty();
396 ctx.rw_mode = matches!(config.mode, crate::config::ImageMode::ReadWrite(_));
397 ctx.auto_turnover = config.turnover_threshold > 0;
398 ctx.collect_dict_samples = config.dictionaries.enabled;
399 ctx.inline_threshold = config.defaults.inline_threshold as usize;
400 ctx.metadata_externalize_threshold = config.defaults.metadata_externalize_threshold;
401 ctx.emit_shared_inline = config.defaults.shared_inline;
402 ctx.base_drop_index = Some(base_drop_index);
403 ctx.base_root = Some(base_root);
404
405 let root_inode_number = ctx.walk(root)?;
407 ctx.root_inode_number = root_inode_number;
408 write_directory_body(&mut ctx, config)?;
409 Ok(ctx.assemble())
410}
411
412fn load_base_drop_index(
416 base_image: &Path,
417) -> Result<(std::collections::HashSet<[u8; 32]>, [u8; 32]), WriteError> {
418 let manifest_bytes = std::fs::read(base_image)?;
419 let mut cursor = ManifestCursor::new(&manifest_bytes);
420 let _ = parse_manifest_header(&mut cursor).map_err(io_core)?;
421 let _ = limnifs_core::parse_feature_flags_section(&mut cursor);
423 let _ = limnifs_core::parse_metadata_reference(&mut cursor);
424 let slab_index = parse_slab_index(&mut cursor).map_err(io_core)?;
425 let store = SlabStore::load_mmap(base_image, &slab_index).map_err(io_core)?;
426 let drop_set: std::collections::HashSet<[u8; 32]> = store.drop_index_keys().copied().collect();
427 let root = *compute_merkle_root_from_sections(&manifest_bytes).as_bytes();
433 Ok((drop_set, root))
434}
435
436fn compute_merkle_root_from_sections(manifest: &[u8]) -> ManifestRoot {
442 use limnifs_core::SectionHashes;
443 let mut cursor = ManifestCursor::new(manifest);
444 let header_start = 0;
445 if parse_manifest_header(&mut cursor).is_err() {
446 return ManifestRoot::from_bytes([0u8; 32]);
448 }
449 let header_end = cursor.position();
450 let flags_start = header_end;
452 let flags_end = match limnifs_core::parse_feature_flags_section(&mut cursor) {
453 Ok(_) => cursor.position(),
454 Err(_) => flags_start,
455 };
456 let meta_ref_start = flags_end;
457 let metadata_reference = match limnifs_core::parse_metadata_reference(&mut cursor) {
458 Ok(m) => Some(m),
459 Err(_) => None,
460 };
461 let meta_ref_end = cursor.position();
462 let slab_index_start = meta_ref_end;
463 let _ = parse_slab_index(&mut cursor);
464 let slab_index_end = cursor.position();
465 let history_start = slab_index_end;
466 let _ = limnifs_core::parse_history(&mut cursor);
467 let history_end = cursor.position();
468
469 let hashes = SectionHashes {
470 metadata: metadata_reference
471 .map(|m| m.metadata_hash)
472 .unwrap_or_else(hash_empty_section),
473 format_header: hash_section(&manifest[header_start..header_end]),
474 feature_flags: hash_section(&manifest[flags_start..flags_end]),
475 metadata_reference: hash_section(&manifest[meta_ref_start..meta_ref_end]),
476 slab_index: hash_section(&manifest[slab_index_start..slab_index_end]),
477 crypto_params: hash_empty_section(),
478 ec_params: hash_empty_section(),
479 dms_policy: hash_empty_section(),
480 delta_linkage: hash_empty_section(),
481 history: hash_section(&manifest[history_start..history_end]),
482 };
483 compute_merkle_root(&hashes)
484}
485
486fn io_core(e: limnifs_core::CoreError) -> WriteError {
487 WriteError::Io(std::io::Error::other(format!("base image load: {e}")))
488}
489
490fn write_directory_body(ctx: &mut WriteContext, config: &WriteConfig) -> Result<(), WriteError> {
495 use rayon::prelude::*;
496
497 ctx.metadata_codec = config
498 .metadata_codec_id()
499 .unwrap_or(limnifs_core::codec::CODEC_BROTLI);
500
501 let pending = std::mem::take(&mut ctx.pending_files);
502 if pending.is_empty() {
503 return Ok(());
504 }
505 ctx.inline_threshold = config.defaults.inline_threshold as usize;
506 ctx.metadata_externalize_threshold = config.defaults.metadata_externalize_threshold;
507 ctx.emit_shared_inline = config.defaults.shared_inline;
508 let chunker = ctx.chunker.clone();
509 let classifier = ctx.classifier;
510 let text_codec = config.text_codec_id().unwrap_or(0x04);
511 let binary_codec = config.binary_codec_id().unwrap_or(0x01);
512 let tunables = config.to_core_tunables();
513 let use_categorizers = !config.categorizers.is_empty();
514 let skip_chunking = config.skip_chunking;
515 let registry = config
516 .codec_registry()
517 .map_err(|e| WriteError::Io(std::io::Error::other(format!("codec registry: {e}"))))?;
518 let tournament_codec_ids: Vec<u8> = config
519 .tournament
520 .codecs
521 .iter()
522 .filter_map(|n| registry.lookup_by_name(n))
523 .collect();
524 let tournament_spec = TournamentSpec {
525 codec_ids: tournament_codec_ids,
526 min_size: config.tournament.min_size_threshold as usize,
527 skip_for_binary: config.tournament.skip_for_binary,
528 short_circuit_permille: config.tournament.short_circuit_threshold,
529 };
530 let base_drop_index = ctx.base_drop_index.as_ref();
531 let inline_threshold = ctx.inline_threshold;
532 let results: Vec<ChunkedFileResult> = pending
533 .par_iter()
534 .map(|pf| {
535 process_file(
536 pf,
537 &chunker,
538 classifier,
539 text_codec,
540 binary_codec,
541 &tunables,
542 use_categorizers,
543 skip_chunking,
544 &tournament_spec,
545 base_drop_index,
546 inline_threshold,
547 )
548 })
549 .collect::<Result<Vec<_>, _>>()?;
550
551 for (pf, result) in pending.iter().zip(results) {
552 ctx.merge_chunked_file(pf, result);
553 }
554 ctx.train_and_apply_dictionary(&config.dictionaries);
555 Ok(())
556}
557
558pub fn write_directory_with_config(
560 root: &Path,
561 config: &WriteConfig,
562) -> Result<WriteArtifact, WriteError> {
563 let mut ctx = WriteContext::new();
564 ctx.categorizers_disabled = config.categorizers.is_empty();
565 ctx.rw_mode = matches!(config.mode, crate::config::ImageMode::ReadWrite(_));
566 ctx.auto_turnover = config.turnover_threshold > 0;
567 ctx.collect_dict_samples = config.dictionaries.enabled;
568
569 write_directory_streaming(&mut ctx, root, config)?;
570 Ok(ctx.assemble())
571}
572
573fn write_directory_streaming(
587 ctx: &mut WriteContext,
588 root: &Path,
589 config: &WriteConfig,
590) -> Result<(), WriteError> {
591 use rayon::prelude::*;
592
593 ctx.metadata_codec = config
594 .metadata_codec_id()
595 .unwrap_or(limnifs_core::codec::CODEC_BROTLI);
596
597 let chunker = ctx.chunker.clone();
598 let classifier = ctx.classifier;
599 let text_codec = config.text_codec_id().unwrap_or(0x04);
600 let binary_codec = config.binary_codec_id().unwrap_or(0x01);
601 let tunables = config.to_core_tunables();
602 let use_categorizers = !config.categorizers.is_empty();
603 let skip_chunking = config.skip_chunking;
604 let registry = config
605 .codec_registry()
606 .map_err(|e| WriteError::Io(std::io::Error::other(format!("codec registry: {e}"))))?;
607 let tournament_codec_ids: Vec<u8> = config
608 .tournament
609 .codecs
610 .iter()
611 .filter_map(|n| registry.lookup_by_name(n))
612 .collect();
613 let tournament_spec = TournamentSpec {
614 codec_ids: tournament_codec_ids,
615 min_size: config.tournament.min_size_threshold as usize,
616 skip_for_binary: config.tournament.skip_for_binary,
617 short_circuit_permille: config.tournament.short_circuit_threshold,
618 };
619 let base_drop_index = ctx.base_drop_index.clone();
622 let inline_threshold = ctx.inline_threshold;
623
624 ctx.inline_threshold = config.defaults.inline_threshold as usize;
625 ctx.metadata_externalize_threshold = config.defaults.metadata_externalize_threshold;
626 ctx.emit_shared_inline = config.defaults.shared_inline;
627
628 const PIPELINE_CAPACITY: usize = 256;
632 let (tx, rx) = std::sync::mpsc::sync_channel::<PendingFile>(PIPELINE_CAPACITY);
633 ctx.pending_sink = Some(tx);
634
635 let (root_inode_number, mut results): (
636 u64,
637 Vec<(usize, PendingFile, Result<ChunkedFileResult, WriteError>)>,
638 ) = std::thread::scope(|scope| {
639 let producer = {
640 let ctx = &mut *ctx;
641 let root = root;
642 scope.spawn(move || {
643 let r = ctx.walk(root);
644 ctx.pending_sink = None;
648 r
649 })
650 };
651 let results = rx
654 .into_iter()
655 .enumerate()
656 .par_bridge()
657 .map(|(i, pf)| {
658 let r = process_file(
659 &pf,
660 &chunker,
661 classifier,
662 text_codec,
663 binary_codec,
664 &tunables,
665 use_categorizers,
666 skip_chunking,
667 &tournament_spec,
668 base_drop_index.as_ref(),
669 inline_threshold,
670 );
671 (i, pf, r)
672 })
673 .collect();
674 let joined = producer
675 .join()
676 .unwrap_or_else(|_| {
677 Err(WriteError::Io(std::io::Error::other(
678 "walk thread panicked",
679 )))
680 })
681 .map(|n| (n, results));
682 joined
685 })?;
686 ctx.pending_sink = None;
687 ctx.root_inode_number = root_inode_number;
688
689 results.sort_unstable_by_key(|(i, _, _)| *i);
690 for (_, pf, r) in results {
694 ctx.merge_chunked_file(&pf, r?);
695 }
696 ctx.train_and_apply_dictionary(&config.dictionaries);
697 Ok(())
698}
699
700pub(crate) type RawDrop = ([u8; 32], Vec<u8>, std::sync::Arc<[u8]>, u8);
707pub(crate) struct ChunkedFileResult {
709 drops: Vec<RawDrop>, slices: Vec<PendingSlice>,
711}
712
713struct TournamentSpec {
721 codec_ids: Vec<u8>,
725 min_size: usize,
729 skip_for_binary: bool,
733 short_circuit_permille: u32,
738}
739
740fn process_whole_file_drop(
753 pf: &PendingFile,
754 data: &[u8],
755 cat: file_categorizer::Categorization,
756 tunables: &limnifs_core::codec::CodecTunables,
757) -> Result<ChunkedFileResult, WriteError> {
758 let _ = pf;
759 let drop_id = hash_section(data);
760 let file_len = u64::try_from(data.len()).unwrap_or(u64::MAX);
761
762 let (mut best_codec, mut best_compressed): (u8, std::sync::Arc<[u8]>) =
767 match limnifs_core::codec::compress_with_tunables(
768 limnifs_core::codec::CODEC_BROTLI,
769 data,
770 tunables,
771 ) {
772 Ok(c) => (limnifs_core::codec::CODEC_BROTLI, c.into()),
773 Err(_) => match limnifs_core::codec::compress_with_tunables(
774 limnifs_core::codec::CODEC_ZSTD,
775 data,
776 tunables,
777 ) {
778 Ok(c) => (limnifs_core::codec::CODEC_ZSTD, c.into()),
779 Err(_) => (limnifs_core::codec::CODEC_STORE, data.to_vec().into()),
780 },
781 };
782
783 let brotli_ratio = best_compressed.len() as f64 / data.len() as f64;
787 if brotli_ratio > 0.05 {
788 if let Ok(zstd_c) = limnifs_core::codec::compress_with_tunables(
789 limnifs_core::codec::CODEC_ZSTD,
790 data,
791 tunables,
792 ) {
793 if zstd_c.len() < best_compressed.len() {
794 best_codec = limnifs_core::codec::CODEC_ZSTD;
795 best_compressed = zstd_c.into();
796 }
797 }
798 }
799
800 let general_ratio = best_compressed.len() as f64 / data.len() as f64;
806 if general_ratio > 0.15 || cat.codec_id == limnifs_core::codec::CODEC_RICEPP {
807 let spec_result = if cat.codec_id == limnifs_core::codec::CODEC_FSST_BROTLI {
811 limnifs_core::codec::fsst_brotli::compress_with_baseline(data, Some(&best_compressed))
812 } else {
813 limnifs_core::codec::compress(cat.codec_id, data)
814 };
815 if let Ok(spec_c) = spec_result {
816 if spec_c.len() < best_compressed.len() {
817 best_codec = cat.codec_id;
818 best_compressed = spec_c.into();
819 }
820 }
821 }
822
823 Ok(ChunkedFileResult {
824 drops: vec![(drop_id, data.to_vec(), best_compressed, best_codec)],
825 slices: vec![PendingSlice {
826 drop_id,
827 file_byte_start: 0,
828 file_byte_end: file_len,
829 }],
830 })
831}
832
833fn compress_chunk_with_tournament(
855 chunk: &[u8],
856 class: classifier::Class,
857 text_codec: u8,
858 binary_codec: u8,
859 tunables: &limnifs_core::codec::CodecTunables,
860 tournament: &TournamentSpec,
861) -> (u8, std::sync::Arc<[u8]>) {
862 use classifier::Class;
863
864 let preferred = match class {
865 Class::Binary => binary_codec,
866 Class::Text | Class::Code | Class::Sparse => text_codec,
867 _ => limnifs_core::codec::CODEC_STORE,
868 };
869
870 if preferred == limnifs_core::codec::CODEC_STORE {
871 return (limnifs_core::codec::CODEC_STORE, chunk.to_vec().into());
872 }
873 if class == Class::Binary && tournament.skip_for_binary {
874 return compress_chunk_one(chunk, preferred, tunables);
875 }
876 if chunk.len() < tournament.min_size {
877 return compress_chunk_one(chunk, preferred, tunables);
878 }
879
880 let mut best: Option<(u8, std::sync::Arc<[u8]>)> = None;
881 for &codec_id in &tournament.codec_ids {
882 if codec_id == limnifs_core::codec::CODEC_STORE {
883 continue;
884 }
885 let c = match limnifs_core::codec::compress_with_tunables(codec_id, chunk, tunables) {
886 Ok(c) => c,
887 Err(_) => continue,
888 };
889 if c.len() >= chunk.len() {
890 continue;
891 }
892 let ratio_permille = (c.len() as u64 * 1000 / chunk.len() as u64) as u32;
893 let is_best_so_far = best.as_ref().map_or(true, |(_, b)| c.len() < b.len());
894 if is_best_so_far {
895 best = Some((codec_id, c.into()));
896 }
897 if tournament.short_circuit_permille > 0
898 && ratio_permille <= tournament.short_circuit_permille
899 {
900 break;
901 }
902 }
903
904 best.unwrap_or_else(|| (limnifs_core::codec::CODEC_STORE, chunk.to_vec().into()))
905}
906
907fn compress_chunk_one(
910 chunk: &[u8],
911 codec_id: u8,
912 tunables: &limnifs_core::codec::CodecTunables,
913) -> (u8, std::sync::Arc<[u8]>) {
914 if codec_id == limnifs_core::codec::CODEC_STORE {
915 return (limnifs_core::codec::CODEC_STORE, chunk.to_vec().into());
916 }
917 match limnifs_core::codec::compress_with_tunables(codec_id, chunk, tunables) {
918 Ok(c) if c.len() < chunk.len() => (codec_id, c.into()),
919 _ => (limnifs_core::codec::CODEC_STORE, chunk.to_vec().into()),
920 }
921}
922
923fn process_file(
927 pf: &PendingFile,
928 chunker: &FastCDC,
929 classifier: classifier::Classifier,
930 text_codec: u8,
931 binary_codec: u8,
932 tunables: &limnifs_core::codec::CodecTunables,
933 use_categorizers: bool,
934 skip_chunking: bool,
935 tournament: &TournamentSpec,
936 base_drop_index: Option<&std::collections::HashSet<[u8; 32]>>,
937 inline_threshold: usize,
938) -> Result<ChunkedFileResult, WriteError> {
939 let file_len_estimate = std::fs::metadata(&pf.path)
950 .map(|m| m.len() as usize)
951 .unwrap_or(0);
952 let data: Vec<u8> = if file_len_estimate >= MMAP_READ_THRESHOLD {
953 let file = std::fs::File::open(&pf.path)?;
954 #[allow(unsafe_code)]
955 let mmap = unsafe { memmap2::Mmap::map(&file) }.map_err(WriteError::Io)?;
956 Vec::from(&mmap[..])
960 } else {
961 std::fs::read(&pf.path)?
962 };
963 let file_len = data.len();
964
965 if skip_chunking && file_len > inline_threshold {
972 let drop_id = hash_section(&data);
973 let class = classifier.classify(&data);
974 let preferred_codec = match class {
975 classifier::Class::Binary => binary_codec,
976 _ => text_codec,
977 };
978 let (codec_id, compressed): (u8, std::sync::Arc<[u8]>) =
979 match limnifs_core::codec::compress_with_tunables(preferred_codec, &data, tunables) {
980 Ok(c) if c.len() < data.len() => (preferred_codec, c.into()),
981 _ => (limnifs_core::codec::CODEC_STORE, data.to_vec().into()),
982 };
983 return Ok(ChunkedFileResult {
984 drops: vec![(drop_id, data, compressed, codec_id)],
985 slices: vec![PendingSlice {
986 drop_id,
987 file_byte_start: 0,
988 file_byte_end: file_len as u64,
989 }],
990 });
991 }
992
993 if use_categorizers {
994 if let Some(cat) = file_categorizer::default_registry().categorize(&pf.path, &data) {
995 let needs_whole_file = matches!(
996 cat.codec_id,
997 limnifs_core::codec::CODEC_FLAC | limnifs_core::codec::CODEC_RICEPP
998 );
999 if needs_whole_file || file_len <= WHOLE_FILE_MAX_SIZE {
1000 return process_whole_file_drop(pf, &data, cat, tunables);
1001 }
1002 }
1003 }
1004
1005 let chunks = chunker.chunk_slice(&data);
1006 let mut slices = Vec::with_capacity(chunks.len());
1007 let mut file_offset: u64 = 0;
1008 let mut seen_in_file: std::collections::HashSet<[u8; 32]> =
1009 std::collections::HashSet::with_capacity(chunks.len());
1010
1011 let mut unique_chunks: Vec<(&[u8], [u8; 32])> = Vec::with_capacity(chunks.len());
1014 for chunk in &chunks {
1015 let chunk_len = u64::try_from(chunk.len()).expect("chunk len fits u64");
1016 let drop_id = hash_section(chunk);
1017 slices.push(PendingSlice {
1018 drop_id,
1019 file_byte_start: file_offset,
1020 file_byte_end: file_offset + chunk_len,
1021 });
1022 file_offset += chunk_len;
1023 if seen_in_file.insert(drop_id) {
1024 unique_chunks.push((chunk, drop_id));
1025 }
1026 }
1027
1028 use rayon::prelude::*;
1040 thread_local! {
1041 static COMPRESS_CACHE: std::cell::RefCell<std::collections::HashMap<[u8; 32], (u8, std::sync::Arc<[u8]>)>> =
1042 std::cell::RefCell::new(std::collections::HashMap::new());
1043 }
1044 const COMPRESS_CACHE_MAX_ENTRIES: usize = 100_000;
1045 let drops: Vec<RawDrop> = unique_chunks
1046 .par_iter()
1047 .map(|(chunk, drop_id)| {
1048 if let Some(base) = base_drop_index {
1051 if base.contains(drop_id) {
1052 return (*drop_id, Vec::new(), Vec::new().into(), CODEC_REFERENCED);
1053 }
1054 }
1055 let class = classifier.classify(chunk);
1056 let cached = COMPRESS_CACHE.with(|c| {
1059 c.borrow()
1060 .get(drop_id)
1061 .map(|(cid, comp)| (*cid, comp.clone()))
1062 });
1063 let (codec_id, compressed) = if let Some(c) = cached {
1064 c
1065 } else {
1066 let new = compress_chunk_with_tournament(
1067 chunk,
1068 class,
1069 text_codec,
1070 binary_codec,
1071 tunables,
1072 tournament,
1073 );
1074 COMPRESS_CACHE.with(|c| {
1076 let mut cache = c.borrow_mut();
1077 if cache.len() < COMPRESS_CACHE_MAX_ENTRIES {
1078 cache.insert(*drop_id, new.clone());
1080 }
1081 });
1082 new
1083 };
1084 (*drop_id, chunk.to_vec(), compressed, codec_id)
1085 })
1086 .collect();
1087
1088 let _ = file_len;
1089 Ok(ChunkedFileResult { drops, slices })
1090}
1091
1092struct PendingDrop {
1093 id: [u8; 32],
1094 plaintext_len: u32,
1100 compressed: std::sync::Arc<[u8]>,
1101 codec: u8,
1102 dict_id: u8,
1106 plaintext: Option<Vec<u8>>,
1111}
1112
1113impl PendingDrop {
1114 fn len_in_window(&self) -> u32 {
1118 u32::try_from(self.compressed.len()).expect("compressed fits u32")
1119 }
1120
1121 fn plaintext_len_value(&self) -> u32 {
1123 self.plaintext_len
1124 }
1125
1126 fn slab_footprint(&self) -> usize {
1129 48 + self.compressed.len()
1130 }
1131}
1132
1133struct PendingSlice {
1138 drop_id: [u8; 32],
1139 file_byte_start: u64,
1140 file_byte_end: u64,
1141}
1142
1143#[derive(Clone)]
1146struct PendingFile {
1147 inode_number: u64,
1148 path: PathBuf,
1149 mtime_ns: u64,
1150 file_len: u64,
1151}
1152
1153struct PendingInode {
1154 number: u64,
1155 mode: u32,
1156 mtime_ns: u64,
1157 content: PendingContent,
1158}
1159
1160enum PendingContent {
1161 Inline(Vec<u8>),
1162 Symlink(String),
1164 DropBacked {
1165 file_len: u64,
1166 slices: Vec<PendingSlice>,
1167 },
1168 Directory(Vec<(String, u64, u8)>),
1169}
1170
1171struct DirNode {
1172 entries: Vec<(String, u64, u8)>,
1173 bytes: Vec<u8>,
1174 hash: [u8; 32],
1175}
1176
1177struct WriteContext {
1178 next_inode: u64,
1179 inodes: Vec<PendingInode>,
1180 dir_nodes: Vec<DirNode>,
1181 drops: Vec<PendingDrop>,
1182 drop_index: HashSet<[u8; 32]>,
1183 pending_files: Vec<PendingFile>,
1184 file_count: usize,
1185 dir_count: usize,
1186 root_inode_number: u64,
1187 chunker: FastCDC,
1188 classifier: classifier::Classifier,
1189 shared_inline_map: HashMap<[u8; 32], usize>,
1190 shared_inline_table: Vec<Vec<u8>>,
1191 profile_name: Option<String>,
1193 metadata_codec: u8,
1196 categorizers_disabled: bool,
1198 rw_mode: bool,
1200 auto_turnover: bool,
1202 collect_dict_samples: bool,
1205 dict_samples_by_class: HashMap<crate::classifier::Class, Vec<Vec<u8>>>,
1212 trained_dicts_by_class: HashMap<crate::classifier::Class, crate::dictionary::TrainedDictionary>,
1217 base_drop_index: Option<HashSet<[u8; 32]>>,
1223 base_root: Option<[u8; 32]>,
1228 metadata_externalize_threshold: usize,
1233 emit_shared_inline: bool,
1239 inline_threshold: usize,
1244 pending_sink: Option<std::sync::mpsc::SyncSender<PendingFile>>,
1250}
1251
1252impl WriteContext {
1253 const MAX_DICT_SAMPLES: usize = 1000;
1256
1257 fn new() -> Self {
1258 Self {
1259 next_inode: 1,
1260 inodes: Vec::new(),
1261 dir_nodes: Vec::new(),
1262 drops: Vec::new(),
1263 drop_index: HashSet::new(),
1264 pending_files: Vec::new(),
1265 file_count: 0,
1266 dir_count: 0,
1267 root_inode_number: 0,
1268 chunker: FastCDC::default(),
1269 classifier: classifier::Classifier,
1270 shared_inline_map: HashMap::new(),
1271 shared_inline_table: Vec::new(),
1272 profile_name: None,
1273 metadata_codec: limnifs_core::codec::CODEC_BROTLI,
1274 categorizers_disabled: false,
1275 rw_mode: false,
1276 auto_turnover: false,
1277 collect_dict_samples: false,
1278 dict_samples_by_class: HashMap::new(),
1279 trained_dicts_by_class: HashMap::new(),
1280 base_drop_index: None,
1281 base_root: None,
1282 pending_sink: None,
1283 inline_threshold: INLINE_THRESHOLD,
1284 metadata_externalize_threshold: METADATA_EXTERNALIZE_THRESHOLD,
1285 emit_shared_inline: true,
1286 }
1287 }
1288
1289 fn alloc_inode(&mut self) -> u64 {
1290 let n = self.next_inode;
1291 self.next_inode += 1;
1292 n
1293 }
1294
1295 fn build_shared_inline_table(&mut self) {
1299 let mut counts: HashMap<[u8; 32], usize> = HashMap::new();
1300 for inode in &self.inodes {
1301 if let PendingContent::Inline(data) = &inode.content {
1302 let h = hash_section(data);
1303 *counts.entry(h).or_default() += 1;
1304 }
1305 }
1306 for inode in &self.inodes {
1308 if let PendingContent::Inline(data) = &inode.content {
1309 let h = hash_section(data);
1310 if counts.get(&h).copied().unwrap_or(0) > 1
1311 && !self.shared_inline_map.contains_key(&h)
1312 {
1313 let idx = self.shared_inline_table.len();
1314 self.shared_inline_table.push(data.clone());
1315 self.shared_inline_map.insert(h, idx);
1316 }
1317 }
1318 }
1319 }
1320
1321 fn merge_chunked_file(&mut self, pf: &PendingFile, result: ChunkedFileResult) {
1324 for (drop_id, plaintext, compressed, codec) in result.drops {
1325 if self.drop_index.insert(drop_id) {
1326 let retain_plaintext =
1332 self.collect_dict_samples && codec == limnifs_core::codec::CODEC_ZSTD;
1333 if retain_plaintext {
1334 let total: usize = self.dict_samples_by_class.values().map(Vec::len).sum();
1335 if total < Self::MAX_DICT_SAMPLES {
1336 let class = self.classifier.classify(&plaintext);
1337 self.dict_samples_by_class
1338 .entry(class)
1339 .or_default()
1340 .push(plaintext.clone());
1341 }
1342 }
1343 self.drops.push(PendingDrop {
1344 id: drop_id,
1345 plaintext_len: u32::try_from(plaintext.len()).unwrap_or(u32::MAX),
1346 compressed,
1347 codec,
1348 dict_id: limnifs_core::drop_record::NO_DICT,
1349 plaintext: if retain_plaintext {
1350 Some(plaintext)
1351 } else {
1352 None
1353 },
1354 });
1355 }
1356 }
1357 self.inodes.push(PendingInode {
1358 number: pf.inode_number,
1359 mode: 0o100_644,
1360 mtime_ns: pf.mtime_ns,
1361 content: PendingContent::DropBacked {
1362 file_len: pf.file_len,
1363 slices: result.slices,
1364 },
1365 });
1366 }
1367
1368 #[allow(dead_code)]
1379 fn deepen_drop(&self, drop_id: [u8; 32], plaintext: &[u8]) -> PendingDrop {
1380 let class = self.classifier.classify(plaintext);
1381 let (codec, compressed): (u8, std::sync::Arc<[u8]>) = match class {
1382 classifier::Class::Text | classifier::Class::Code | classifier::Class::Binary => {
1383 let c = limnifs_core::codec::compress_lz4_with_size(plaintext);
1384 (limnifs_core::codec::CODEC_LZ4, c.into())
1385 }
1386 _ => (limnifs_core::codec::CODEC_STORE, plaintext.to_vec().into()),
1387 };
1388 PendingDrop {
1389 id: drop_id,
1390 plaintext_len: u32::try_from(plaintext.len()).unwrap_or(u32::MAX),
1391 compressed,
1392 codec,
1393 dict_id: limnifs_core::drop_record::NO_DICT,
1394 plaintext: None,
1395 }
1396 }
1397
1398 fn walk(&mut self, path: &Path) -> Result<u64, WriteError> {
1399 let meta = std::fs::symlink_metadata(path)?;
1400 let file_type = meta.file_type();
1401 let mtime_ns = meta
1402 .modified()
1403 .ok()
1404 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1405 .map_or(0u128, |d| d.as_nanos());
1406 let mtime_ns: u64 = mtime_ns.try_into().unwrap_or(0);
1407
1408 if file_type.is_dir() {
1409 self.dir_count += 1;
1410 let inode_number = self.alloc_inode();
1411 let mut entries: Vec<(String, u64, u8)> = Vec::new();
1412
1413 for entry in std::fs::read_dir(path)? {
1414 let entry = entry?;
1415 let name = entry.file_name().to_string_lossy().into_owned();
1416 let child_path = entry.path();
1417 let child_inode = self.walk(&child_path)?;
1418 let ft = entry.file_type()?;
1421 let entry_type = if ft.is_symlink() {
1422 0x03
1423 } else if ft.is_dir() {
1424 0x02
1425 } else {
1426 0x01
1427 };
1428 entries.push((name, child_inode, entry_type));
1429 }
1430
1431 entries.sort_by(|a, b| a.0.cmp(&b.0));
1432 let dir_node = encode_dir_node(&entries);
1433 self.dir_nodes.push(dir_node);
1434 self.inodes.push(PendingInode {
1435 number: inode_number,
1436 mode: 0o040_755,
1437 mtime_ns,
1438 content: PendingContent::Directory(entries),
1439 });
1440 Ok(inode_number)
1441 } else if file_type.is_file() {
1442 self.file_count += 1;
1443 let inode_number = self.alloc_inode();
1444 let file_len = meta.len();
1445
1446 if file_len <= u64::try_from(self.inline_threshold).unwrap_or(u64::MAX) {
1447 let data = std::fs::read(path)?;
1448 self.inodes.push(PendingInode {
1449 number: inode_number,
1450 mode: 0o100_644,
1451 mtime_ns,
1452 content: PendingContent::Inline(data),
1453 });
1454 } else {
1455 let pf = PendingFile {
1457 inode_number,
1458 path: path.to_path_buf(),
1459 mtime_ns,
1460 file_len,
1461 };
1462 if let Some(sink) = &self.pending_sink {
1463 sink.send(pf).map_err(|_| {
1469 WriteError::Io(std::io::Error::other("walk: compress pipeline shut down"))
1470 })?;
1471 } else {
1472 self.pending_files.push(pf);
1473 }
1474 }
1475 Ok(inode_number)
1476 } else if file_type.is_symlink() {
1477 let inode_number = self.alloc_inode();
1485 let target = std::fs::read_link(path)?;
1486 let target = target
1487 .to_str()
1488 .ok_or_else(|| WriteError::UnsupportedFileType {
1489 path: path.to_path_buf(),
1490 kind: format!("symlink with non-UTF-8 target ({})", target.display()),
1491 })?
1492 .to_owned();
1493 self.inodes.push(PendingInode {
1494 number: inode_number,
1495 mode: limnifs_core::inode::S_IFLNK | 0o777,
1496 mtime_ns,
1497 content: PendingContent::Symlink(target),
1498 });
1499 Ok(inode_number)
1500 } else {
1501 #[cfg(unix)]
1502 let kind = {
1503 use std::os::unix::fs::FileTypeExt;
1504 if file_type.is_fifo() {
1505 "fifo".to_owned()
1506 } else if file_type.is_socket() {
1507 "socket".to_owned()
1508 } else if file_type.is_block_device() {
1509 "block device".to_owned()
1510 } else if file_type.is_char_device() {
1511 "character device".to_owned()
1512 } else {
1513 "unknown".to_owned()
1514 }
1515 };
1516 #[cfg(not(unix))]
1517 let kind = "unknown".to_owned();
1518 Err(WriteError::UnsupportedFileType {
1519 path: path.to_path_buf(),
1520 kind,
1521 })
1522 }
1523 }
1524
1525 fn train_and_apply_dictionary(&mut self, dictionaries: &crate::config::DictionaryConfig) {
1540 let cleanup = |ctx: &mut Self| {
1541 for d in &mut ctx.drops {
1542 d.plaintext = None;
1543 }
1544 ctx.dict_samples_by_class.clear();
1545 };
1546
1547 if !dictionaries.enabled {
1548 cleanup(self);
1549 return;
1550 }
1551
1552 let target = usize::try_from(dictionaries.max_dict_size).unwrap_or(65_536);
1553 let min_class = usize::try_from(dictionaries.min_class_size).unwrap_or(0);
1554
1555 let text_classes = [
1559 crate::classifier::Class::Text,
1560 crate::classifier::Class::Code,
1561 crate::classifier::Class::Sparse,
1562 ];
1563 let binary_classes = [crate::classifier::Class::Binary];
1564
1565 let text_samples: Vec<&[u8]> = text_classes
1567 .iter()
1568 .flat_map(|c| self.dict_samples_by_class.get(c).into_iter().flatten())
1569 .map(Vec::as_slice)
1570 .collect();
1571 let trainer = crate::dictionary::TrainerKind::from_config_str(&dictionaries.trainer);
1572 if text_samples.len() >= min_class {
1573 if let Some(dict) =
1574 crate::dictionary::train_zstd_with_trainer(0, &text_samples, target, trainer)
1575 {
1576 self.trained_dicts_by_class
1577 .insert(crate::classifier::Class::Text, dict);
1578 }
1579 }
1580 let binary_samples: Vec<&[u8]> = binary_classes
1581 .iter()
1582 .flat_map(|c| self.dict_samples_by_class.get(c).into_iter().flatten())
1583 .map(Vec::as_slice)
1584 .collect();
1585 if binary_samples.len() >= min_class {
1586 if let Some(dict) =
1587 crate::dictionary::train_zstd_with_trainer(1, &binary_samples, target, trainer)
1588 {
1589 self.trained_dicts_by_class
1590 .insert(crate::classifier::Class::Binary, dict);
1591 }
1592 }
1593
1594 for d in self.drops.iter_mut() {
1597 if d.codec != limnifs_core::codec::CODEC_ZSTD {
1598 continue;
1599 }
1600 let Some(plaintext) = d.plaintext.clone() else {
1601 continue;
1602 };
1603 let class = self.classifier.classify(&plaintext);
1604 let dict_class = if text_classes.contains(&class) {
1605 crate::classifier::Class::Text
1606 } else if binary_classes.contains(&class) {
1607 crate::classifier::Class::Binary
1608 } else {
1609 continue;
1610 };
1611 let Some(dict) = self.trained_dicts_by_class.get(&dict_class) else {
1612 continue;
1613 };
1614 let Ok(dict_compressed) = dict.compress(&plaintext) else {
1615 continue;
1616 };
1617 if dict_compressed.len() < d.compressed.len() {
1618 d.compressed = dict_compressed.into();
1619 d.dict_id = dict.id;
1620 }
1621 }
1622
1623 cleanup(self);
1624 }
1625
1626 fn trace_phase(label: &str, start: std::time::Instant) {
1628 if std::env::var_os("LIMNIFS_TRACE_ASSEMBLE").is_some() {
1629 eprintln!("[assemble] {label}: {:?}", start.elapsed());
1630 }
1631 }
1632
1633 fn assemble(mut self) -> WriteArtifact {
1634 let t_assemble = std::time::Instant::now();
1635 let inode_count = self.inodes.len();
1636 let dir_count = self.dir_count;
1637 let drop_count = self.drops.len();
1638
1639 let t = std::time::Instant::now();
1645 let slabs = pack_slabs(&self.drops);
1646 Self::trace_phase("pack_slabs", t);
1647
1648 let t = std::time::Instant::now();
1652 if self.emit_shared_inline {
1653 self.build_shared_inline_table();
1654 }
1655 Self::trace_phase("shared_inline_table", t);
1656
1657 let mut metadata_blob = Vec::new();
1658 metadata_blob.extend_from_slice(&u32::try_from(self.inodes.len()).unwrap().to_le_bytes());
1659 for inode in &self.inodes {
1660 self.encode_inode(&mut metadata_blob, inode);
1661 }
1662 metadata_blob
1663 .extend_from_slice(&u32::try_from(self.dir_nodes.len()).unwrap().to_le_bytes());
1664 for node in &self.dir_nodes {
1665 metadata_blob.extend_from_slice(&node.bytes);
1666 }
1667 if !self.shared_inline_table.is_empty() {
1670 metadata_blob.extend_from_slice(
1671 &u32::try_from(self.shared_inline_table.len())
1672 .unwrap()
1673 .to_le_bytes(),
1674 );
1675 for entry in &self.shared_inline_table {
1676 let len = u32::try_from(entry.len()).expect("shared entry fits u32");
1677 metadata_blob.extend_from_slice(&len.to_le_bytes());
1678 metadata_blob.extend_from_slice(entry);
1679 }
1680 }
1681
1682 Self::trace_phase("metadata_encode", t);
1683 let uncompressed_len =
1690 u32::try_from(metadata_blob.len()).expect("metadata blob length fits u32");
1691 let t = std::time::Instant::now();
1692 let metadata_hash = hash_section(&metadata_blob);
1693 let metadata_codec = self.metadata_codec;
1694 let metadata_quality = if metadata_blob.len() > METADATA_LARGE_BLOB_THRESHOLD {
1695 METADATA_LARGE_BLOB_QUALITY
1696 } else {
1697 METADATA_SMALL_BLOB_QUALITY
1698 };
1699 let compressed_blob = if metadata_codec == limnifs_core::codec::CODEC_BROTLI {
1700 limnifs_core::codec::compress_brotli_with_quality(&metadata_blob, metadata_quality)
1701 .unwrap_or_else(|_| metadata_blob.clone())
1702 } else {
1703 limnifs_core::codec::compress(metadata_codec, &metadata_blob)
1704 .unwrap_or_else(|_| metadata_blob.clone())
1705 };
1706 Self::trace_phase("metadata_compress", t);
1707 let (on_wire_codec, on_wire_blob) = if compressed_blob.len() < metadata_blob.len() {
1708 (metadata_codec, compressed_blob)
1709 } else {
1710 (limnifs_core::codec::CODEC_STORE, metadata_blob.clone())
1711 };
1712
1713 let externalize_at = self
1717 .metadata_externalize_threshold
1718 .min(limnifs_core::metadata_reference::DEFAULT_INLINE_METADATA_MAX_BYTES as usize);
1719 let (metadata_sidecar, inline_data, metadata_locator_count) =
1720 if on_wire_blob.len() > externalize_at {
1721 let locator = "file:metadata.bin".to_owned();
1722 let sidecar = MetadataSidecar {
1723 bytes: on_wire_blob.clone(),
1724 locator,
1725 };
1726 (Some(sidecar), None, 1u32)
1727 } else {
1728 (None, Some(on_wire_blob.clone()), 0u32)
1729 };
1730
1731 let mut manifest = Vec::new();
1732
1733 let header_start = manifest.len();
1734 manifest.extend_from_slice(&ManifestHeader::current().to_bytes());
1735 let header_end = manifest.len();
1736
1737 let flags_start = manifest.len();
1738 manifest.push(FEATURE_FLAGS_SECTION_VERSION);
1739 manifest.extend_from_slice(&0u32.to_le_bytes());
1740 let flags_end = manifest.len();
1741
1742 let meta_ref_start = manifest.len();
1745 manifest.push(METADATA_REFERENCE_SECTION_VERSION_2);
1746 manifest.extend_from_slice(&metadata_hash);
1747 manifest.extend_from_slice(&uncompressed_len.to_le_bytes());
1748 manifest.push(on_wire_codec);
1749 manifest.extend_from_slice(&metadata_locator_count.to_le_bytes());
1750 if let Some(sidecar) = &metadata_sidecar {
1751 let loc_bytes = sidecar.locator.as_bytes();
1752 let loc_len = u32::try_from(loc_bytes.len()).expect("locator fits u32");
1753 manifest.extend_from_slice(&loc_len.to_le_bytes());
1754 manifest.extend_from_slice(loc_bytes);
1755 }
1756 match &inline_data {
1757 Some(blob) => {
1758 let inline_len = u32::try_from(blob.len()).expect("metadata fits u32");
1759 manifest.extend_from_slice(&inline_len.to_le_bytes());
1760 manifest.extend_from_slice(blob);
1761 }
1762 None => {
1763 manifest.extend_from_slice(&0u32.to_le_bytes());
1764 }
1765 }
1766 let meta_ref_end = manifest.len();
1767
1768 let slab_index_start = manifest.len();
1769 manifest.push(SLAB_INDEX_SECTION_VERSION);
1770 manifest.extend_from_slice(&u32::try_from(slabs.len()).unwrap().to_le_bytes());
1771 for slab in &slabs {
1772 manifest.extend_from_slice(&slab.id.to_bytes());
1773 manifest.extend_from_slice(&1u32.to_le_bytes());
1774 let loc_bytes = slab.locator.as_bytes();
1775 let loc_len = u32::try_from(loc_bytes.len()).expect("locator fits u32");
1776 manifest.extend_from_slice(&loc_len.to_le_bytes());
1777 manifest.extend_from_slice(loc_bytes);
1778 }
1779 let slab_index_end = manifest.len();
1780
1781 let history_start = manifest.len();
1782 manifest.push(HISTORY_SECTION_VERSION);
1783 manifest.extend_from_slice(&1u32.to_le_bytes());
1784 manifest.push(0x01);
1785 manifest.extend_from_slice(&0u64.to_le_bytes());
1786 manifest.extend_from_slice(&0u32.to_le_bytes());
1787 manifest.extend_from_slice(&0u32.to_le_bytes());
1788 let history_end = manifest.len();
1789
1790 let profile_desc_start = manifest.len();
1795 if let Some(ref name) = self.profile_name {
1796 let desc = limnifs_core::profile_descriptor::ProfileDescriptor {
1797 version: limnifs_core::profile_descriptor::PROFILE_DESCRIPTOR_SECTION_VERSION,
1798 profile_name: Some(name.clone()),
1799 blake3_hashing: true,
1800 cross_file_dedup: true,
1801 content_classification: !self.categorizers_disabled,
1802 integrity_verify: true,
1803 read_write: self.rw_mode,
1804 auto_turnover: self.auto_turnover,
1805 };
1806 limnifs_core::profile_descriptor::encode_profile_descriptor(&desc, &mut manifest);
1807 }
1808 let profile_desc_end = manifest.len();
1809
1810 if !self.trained_dicts_by_class.is_empty() {
1816 let dicts: Vec<_> = self
1817 .trained_dicts_by_class
1818 .values()
1819 .map(|d| limnifs_core::dictionary_section::Dictionary {
1820 codec_id: d.codec,
1821 class_id: d.id,
1822 data: d.content.clone(),
1823 })
1824 .collect();
1825 let section = limnifs_core::dictionary_section::DictionarySection {
1826 version: limnifs_core::dictionary_section::DICTIONARY_SECTION_VERSION,
1827 dicts,
1828 };
1829 limnifs_core::dictionary_section::encode_dictionary_section(§ion, &mut manifest);
1830 }
1831
1832 let dictionary_end = manifest.len();
1833
1834 let delta_linkage_hash = if let Some(base_root) = self.base_root {
1840 let delta_start = manifest.len();
1841 manifest.push(limnifs_core::delta_linkage::DELTA_LINKAGE_SECTION_VERSION);
1847 manifest.extend_from_slice(&base_root);
1848 manifest.extend_from_slice(&0u32.to_le_bytes());
1849 hash_section(&manifest[delta_start..])
1850 } else {
1851 hash_empty_section()
1852 };
1853 let _ = dictionary_end;
1854
1855 let hashes = SectionHashes {
1856 metadata: metadata_hash,
1857 format_header: hash_section(&manifest[header_start..header_end]),
1858 feature_flags: hash_section(&manifest[flags_start..flags_end]),
1859 metadata_reference: hash_section(&manifest[meta_ref_start..meta_ref_end]),
1860 slab_index: hash_section(&manifest[slab_index_start..slab_index_end]),
1861 crypto_params: hash_empty_section(),
1862 ec_params: hash_empty_section(),
1863 dms_policy: hash_empty_section(),
1864 delta_linkage: delta_linkage_hash,
1865 history: hash_section(&manifest[history_start..history_end]),
1866 };
1878 let merkle_root = compute_merkle_root(&hashes);
1879
1880 WriteArtifact {
1881 bytes: manifest,
1882 merkle_root,
1883 slabs,
1884 metadata_sidecar,
1885 inode_count,
1886 file_count: self.file_count,
1887 dir_count,
1888 drop_count,
1889 root_inode_number: self.root_inode_number,
1890 }
1891 }
1892
1893 fn encode_inode(&self, out: &mut Vec<u8>, inode: &PendingInode) {
1894 out.extend_from_slice(&inode.number.to_le_bytes());
1895 out.extend_from_slice(&inode.mode.to_le_bytes());
1896 out.extend_from_slice(&0u32.to_le_bytes());
1897 out.extend_from_slice(&0u32.to_le_bytes());
1898 out.extend_from_slice(&inode.mtime_ns.to_le_bytes());
1899 out.extend_from_slice(&inode.mtime_ns.to_le_bytes());
1900 out.extend_from_slice(&1u32.to_le_bytes());
1901 match &inode.content {
1902 PendingContent::Inline(data) => {
1903 let h = hash_section(data);
1904 if let Some(&idx) = self.shared_inline_map.get(&h) {
1905 out.push(INODE_FLAG_INLINE_DATA | INODE_FLAG_SHARED_INLINE);
1907 out.extend_from_slice(&(idx as u32).to_le_bytes());
1908 } else {
1909 out.push(INODE_FLAG_INLINE_DATA);
1910 let len = u32::try_from(data.len()).expect("data fits u32");
1911 out.extend_from_slice(&len.to_le_bytes());
1912 out.extend_from_slice(data);
1913 }
1914 }
1915 PendingContent::DropBacked { file_len, slices } => {
1916 out.push(0x00);
1917 let slice_count = u32::try_from(slices.len()).expect("slice count fits u32");
1918 out.extend_from_slice(&slice_count.to_le_bytes());
1919 for slice in slices {
1920 out.extend_from_slice(&slice.file_byte_start.to_le_bytes());
1921 out.extend_from_slice(&slice.file_byte_end.to_le_bytes());
1922 out.extend_from_slice(&slice.drop_id);
1923 out.extend_from_slice(&0u32.to_le_bytes());
1925 let drop_byte_len = u32::try_from(slice.file_byte_end - slice.file_byte_start)
1929 .expect("slice range fits u32");
1930 out.extend_from_slice(&drop_byte_len.to_le_bytes());
1931 }
1932 let _ = file_len;
1933 }
1934 PendingContent::Symlink(target) => {
1935 out.push(0x00);
1939 let t = target.as_bytes();
1940 let len = u32::try_from(t.len()).expect("target fits u32");
1941 out.extend_from_slice(&len.to_le_bytes());
1942 out.extend_from_slice(t);
1943 }
1944 PendingContent::Directory(entries) => {
1945 out.push(0x00);
1946 let node = self
1947 .dir_nodes
1948 .iter()
1949 .find(|n| n.entries == *entries)
1950 .expect("directory node must exist");
1951 out.extend_from_slice(&node.hash);
1952 }
1953 }
1954 }
1955}
1956
1957fn sidecar_name(locator: &str) -> Result<&str, WriteError> {
1963 limnifs_core::locator::local_sidecar_name(locator)
1964 .map_err(|e| WriteError::Io(std::io::Error::other(format!("{e}"))))
1965}
1966
1967fn encode_dir_node(entries: &[(String, u64, u8)]) -> DirNode {
1968 let mut bytes = Vec::new();
1969 bytes.push(1u8);
1970 let count = u32::try_from(entries.len()).expect("entry count fits u32");
1971 bytes.extend_from_slice(&count.to_le_bytes());
1972 for (name, inode_number, entry_type) in entries {
1973 let name_bytes = name.as_bytes();
1974 let name_len = u32::try_from(name_bytes.len()).expect("name fits u32");
1975 bytes.extend_from_slice(&name_len.to_le_bytes());
1976 bytes.extend_from_slice(name_bytes);
1977 bytes.extend_from_slice(&inode_number.to_le_bytes());
1978 bytes.push(*entry_type);
1979 }
1980 let hash = hash_section(&bytes);
1981 DirNode {
1982 entries: entries.to_vec(),
1983 bytes,
1984 hash,
1985 }
1986}
1987
1988fn pack_slabs(drops: &[PendingDrop]) -> Vec<SlabArtifact> {
1997 let local_drops: Vec<&PendingDrop> = drops
2002 .iter()
2003 .filter(|d| d.codec != limnifs_core::codec::CODEC_REFERENCED)
2004 .collect();
2005 if local_drops.is_empty() {
2006 return Vec::new();
2007 }
2008
2009 let max_content = MAX_SLAB_TOTAL_BYTES.saturating_sub(SLAB_HEADER_LEN);
2010
2011 let mut slab_groups: Vec<Vec<&PendingDrop>> = Vec::new();
2016 let mut current: Vec<&PendingDrop> = Vec::new();
2017 let mut current_size: usize = 0;
2018
2019 for drop in &local_drops {
2020 let footprint = drop.slab_footprint();
2021 if !current.is_empty() && current_size + footprint > max_content {
2022 slab_groups.push(std::mem::take(&mut current));
2023 current_size = 0;
2024 }
2025 current.push(*drop);
2026 current_size += footprint;
2027 }
2028 if !current.is_empty() {
2029 slab_groups.push(current);
2030 }
2031
2032 use rayon::prelude::*;
2038 slab_groups
2039 .par_iter()
2040 .enumerate()
2041 .map(|(ordinal, group)| {
2042 let ordinal_u64 = u64::try_from(ordinal).expect("slab count fits u64");
2043 encode_slab(ordinal_u64, group)
2044 })
2045 .collect()
2046}
2047
2048fn encode_slab(ordinal: u64, drops: &[&PendingDrop]) -> SlabArtifact {
2052 const DROP_RECORD_LEN: usize = 49;
2058 let mut drop_records = Vec::with_capacity(drops.len() * DROP_RECORD_LEN);
2059 let mut solid_window = Vec::new();
2060 let mut drop_ids = Vec::with_capacity(drops.len());
2061 let mut offset_in_window: u32 = 0;
2062
2063 for drop in drops {
2064 let plaintext_len = drop.plaintext_len_value();
2065 let window_len = drop.len_in_window();
2066 drop_records.extend_from_slice(&drop.id);
2067 drop_records.extend_from_slice(&plaintext_len.to_le_bytes());
2068 drop_records.extend_from_slice(&[drop.codec, 0x00, 0x00]);
2070 drop_records.push(0x00); drop_records.extend_from_slice(&offset_in_window.to_le_bytes());
2072 drop_records.extend_from_slice(&window_len.to_le_bytes());
2073 drop_records.push(drop.dict_id); solid_window.extend_from_slice(&drop.compressed);
2075 drop_ids.push(drop.id);
2076 offset_in_window = offset_in_window
2077 .checked_add(window_len)
2078 .expect("slab window size fits u32");
2079 }
2080
2081 let slab_content = [&drop_records[..], &solid_window[..]].concat();
2082 let slab_hash = hash_section(&slab_content);
2083 let slab_id = SlabId::new(ordinal, slab_hash);
2084
2085 let total_length = SLAB_HEADER_LEN + slab_content.len();
2086 let mut slab_bytes = Vec::with_capacity(total_length);
2087 slab_bytes.extend_from_slice(b"LIM1");
2088 slab_bytes.extend_from_slice(&1u16.to_le_bytes());
2089 slab_bytes.extend_from_slice(&slab_id.to_bytes());
2090 slab_bytes.extend_from_slice(
2091 &u64::try_from(total_length)
2092 .unwrap_or(u64::MAX)
2093 .to_le_bytes(),
2094 );
2095 slab_bytes.push(0x00);
2096 slab_bytes.push(0x00);
2097 slab_bytes.extend_from_slice(&slab_content);
2098
2099 let locator = format!("file:slab-{ordinal}.bin");
2100
2101 SlabArtifact {
2102 id: slab_id,
2103 bytes: slab_bytes,
2104 locator,
2105 drop_ids,
2106 }
2107}
2108
2109#[cfg(test)]
2110fn pseudo_random_bytes(seed: u64, count: usize) -> Vec<u8> {
2111 let mut state = seed;
2112 let mut out = Vec::with_capacity(count);
2113 for _ in 0..count {
2114 state = state
2115 .wrapping_mul(6_364_136_223_846_793_005)
2116 .wrapping_add(1_442_695_040_888_963_407);
2117 out.push(u8::try_from(state >> 56).expect("fits u8"));
2118 }
2119 out
2120}
2121
2122#[cfg(test)]
2123mod tests {
2124 use super::*;
2125 use limnifs_core::ManifestCursor;
2126
2127 #[test]
2128 fn write_stream_packs_single_named_stream() {
2129 let temp = std::env::temp_dir().join(format!(
2133 "limnifs-write-stream-test-{}-{}",
2134 std::process::id(),
2135 std::time::SystemTime::now()
2136 .duration_since(std::time::UNIX_EPOCH)
2137 .unwrap()
2138 .as_nanos()
2139 ));
2140 std::fs::create_dir_all(&temp).expect("create temp dir");
2141
2142 let content = b"stream test content line\n".repeat(10_000); let cursor = std::io::Cursor::new(content.clone());
2144 let config = WriteConfig::default_v0_1();
2145 let artifact = write_stream("streamed.txt", cursor, &config).expect("write_stream");
2146
2147 assert!(!artifact.bytes.is_empty(), "manifest bytes non-empty");
2148 assert!(!artifact.slabs.is_empty(), "at least one slab produced");
2149 let total_drop_bytes: usize = artifact.slabs.iter().map(|s| s.bytes.len()).sum();
2154 assert!(total_drop_bytes > 0, "drops non-empty");
2155
2156 let _ = std::fs::remove_dir_all(&temp);
2157 }
2158
2159 #[test]
2160 fn write_layer_references_base_drops() {
2161 let temp = std::env::temp_dir().join(format!(
2167 "limnifs-write-layer-test-{}-{}",
2168 std::process::id(),
2169 std::time::SystemTime::now()
2170 .duration_since(std::time::UNIX_EPOCH)
2171 .unwrap()
2172 .as_nanos()
2173 ));
2174 std::fs::create_dir_all(&temp).expect("create temp dir");
2175
2176 let base_dir = temp.join("base");
2178 std::fs::create_dir_all(&base_dir).expect("base dir");
2179 let text = b"layer test content line\n".repeat(50_000); std::fs::write(base_dir.join("shared.txt"), &text).expect("write shared");
2181 std::fs::write(base_dir.join("base-only.txt"), b"only in base").expect("write base-only");
2182
2183 let config = WriteConfig::default_v0_1();
2184 let base_artifact = write_directory_with_config(&base_dir, &config).expect("base");
2185
2186 let base_manifest = temp.join("base.lim");
2187 std::fs::write(&base_manifest, &base_artifact.bytes).expect("write base manifest");
2188 for slab in &base_artifact.slabs {
2189 let slab_name = sidecar_name(&slab.locator).expect("slab locator");
2190 std::fs::write(temp.join(slab_name), &slab.bytes).expect("write base slab");
2191 }
2192
2193 let layer_dir = temp.join("layer");
2195 std::fs::create_dir_all(&layer_dir).expect("layer dir");
2196 std::fs::write(layer_dir.join("shared.txt"), &text).expect("write shared in layer");
2197 std::fs::write(layer_dir.join("new.txt"), b"fresh content in layer").expect("write new");
2198
2199 let layer_artifact = write_layer(&base_manifest, &layer_dir, &config).expect("layer");
2200
2201 let layer_slab_bytes: usize = layer_artifact.slabs.iter().map(|s| s.bytes.len()).sum();
2204 let base_slab_bytes: usize = base_artifact.slabs.iter().map(|s| s.bytes.len()).sum();
2205 assert!(
2206 layer_slab_bytes < base_slab_bytes / 4,
2207 "layer slabs ({}) should be much smaller than base ({}) — layering failed",
2208 layer_slab_bytes,
2209 base_slab_bytes
2210 );
2211
2212 let base_root = base_artifact.merkle_root.as_bytes();
2215 assert!(
2216 layer_artifact
2217 .bytes
2218 .windows(32)
2219 .any(|w| w == base_root.as_slice()),
2220 "layer manifest must contain base's ManifestRoot bytes"
2221 );
2222
2223 let _ = std::fs::remove_dir_all(&temp);
2224 }
2225
2226 #[test]
2227 fn tournament_short_circuits_on_highly_compressible_chunk() {
2228 let chunk = b"hello world ".repeat(500);
2231 let tunables = limnifs_core::codec::CodecTunables::default();
2232 let tournament = TournamentSpec {
2233 codec_ids: vec![
2234 limnifs_core::codec::CODEC_LZ4,
2235 limnifs_core::codec::CODEC_BROTLI,
2236 ],
2237 min_size: 16,
2238 skip_for_binary: false,
2239 short_circuit_permille: 250,
2240 };
2241 let (codec_id, compressed) = compress_chunk_with_tournament(
2242 &chunk,
2243 classifier::Class::Text,
2244 limnifs_core::codec::CODEC_BROTLI,
2245 limnifs_core::codec::CODEC_LZ4,
2246 &tunables,
2247 &tournament,
2248 );
2249 assert_eq!(codec_id, limnifs_core::codec::CODEC_LZ4);
2250 assert!(compressed.len() < chunk.len());
2251 }
2252
2253 #[test]
2254 fn tournament_runs_all_codecs_when_short_circuit_disabled() {
2255 let chunk = b"hello world ".repeat(500);
2261 let tunables = limnifs_core::codec::CodecTunables::default();
2262 let tournament = TournamentSpec {
2263 codec_ids: vec![
2264 limnifs_core::codec::CODEC_LZ4,
2265 limnifs_core::codec::CODEC_BROTLI,
2266 limnifs_core::codec::CODEC_ZSTD,
2267 ],
2268 min_size: 16,
2269 skip_for_binary: false,
2270 short_circuit_permille: 0,
2271 };
2272 let (codec_id, compressed) = compress_chunk_with_tournament(
2273 &chunk,
2274 classifier::Class::Text,
2275 limnifs_core::codec::CODEC_BROTLI,
2276 limnifs_core::codec::CODEC_LZ4,
2277 &tunables,
2278 &tournament,
2279 );
2280 assert!(
2284 codec_id == limnifs_core::codec::CODEC_ZSTD
2285 || codec_id == limnifs_core::codec::CODEC_BROTLI,
2286 "expected ZSTD or Brotli to win, got codec {codec_id}"
2287 );
2288 assert!(compressed.len() < chunk.len());
2289 }
2290
2291 #[test]
2292 fn tournament_skips_for_binary_when_configured() {
2293 let chunk = vec![0u8; 4096];
2294 let tunables = limnifs_core::codec::CodecTunables::default();
2295 let tournament = TournamentSpec {
2296 codec_ids: vec![limnifs_core::codec::CODEC_BROTLI],
2297 min_size: 16,
2298 skip_for_binary: true,
2299 short_circuit_permille: 250,
2300 };
2301 let (codec_id, _compressed) = compress_chunk_with_tournament(
2302 &chunk,
2303 classifier::Class::Binary,
2304 limnifs_core::codec::CODEC_BROTLI,
2305 limnifs_core::codec::CODEC_LZ4,
2306 &tunables,
2307 &tournament,
2308 );
2309 assert_eq!(codec_id, limnifs_core::codec::CODEC_LZ4);
2311 }
2312
2313 #[test]
2314 fn tournament_small_chunk_uses_preferred_codec() {
2315 let chunk = b"tiny";
2316 let tunables = limnifs_core::codec::CodecTunables::default();
2317 let tournament = TournamentSpec {
2318 codec_ids: vec![limnifs_core::codec::CODEC_BROTLI],
2319 min_size: 1024,
2320 skip_for_binary: false,
2321 short_circuit_permille: 0,
2322 };
2323 let (codec_id, _compressed) = compress_chunk_with_tournament(
2324 chunk,
2325 classifier::Class::Text,
2326 limnifs_core::codec::CODEC_BROTLI,
2327 limnifs_core::codec::CODEC_LZ4,
2328 &tunables,
2329 &tournament,
2330 );
2331 assert_eq!(codec_id, limnifs_core::codec::CODEC_STORE);
2333 }
2334
2335 #[test]
2336 fn tournament_falls_back_to_store_when_no_codec_compresses() {
2337 let chunk = pseudo_random_bytes(42, 4096);
2341 let tunables = limnifs_core::codec::CodecTunables::default();
2342 let tournament = TournamentSpec {
2343 codec_ids: vec![limnifs_core::codec::CODEC_LZ4],
2344 min_size: 16,
2345 skip_for_binary: false,
2346 short_circuit_permille: 0,
2347 };
2348 let (codec_id, compressed) = compress_chunk_with_tournament(
2349 &chunk,
2350 classifier::Class::Binary,
2351 limnifs_core::codec::CODEC_BROTLI,
2352 limnifs_core::codec::CODEC_LZ4,
2353 &tunables,
2354 &tournament,
2355 );
2356 assert_eq!(codec_id, limnifs_core::codec::CODEC_STORE);
2357 assert_eq!(compressed.len(), chunk.len());
2358 }
2359
2360 #[test]
2361 fn dictionaries_enabled_emits_dictionary_section_when_enough_samples() {
2362 let temp = std::env::temp_dir().join(format!(
2367 "limnifs-write-test-{}-dict-{}",
2368 std::process::id(),
2369 std::time::SystemTime::now()
2370 .duration_since(std::time::UNIX_EPOCH)
2371 .map(|d| d.as_nanos() as u64)
2372 .unwrap_or(0),
2373 ));
2374 let _ = std::fs::remove_dir_all(&temp);
2375 std::fs::create_dir_all(&temp).expect("mkdir");
2376
2377 for i in 0..200 {
2380 let content = format!(
2382 "function test_case_{i}() {{ return constant + {i}; }}\n\
2383 // shared comment line {i}\n\
2384 struct Foo {{ x: i32 }} // type {i}\n"
2385 )
2386 .repeat(5);
2387 let path = temp.join(format!("file_{i:04}.txt"));
2388 std::fs::write(&path, content.as_bytes()).expect("write");
2389 }
2390
2391 let mut config = crate::profile::balanced();
2392 config.defaults.text_codec = "zstd".into();
2394 config.defaults.metadata_codec = "zstd".into();
2398 config.dictionaries.enabled = true;
2399 config.dictionaries.min_class_size = 50;
2400 config.dictionaries.max_dict_size = 8192;
2401
2402 let artifact = write_directory_with_config(&temp, &config).expect("write");
2403 std::fs::remove_dir_all(&temp).ok();
2404
2405 let mut cursor = ManifestCursor::new(&artifact.bytes);
2410 let _ = limnifs_core::parse_manifest_header(&mut cursor).expect("header");
2411 let _ = limnifs_core::parse_feature_flags_section(&mut cursor).expect("flags");
2412 let _ = limnifs_core::parse_metadata_reference(&mut cursor).expect("meta_ref");
2413 let _ = limnifs_core::parse_slab_index(&mut cursor).expect("slab_index");
2414 let _ = limnifs_core::parse_history(&mut cursor).expect("history");
2415 let _remaining = cursor.remaining_len();
2418 }
2419
2420 #[test]
2421 fn write_empty_directory() {
2422 let temp =
2423 std::env::temp_dir().join(format!("limnifs-write-test-{}-empty", std::process::id()));
2424 std::fs::create_dir_all(&temp).expect("create temp dir");
2425 let artifact = write_directory(&temp).expect("write succeeds");
2426 std::fs::remove_dir_all(&temp).ok();
2427 assert!(artifact.inode_count >= 1);
2428 assert_eq!(artifact.file_count, 0);
2429 assert_eq!(artifact.dir_count, 1);
2430 assert!(artifact.slabs.is_empty());
2431 }
2432
2433 #[test]
2434 fn write_small_file_inline() {
2435 let temp =
2436 std::env::temp_dir().join(format!("limnifs-write-test-{}-small", std::process::id()));
2437 std::fs::create_dir_all(&temp).expect("create temp dir");
2438 std::fs::write(temp.join("hello.txt"), b"hello world").expect("write file");
2439 let artifact = write_directory(&temp).expect("write succeeds");
2440 std::fs::remove_dir_all(&temp).ok();
2441 assert_eq!(artifact.file_count, 1);
2442 assert!(artifact.slabs.is_empty());
2443 assert_eq!(artifact.drop_count, 0);
2444 }
2445
2446 #[test]
2447 fn write_large_file_uses_slab() {
2448 let temp =
2449 std::env::temp_dir().join(format!("limnifs-write-test-{}-large", std::process::id()));
2450 std::fs::create_dir_all(&temp).expect("create temp dir");
2451 let large_data = vec![0xABu8; INLINE_THRESHOLD + 100];
2452 std::fs::write(temp.join("big.bin"), &large_data).expect("write big");
2453 let artifact = write_directory(&temp).expect("write succeeds");
2454 std::fs::remove_dir_all(&temp).ok();
2455 assert_eq!(artifact.drop_count, 1);
2456 assert_eq!(artifact.slabs.len(), 1);
2457 }
2458
2459 #[test]
2460 fn write_mixed_inline_and_large() {
2461 let temp =
2462 std::env::temp_dir().join(format!("limnifs-write-test-{}-mix", std::process::id()));
2463 std::fs::create_dir_all(&temp).expect("create temp dir");
2464 std::fs::write(temp.join("small.txt"), b"tiny").expect("write small");
2465 std::fs::write(temp.join("large.bin"), vec![0xCDu8; INLINE_THRESHOLD * 2])
2466 .expect("write large");
2467 let artifact = write_directory(&temp).expect("write succeeds");
2468 std::fs::remove_dir_all(&temp).ok();
2469 assert_eq!(artifact.file_count, 2);
2470 assert_eq!(artifact.drop_count, 1);
2471 assert_eq!(artifact.slabs.len(), 1);
2472 }
2473
2474 #[test]
2475 fn deduplicates_identical_large_files() {
2476 let temp =
2477 std::env::temp_dir().join(format!("limnifs-write-test-{}-dedup", std::process::id()));
2478 std::fs::create_dir_all(&temp).expect("create temp dir");
2479 let data = vec![0x77u8; INLINE_THRESHOLD + 10];
2480 std::fs::write(temp.join("a.bin"), &data).expect("write a");
2481 std::fs::write(temp.join("b.bin"), &data).expect("write b");
2482 let artifact = write_directory(&temp).expect("write succeeds");
2483 std::fs::remove_dir_all(&temp).ok();
2484 assert_eq!(artifact.drop_count, 1);
2485 }
2486
2487 #[test]
2488 fn write_and_verify_roundtrip() {
2489 let temp = std::env::temp_dir().join(format!(
2490 "limnifs-write-test-{}-roundtrip",
2491 std::process::id()
2492 ));
2493 std::fs::create_dir_all(&temp).expect("create temp dir");
2494 std::fs::write(temp.join("a.txt"), b"aaa").expect("write a");
2495 std::fs::write(temp.join("b.txt"), b"bbb").expect("write b");
2496 std::fs::create_dir_all(temp.join("sub")).expect("create sub");
2497 std::fs::write(temp.join("sub").join("c.txt"), b"ccc").expect("write c");
2498 let artifact = write_directory(&temp).expect("write succeeds");
2499 std::fs::remove_dir_all(&temp).ok();
2500 assert_eq!(artifact.file_count, 3);
2501 assert_eq!(artifact.dir_count, 2);
2502
2503 let mut cursor = ManifestCursor::new(&artifact.bytes);
2504 limnifs_core::parse_manifest_header(&mut cursor).expect("header");
2505 limnifs_core::parse_feature_flags_section(&mut cursor).expect("flags");
2506 let meta_ref = limnifs_core::parse_metadata_reference(&mut cursor).expect("meta ref");
2507 assert!(meta_ref.is_inlined());
2508 let slab_index = limnifs_core::parse_slab_index(&mut cursor).expect("slab index");
2509 assert_eq!(slab_index.len(), 0);
2510 limnifs_core::parse_history(&mut cursor).expect("history");
2511 }
2512
2513 #[test]
2514 fn write_deterministic() {
2515 let temp =
2516 std::env::temp_dir().join(format!("limnifs-write-test-{}-det", std::process::id()));
2517 std::fs::create_dir_all(&temp).expect("create temp dir");
2518 std::fs::write(temp.join("x.txt"), b"xxx").expect("write x");
2519
2520 let a1 = write_directory(&temp).expect("first write");
2521 let a2 = write_directory(&temp).expect("second write");
2522 std::fs::remove_dir_all(&temp).ok();
2523
2524 assert_eq!(a1.bytes, a2.bytes);
2525 assert_eq!(a1.merkle_root, a2.merkle_root);
2526 }
2527
2528 #[test]
2529 fn slab_parses_correctly() {
2530 let temp =
2531 std::env::temp_dir().join(format!("limnifs-write-test-{}-slab", std::process::id()));
2532 std::fs::create_dir_all(&temp).expect("create temp dir");
2533 std::fs::write(temp.join("big.bin"), vec![0x11u8; INLINE_THRESHOLD + 1])
2534 .expect("write big");
2535 let artifact = write_directory(&temp).expect("write succeeds");
2536 std::fs::remove_dir_all(&temp).ok();
2537
2538 let slab_bytes = &artifact.slabs[0].bytes;
2539 let mut cursor = ManifestCursor::new(slab_bytes);
2540 let slab_header = limnifs_core::parse_slab_header(&mut cursor).expect("slab header parses");
2541 assert_eq!(slab_header.format_version, 1);
2542 assert!(!slab_header.is_sealed());
2543 assert!(!slab_header.has_erasure_coding());
2544
2545 let drop_record =
2546 limnifs_core::parse_drop_record(&mut cursor, &slab_header).expect("drop record parses");
2547 assert_eq!(drop_record.plaintext_len as usize, INLINE_THRESHOLD + 1);
2548 }
2549
2550 #[test]
2551 fn fastcdc_produces_multiple_chunks_for_large_files() {
2552 let temp = std::env::temp_dir().join(format!(
2555 "limnifs-write-test-{}-cdc-multi",
2556 std::process::id()
2557 ));
2558 std::fs::create_dir_all(&temp).expect("create temp dir");
2559 let data = pseudo_random_bytes(42, 1024 * 1024);
2560 std::fs::write(temp.join("big.bin"), &data).expect("write big");
2561 let artifact = write_directory(&temp).expect("write succeeds");
2562 std::fs::remove_dir_all(&temp).ok();
2563 assert!(
2564 artifact.drop_count > 1,
2565 "expected FastCDC to produce multiple drops for 1 MiB input, got {}",
2566 artifact.drop_count
2567 );
2568 }
2569
2570 #[test]
2571 fn fastcdc_deduplicates_shared_substrings() {
2572 let temp = std::env::temp_dir().join(format!(
2576 "limnifs-write-test-{}-cdc-dedup",
2577 std::process::id()
2578 ));
2579 std::fs::create_dir_all(&temp).expect("create temp dir");
2580 let shared = pseudo_random_bytes(7, 512 * 1024);
2581 let mut a = Vec::with_capacity(shared.len() + 1024);
2582 a.extend_from_slice(&pseudo_random_bytes(1, 1024));
2583 a.extend_from_slice(&shared);
2584 let mut b = Vec::with_capacity(shared.len() + 2048);
2585 b.extend_from_slice(&pseudo_random_bytes(2, 2048));
2586 b.extend_from_slice(&shared);
2587 std::fs::write(temp.join("a.bin"), &a).expect("write a");
2588 std::fs::write(temp.join("b.bin"), &b).expect("write b");
2589
2590 let temp_a = std::env::temp_dir().join(format!(
2592 "limnifs-write-test-{}-cdc-dedup-a",
2593 std::process::id()
2594 ));
2595 std::fs::create_dir_all(&temp_a).expect("create temp_a");
2596 std::fs::write(temp_a.join("a.bin"), &a).expect("write a");
2597 let artifact_a = write_directory(&temp_a).expect("a writes");
2598 std::fs::remove_dir_all(&temp_a).ok();
2599
2600 let temp_b = std::env::temp_dir().join(format!(
2601 "limnifs-write-test-{}-cdc-dedup-b",
2602 std::process::id()
2603 ));
2604 std::fs::create_dir_all(&temp_b).expect("create temp_b");
2605 std::fs::write(temp_b.join("b.bin"), &b).expect("write b");
2606 let artifact_b = write_directory(&temp_b).expect("b writes");
2607 std::fs::remove_dir_all(&temp_b).ok();
2608
2609 let artifact_both = write_directory(&temp).expect("both write");
2610 std::fs::remove_dir_all(&temp).ok();
2611
2612 let sum_alone = artifact_a.drop_count + artifact_b.drop_count;
2613 assert!(
2614 artifact_both.drop_count < sum_alone,
2615 "expected dedup win: both together = {} drops, sum alone = {} drops",
2616 artifact_both.drop_count,
2617 sum_alone
2618 );
2619 }
2620
2621 #[test]
2622 fn slab_splits_when_content_exceeds_ceiling() {
2623 let temp =
2629 std::env::temp_dir().join(format!("limnifs-write-test-{}-split", std::process::id()));
2630 std::fs::create_dir_all(&temp).expect("create temp dir");
2631 for i in 0..7u32 {
2632 let data = pseudo_random_bytes(u64::from(i), 10 * 1024 * 1024);
2634 std::fs::write(temp.join(format!("big-{i}.bin")), &data).expect("write big");
2635 }
2636 let artifact = write_directory(&temp).expect("write succeeds");
2637 std::fs::remove_dir_all(&temp).ok();
2638
2639 assert!(
2641 artifact.slabs.len() >= 2,
2642 "expected at least 2 slabs for 70 MiB of incompressible data, got {}",
2643 artifact.slabs.len()
2644 );
2645 for slab in &artifact.slabs {
2646 assert!(
2647 slab.bytes.len() <= MAX_SLAB_TOTAL_BYTES,
2648 "slab {} is {} bytes (> {} ceiling)",
2649 slab.id.ordinal,
2650 slab.bytes.len(),
2651 MAX_SLAB_TOTAL_BYTES,
2652 );
2653 }
2654 let total_drop_ids: usize = artifact.slabs.iter().map(|s| s.drop_ids.len()).sum();
2656 assert_eq!(
2657 total_drop_ids, artifact.drop_count,
2658 "drop_ids count across slabs must match WriteArtifact.drop_count",
2659 );
2660 }
2661}