Skip to main content

wire/
native_pack.rs

1// SPDX-License-Identifier: Apache-2.0
2use std::{
3    fs::{self, File, OpenOptions},
4    io::{Read, Write},
5    path::{Path, PathBuf},
6    time::{SystemTime, UNIX_EPOCH},
7};
8
9use objects::store::{
10    CompressionConfig, ObjectStore,
11    pack::{PackBuilder, PackObjectId, PackReader, StreamingPackBuilder},
12};
13
14use crate::{
15    ObjectData, ObjectId, ObjectInfo, ObjectType, ProtocolError, Result, load_object_data,
16};
17
18/// Maximum hosted native-pack body accepted by the receive primitive.
19///
20/// Native sync packs are produced from bounded state-closure wants and
21/// each decoded pack object is separately capped at 1 GiB in the pack
22/// reader. A 2 GiB compressed pack is materially above normal hosted
23/// sync use while still preventing an untrusted server from growing the
24/// in-memory receive buffer without limit. The receive path can now move
25/// to temp-file spooling plus `install_pack_streaming` — that install API
26/// reports the installed ids the receiver needs, so only the spooling of
27/// the receive buffer itself remains.
28pub const MAX_RECEIVED_PACK_SIZE: u64 = 2 * 1024 * 1024 * 1024;
29
30/// Maximum hosted native-pack index accepted by the receive primitive.
31///
32/// Pack indexes are proportional to object count, not object payload
33/// size. 256 MiB leaves room for millions of entries while bounding the
34/// second in-memory buffer controlled by the remote sender.
35pub const MAX_RECEIVED_PACK_INDEX_SIZE: u64 = 256 * 1024 * 1024;
36
37/// Maximum hosted Git pack accepted by the Git-lane transfer primitive.
38///
39/// Git-overlay sync sends Git-shaped data as raw Git packs. The sender and
40/// receiver still stream those bytes in bounded chunks, but the declared pack
41/// size is untrusted wire input and needs a hard ceiling before buffering or
42/// spooling work begins.
43pub const MAX_RECEIVED_GIT_PACK_SIZE: u64 = 2 * 1024 * 1024 * 1024;
44
45#[derive(Debug, Clone)]
46pub struct NativePackBundle {
47    pub pack_data: Vec<u8>,
48    pub index_data: Vec<u8>,
49}
50
51#[derive(Debug)]
52pub struct NativePackFileBundle {
53    dir: PathBuf,
54    pub pack_path: PathBuf,
55    pub index_path: PathBuf,
56    pub pack_len: u64,
57    pub index_len: u64,
58}
59
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub struct ReusedNativePackStats {
62    pub object_count: usize,
63    pub encoded_bytes_copied: u64,
64}
65
66/// Build a hosted transport pack by reusing non-delta encoded entries from an
67/// authoritative local pack. `Ok(None)` means the caller must use the normal
68/// object-loading writer.
69pub fn reuse_native_pack_encoded_subset_in(
70    root: &Path,
71    source_pack_path: &Path,
72    objects: &[ObjectInfo],
73) -> Result<Option<(NativePackFileBundle, ReusedNativePackStats)>> {
74    if objects.is_empty()
75        || objects
76            .iter()
77            .any(|object| !object.obj_type.packable_for_push())
78    {
79        return Ok(None);
80    }
81    let source_index_path = source_pack_path.with_extension("idx");
82    if !source_pack_path.is_file() || !source_index_path.is_file() {
83        return Ok(None);
84    }
85    let reader = PackReader::open(source_pack_path, &source_index_path, root)?;
86    let expected = objects
87        .iter()
88        .map(|object| {
89            Ok((
90                to_pack_object_id(&object.id, object.obj_type),
91                object.obj_type.pack_object_type()?,
92                object.size,
93            ))
94        })
95        .collect::<Result<Vec<_>>>()?;
96    let Some(reused) = reader.copy_hosted_encoded_subset(&expected)? else {
97        return Ok(None);
98    };
99
100    let base = root.join("transfer-spool");
101    fs::create_dir_all(&base)?;
102    let dir = unique_spool_dir(&base)?;
103    let pack_path = dir.join("pack");
104    let index_path = dir.join("idx");
105    let write_result = (|| -> Result<(u64, u64)> {
106        fs::write(&pack_path, &reused.pack_data)?;
107        fs::write(&index_path, &reused.index_data)?;
108        Ok((
109            u64::try_from(reused.pack_data.len()).map_err(|_| {
110                ProtocolError::InvalidState("reused pack length exceeds u64".to_string())
111            })?,
112            u64::try_from(reused.index_data.len()).map_err(|_| {
113                ProtocolError::InvalidState("reused pack index length exceeds u64".to_string())
114            })?,
115        ))
116    })();
117    let (pack_len, index_len) = match write_result {
118        Ok(lengths) => lengths,
119        Err(error) => {
120            let _ = fs::remove_dir_all(&dir);
121            return Err(error);
122        }
123    };
124    Ok(Some((
125        NativePackFileBundle {
126            dir,
127            pack_path,
128            index_path,
129            pack_len,
130            index_len,
131        },
132        ReusedNativePackStats {
133            object_count: objects.len(),
134            encoded_bytes_copied: reused.encoded_bytes_copied,
135        },
136    )))
137}
138
139impl Drop for NativePackFileBundle {
140    fn drop(&mut self) {
141        let _ = fs::remove_dir_all(&self.dir);
142    }
143}
144
145#[derive(Debug)]
146pub struct PackFileChunkReader {
147    file: File,
148    total_len: u64,
149    chunk_size: usize,
150    offset: u64,
151    chunk_index: u32,
152}
153
154pub type NativePackFileChunk = (u64, u32, Vec<u8>, bool);
155
156impl PackFileChunkReader {
157    pub fn open(path: &Path, chunk_size: usize) -> Result<Self> {
158        let file = File::open(path)?;
159        let total_len = file.metadata()?.len();
160        Ok(Self {
161            file,
162            total_len,
163            chunk_size: chunk_size.max(1),
164            offset: 0,
165            chunk_index: 0,
166        })
167    }
168
169    pub fn next_chunk(&mut self) -> Result<Option<NativePackFileChunk>> {
170        if self.offset >= self.total_len {
171            return Ok(None);
172        }
173        let remaining = self.total_len - self.offset;
174        let len = remaining.min(self.chunk_size as u64);
175        let len = usize::try_from(len).map_err(|_| {
176            ProtocolError::InvalidState("native pack file chunk length exceeds usize".to_string())
177        })?;
178        let mut data = vec![0u8; len];
179        self.file.read_exact(&mut data)?;
180
181        let offset = self.offset;
182        let chunk_index = self.chunk_index;
183        self.offset = self.offset.checked_add(len as u64).ok_or_else(|| {
184            ProtocolError::InvalidState("native pack file chunk offset overflow".to_string())
185        })?;
186        self.chunk_index = self.chunk_index.checked_add(1).ok_or_else(|| {
187            ProtocolError::InvalidState("native pack file chunk index overflow".to_string())
188        })?;
189        Ok(Some((
190            offset,
191            chunk_index,
192            data,
193            self.offset == self.total_len,
194        )))
195    }
196}
197
198#[derive(Debug)]
199pub struct GrowingPackChunkReader {
200    file: File,
201    chunk_size: usize,
202    offset: u64,
203    chunk_index: u32,
204}
205
206impl GrowingPackChunkReader {
207    pub fn open(path: &Path, chunk_size: usize) -> Result<Self> {
208        Ok(Self {
209            file: File::open(path)?,
210            chunk_size: chunk_size.max(1),
211            offset: 0,
212            chunk_index: 0,
213        })
214    }
215
216    pub fn next_available_chunk(
217        &mut self,
218        final_stream: bool,
219    ) -> Result<Option<NativePackFileChunk>> {
220        let total_len = self.file.metadata()?.len();
221        if self.offset >= total_len {
222            return Ok(None);
223        }
224        let available = total_len - self.offset;
225        if !final_stream && available < self.chunk_size as u64 {
226            return Ok(None);
227        }
228
229        let len = available.min(self.chunk_size as u64);
230        let len = usize::try_from(len).map_err(|_| {
231            ProtocolError::InvalidState(
232                "growing native pack chunk length exceeds usize".to_string(),
233            )
234        })?;
235        let mut data = vec![0u8; len];
236        self.file.read_exact(&mut data)?;
237
238        let offset = self.offset;
239        let chunk_index = self.chunk_index;
240        self.offset = self.offset.checked_add(len as u64).ok_or_else(|| {
241            ProtocolError::InvalidState("growing native pack chunk offset overflow".to_string())
242        })?;
243        self.chunk_index = self.chunk_index.checked_add(1).ok_or_else(|| {
244            ProtocolError::InvalidState("growing native pack chunk index overflow".to_string())
245        })?;
246        Ok(Some((
247            offset,
248            chunk_index,
249            data,
250            final_stream && self.offset == total_len,
251        )))
252    }
253}
254
255pub struct NativePackStreamingWriter {
256    dir: Option<PathBuf>,
257    pack_path: PathBuf,
258    index_path: PathBuf,
259    builder: Option<StreamingPackBuilder<File>>,
260}
261
262impl NativePackStreamingWriter {
263    pub fn new_in(root: &Path, object_count: u64) -> Result<Self> {
264        let base = root.join("transfer-spool");
265        fs::create_dir_all(&base)?;
266        let dir = unique_spool_dir(&base)?;
267        let pack_path = dir.join("pack");
268        let index_path = dir.join("idx");
269        let bucket_dir = dir.join("buckets");
270        let pack_file = OpenOptions::new()
271            .read(true)
272            .write(true)
273            .create_new(true)
274            .open(&pack_path)?;
275        let builder = StreamingPackBuilder::new_with_object_count_ephemeral(
276            pack_file,
277            index_path.clone(),
278            sync_pack_compression(),
279            bucket_dir,
280            object_count,
281        )
282        .map_err(ProtocolError::from)?;
283
284        Ok(Self {
285            dir: Some(dir),
286            pack_path,
287            index_path,
288            builder: Some(builder),
289        })
290    }
291
292    pub fn pack_path(&self) -> &Path {
293        &self.pack_path
294    }
295
296    pub fn index_path(&self) -> &Path {
297        &self.index_path
298    }
299
300    pub fn add_object_data(&mut self, object: ObjectData) -> Result<()> {
301        if !is_native_packable_object_type(object.obj_type) {
302            return Err(ProtocolError::InvalidState(format!(
303                "{:?} sidecar records cannot be packed into the content-addressed object pack",
304                object.obj_type
305            )));
306        }
307        let builder = self.builder.as_mut().ok_or_else(|| {
308            ProtocolError::InvalidState("native pack streaming writer is finalized".to_string())
309        })?;
310        let pack_id = to_pack_object_id(&object.id, object.obj_type);
311        builder
312            .add_id(pack_id, object.obj_type.pack_object_type()?, object.data)
313            .map_err(ProtocolError::from)
314    }
315
316    pub fn flush_pack(&mut self) -> Result<()> {
317        let builder = self.builder.as_mut().ok_or_else(|| {
318            ProtocolError::InvalidState("native pack streaming writer is finalized".to_string())
319        })?;
320        builder.flush_pack().map_err(ProtocolError::from)
321    }
322
323    pub fn finish(mut self) -> Result<NativePackFileBundle> {
324        let builder = self.builder.take().ok_or_else(|| {
325            ProtocolError::InvalidState("native pack streaming writer is finalized".to_string())
326        })?;
327        let (mut file, _) = builder.finalize().map_err(ProtocolError::from)?;
328        file.flush()?;
329        drop(file);
330        let pack_len = fs::metadata(&self.pack_path)?.len();
331        let index_len = fs::metadata(&self.index_path)?.len();
332        let dir = self.dir.take().ok_or_else(|| {
333            ProtocolError::InvalidState("native pack streaming writer lost spool dir".to_string())
334        })?;
335        Ok(NativePackFileBundle {
336            dir,
337            pack_path: self.pack_path.clone(),
338            index_path: self.index_path.clone(),
339            pack_len,
340            index_len,
341        })
342    }
343}
344
345impl Drop for NativePackStreamingWriter {
346    fn drop(&mut self) {
347        if let Some(dir) = self.dir.take() {
348            let _ = fs::remove_dir_all(dir);
349        }
350    }
351}
352
353#[derive(Debug, Default, Clone)]
354pub struct PackChunkState {
355    pub pack_data: Vec<u8>,
356    pub index_data: Vec<u8>,
357    pack_progress: (u64, u32),
358    index_progress: (u64, u32),
359    pack_complete: bool,
360    index_complete: bool,
361}
362
363impl PackChunkState {
364    pub fn is_complete(&self) -> bool {
365        self.pack_complete && self.index_complete
366    }
367}
368
369#[derive(Debug, Default, Clone)]
370pub struct GitPackChunkState {
371    transfer_id: Option<String>,
372    pack_size: Option<u64>,
373    next_offset: u64,
374    next_chunk_index: u32,
375    pack_data: Vec<u8>,
376}
377
378impl GitPackChunkState {
379    pub fn is_idle(&self) -> bool {
380        self.transfer_id.is_none()
381            && self.pack_size.is_none()
382            && self.next_offset == 0
383            && self.next_chunk_index == 0
384            && self.pack_data.is_empty()
385    }
386
387    pub fn ensure_idle(&self) -> Result<()> {
388        if self.is_idle() {
389            Ok(())
390        } else {
391            Err(ProtocolError::InvalidState(
392                "Git pack transfer ended before final chunk".to_string(),
393            ))
394        }
395    }
396
397    pub fn receive_chunk(
398        &mut self,
399        transfer_id: &str,
400        offset: u64,
401        chunk_index: u32,
402        is_final_chunk: bool,
403        pack_size: u64,
404        data: &[u8],
405    ) -> Result<Option<Vec<u8>>> {
406        if transfer_id.is_empty() {
407            return Err(ProtocolError::InvalidState(
408                "Git pack transfer_id is required".to_string(),
409            ));
410        }
411        if pack_size > MAX_RECEIVED_GIT_PACK_SIZE {
412            return Err(ProtocolError::InvalidState(format!(
413                "Git pack exceeds maximum transfer size of {MAX_RECEIVED_GIT_PACK_SIZE} bytes"
414            )));
415        }
416        if data.is_empty() {
417            return Err(ProtocolError::InvalidState(
418                "Git pack chunk must not be empty".to_string(),
419            ));
420        }
421        match self.transfer_id.as_ref() {
422            Some(current) if current != transfer_id => {
423                return Err(ProtocolError::InvalidState(format!(
424                    "Git pack transfer id changed from {current:?} to {transfer_id:?}"
425                )));
426            }
427            Some(_) => {}
428            None => {
429                self.transfer_id = Some(transfer_id.to_string());
430                self.pack_size = Some(pack_size);
431            }
432        }
433        if self.pack_size != Some(pack_size) {
434            return Err(ProtocolError::InvalidState(
435                "Git pack size changed during transfer".to_string(),
436            ));
437        }
438        if offset != self.next_offset {
439            return Err(ProtocolError::InvalidState(format!(
440                "Git pack offset mismatch: expected {}, got {}",
441                self.next_offset, offset
442            )));
443        }
444        if chunk_index != self.next_chunk_index {
445            return Err(ProtocolError::InvalidState(format!(
446                "Git pack chunk index mismatch: expected {}, got {}",
447                self.next_chunk_index, chunk_index
448            )));
449        }
450        let chunk_len = u64::try_from(data.len()).map_err(|_| {
451            ProtocolError::InvalidState("Git pack chunk length exceeds u64".to_string())
452        })?;
453        let next_offset = self
454            .next_offset
455            .checked_add(chunk_len)
456            .ok_or_else(|| ProtocolError::InvalidState("Git pack offset overflow".to_string()))?;
457        if next_offset > pack_size {
458            return Err(ProtocolError::InvalidState(
459                "Git pack chunk exceeds declared pack size".to_string(),
460            ));
461        }
462        self.pack_data.extend_from_slice(data);
463        self.next_offset = next_offset;
464        self.next_chunk_index = self.next_chunk_index.checked_add(1).ok_or_else(|| {
465            ProtocolError::InvalidState("Git pack chunk index overflow".to_string())
466        })?;
467        if is_final_chunk {
468            if self.next_offset != pack_size {
469                return Err(ProtocolError::InvalidState(format!(
470                    "Git pack final size mismatch: declared {}, received {}",
471                    pack_size, self.next_offset
472                )));
473            }
474            let pack_data = std::mem::take(&mut self.pack_data);
475            self.transfer_id = None;
476            self.pack_size = None;
477            self.next_offset = 0;
478            self.next_chunk_index = 0;
479            return Ok(Some(pack_data));
480        }
481        if self.next_offset == pack_size {
482            return Err(ProtocolError::InvalidState(
483                "Git pack reached declared size without final chunk marker".to_string(),
484            ));
485        }
486        Ok(None)
487    }
488}
489
490#[derive(Debug)]
491pub struct PackChunkSpool {
492    dir: PathBuf,
493    pack: PackStreamSpool,
494    index: PackStreamSpool,
495}
496
497impl PackChunkSpool {
498    pub fn new_in(root: &Path) -> Result<Self> {
499        let base = root.join("transfer-spool");
500        fs::create_dir_all(&base)?;
501        let dir = unique_spool_dir(&base)?;
502        let pack = PackStreamSpool::new(dir.join("pack"))?;
503        let index = PackStreamSpool::new(dir.join("idx"))?;
504        Ok(Self { dir, pack, index })
505    }
506
507    pub fn is_complete(&self) -> bool {
508        self.pack.complete && self.index.complete
509    }
510
511    #[allow(clippy::too_many_arguments)]
512    pub fn receive_chunk(
513        &mut self,
514        is_index: bool,
515        resume_offset: u64,
516        chunk_index: u32,
517        is_complete: bool,
518        data: &[u8],
519        is_final_chunk: bool,
520    ) -> Result<()> {
521        let max_bytes = if is_index {
522            MAX_RECEIVED_PACK_INDEX_SIZE
523        } else {
524            MAX_RECEIVED_PACK_SIZE
525        };
526        let stream = if is_index {
527            &mut self.index
528        } else {
529            &mut self.pack
530        };
531        receive_pack_chunk_to_spool(
532            stream,
533            is_index,
534            resume_offset,
535            chunk_index,
536            is_complete,
537            data,
538            is_final_chunk,
539            max_bytes,
540        )
541    }
542
543    pub fn install_into(
544        &mut self,
545        store: &impl ObjectStore,
546    ) -> Result<objects::store::pack::PackInventory> {
547        if !self.is_complete() {
548            return Err(ProtocolError::InvalidState(
549                "native pack spool is incomplete".to_string(),
550            ));
551        }
552        self.pack.close()?;
553        self.index.close()?;
554        store
555            .install_pack_streaming(&self.pack.path, &self.index.path)
556            .map_err(ProtocolError::from)
557    }
558}
559
560impl Drop for PackChunkSpool {
561    fn drop(&mut self) {
562        let _ = fs::remove_dir_all(&self.dir);
563    }
564}
565
566#[derive(Debug)]
567struct PackStreamSpool {
568    path: PathBuf,
569    file: Option<File>,
570    progress: (u64, u32),
571    complete: bool,
572}
573
574impl PackStreamSpool {
575    fn new(path: PathBuf) -> Result<Self> {
576        let file = File::create(&path)?;
577        Ok(Self {
578            path,
579            file: Some(file),
580            progress: (0, 0),
581            complete: false,
582        })
583    }
584
585    fn write_all(&mut self, data: &[u8]) -> Result<()> {
586        let Some(file) = self.file.as_mut() else {
587            return Err(ProtocolError::InvalidState(
588                "native pack spool stream is already closed".to_string(),
589            ));
590        };
591        file.write_all(data)?;
592        Ok(())
593    }
594
595    fn close(&mut self) -> Result<()> {
596        if let Some(mut file) = self.file.take() {
597            file.flush()?;
598            objects::fs_atomic::sync_file(&file, &self.path)?;
599        }
600        Ok(())
601    }
602}
603
604pub fn native_pack_excluded_object_types() -> &'static [ObjectType] {
605    &[
606        ObjectType::Redaction,
607        ObjectType::StateVisibility,
608        ObjectType::KeyBinding,
609    ]
610}
611
612pub fn is_native_packable_object_type(obj_type: ObjectType) -> bool {
613    obj_type.packable()
614}
615
616pub fn build_native_pack(
617    store: &impl ObjectStore,
618    objects: &[ObjectInfo],
619) -> Result<NativePackBundle> {
620    let mut builder = PackBuilder::new(sync_pack_compression());
621
622    for info in objects {
623        // Sidecar records (redaction + state-visibility) live outside
624        // `.heddle/objects/` so GC cannot touch them, and must not be
625        // folded into the content-addressed pack. They ship via the
626        // per-object transfer path instead; callers split them out before
627        // packing.
628        if !is_native_packable_object_type(info.obj_type) {
629            continue;
630        }
631        let object = load_object_data(store, &info.id, info.obj_type)?;
632        let pack_id = to_pack_object_id(&object.id, object.obj_type);
633        builder.add_id(pack_id, object.obj_type.pack_object_type()?, object.data);
634    }
635
636    let (pack_data, index_data, _) = builder.build()?;
637    Ok(NativePackBundle {
638        pack_data,
639        index_data,
640    })
641}
642
643fn sync_pack_compression() -> CompressionConfig {
644    CompressionConfig {
645        level: 1,
646        min_size: 1024,
647        max_delta_size: 0,
648        ..CompressionConfig::default()
649    }
650}
651
652pub fn install_received_pack(
653    store: &impl ObjectStore,
654    pack_data: &[u8],
655    index_data: &[u8],
656) -> Result<Vec<PackObjectId>> {
657    store
658        .install_pack(pack_data, index_data)
659        .map_err(ProtocolError::from)
660}
661
662pub fn next_pack_chunk(
663    data: &[u8],
664    chunk_size: usize,
665    chunk_index: usize,
666) -> Option<(usize, Vec<u8>, bool)> {
667    let (start, len) = crate::chunk_bounds(data.len(), chunk_size.max(1), chunk_index)?;
668    let is_final = start + len == data.len();
669    Some((start, data[start..start + len].to_vec(), is_final))
670}
671
672pub fn receive_pack_chunk(
673    state: &mut PackChunkState,
674    is_index: bool,
675    resume_offset: u64,
676    chunk_index: u32,
677    is_complete: bool,
678    data: &[u8],
679    is_final_chunk: bool,
680) -> Result<()> {
681    let max_bytes = if is_index {
682        MAX_RECEIVED_PACK_INDEX_SIZE
683    } else {
684        MAX_RECEIVED_PACK_SIZE
685    };
686    receive_pack_chunk_with_limit(
687        state,
688        is_index,
689        resume_offset,
690        chunk_index,
691        is_complete,
692        data,
693        is_final_chunk,
694        max_bytes,
695    )
696}
697
698#[allow(clippy::too_many_arguments)]
699fn receive_pack_chunk_with_limit(
700    state: &mut PackChunkState,
701    is_index: bool,
702    resume_offset: u64,
703    chunk_index: u32,
704    is_complete: bool,
705    data: &[u8],
706    is_final_chunk: bool,
707    max_bytes: u64,
708) -> Result<()> {
709    let (buffer, progress, complete) = if is_index {
710        (
711            &mut state.index_data,
712            &mut state.index_progress,
713            &mut state.index_complete,
714        )
715    } else {
716        (
717            &mut state.pack_data,
718            &mut state.pack_progress,
719            &mut state.pack_complete,
720        )
721    };
722
723    let next_progress = validate_pack_chunk(
724        *progress,
725        is_index,
726        resume_offset,
727        chunk_index,
728        data,
729        max_bytes,
730    )?;
731
732    buffer.extend_from_slice(data);
733    *progress = next_progress;
734    if is_final_chunk || is_complete {
735        *complete = true;
736    }
737    Ok(())
738}
739
740#[allow(clippy::too_many_arguments)]
741fn receive_pack_chunk_to_spool(
742    stream: &mut PackStreamSpool,
743    is_index: bool,
744    resume_offset: u64,
745    chunk_index: u32,
746    is_complete: bool,
747    data: &[u8],
748    is_final_chunk: bool,
749    max_bytes: u64,
750) -> Result<()> {
751    let next_progress = validate_pack_chunk(
752        stream.progress,
753        is_index,
754        resume_offset,
755        chunk_index,
756        data,
757        max_bytes,
758    )?;
759    stream.write_all(data)?;
760    stream.progress = next_progress;
761    if is_final_chunk || is_complete {
762        stream.complete = true;
763    }
764    Ok(())
765}
766
767fn validate_pack_chunk(
768    progress: (u64, u32),
769    is_index: bool,
770    resume_offset: u64,
771    chunk_index: u32,
772    data: &[u8],
773    max_bytes: u64,
774) -> Result<(u64, u32)> {
775    if resume_offset != progress.0 {
776        return Err(ProtocolError::InvalidState(format!(
777            "native pack chunk resume offset mismatch: expected {}, got {}",
778            progress.0, resume_offset
779        )));
780    }
781    if chunk_index != progress.1 {
782        return Err(ProtocolError::InvalidState(format!(
783            "native pack chunk index mismatch: expected {}, got {}",
784            progress.1, chunk_index
785        )));
786    }
787
788    let data_len = u64::try_from(data.len()).map_err(|_| {
789        ProtocolError::InvalidState("native pack chunk length does not fit in u64".to_string())
790    })?;
791    let next_offset = progress.0.checked_add(data_len).ok_or_else(|| {
792        ProtocolError::InvalidState("native pack chunk offset overflow".to_string())
793    })?;
794    if next_offset > max_bytes {
795        let stream_name = if is_index { "index" } else { "body" };
796        return Err(ProtocolError::InvalidState(format!(
797            "native pack {stream_name} exceeds receive size limit: {next_offset} bytes (max {max_bytes})"
798        )));
799    }
800    let next_chunk = progress.1.checked_add(1).ok_or_else(|| {
801        ProtocolError::InvalidState("native pack chunk index overflow".to_string())
802    })?;
803
804    Ok((next_offset, next_chunk))
805}
806
807pub(crate) fn unique_spool_dir(base: &Path) -> Result<PathBuf> {
808    let stamp = SystemTime::now()
809        .duration_since(UNIX_EPOCH)
810        .map_err(|err| {
811            ProtocolError::InvalidState(format!("system clock before UNIX epoch: {err}"))
812        })?
813        .as_nanos();
814    for attempt in 0..100u32 {
815        let dir = base.join(format!("pack-{}-{stamp}-{attempt}", std::process::id()));
816        match fs::create_dir(&dir) {
817            Ok(()) => return Ok(dir),
818            Err(err) if err.kind() == std::io::ErrorKind::AlreadyExists => continue,
819            Err(err) => return Err(ProtocolError::Io(err)),
820        }
821    }
822    Err(ProtocolError::InvalidState(
823        "failed to allocate native pack spool directory".to_string(),
824    ))
825}
826
827fn to_pack_object_id(id: &ObjectId, object_type: ObjectType) -> PackObjectId {
828    match (id, object_type) {
829        (ObjectId::Hash(hash), ObjectType::AnnotatedTag) => PackObjectId::AnnotatedTag(*hash),
830        (ObjectId::Hash(hash), _) => PackObjectId::Hash(*hash),
831        (ObjectId::StateId(state_id), _) => PackObjectId::StateId(*state_id),
832        (ObjectId::StateAttachment { id, .. }, _) => PackObjectId::Hash(*id.as_hash()),
833    }
834}
835
836#[cfg(test)]
837mod tests {
838    use objects::{
839        object::{AnnotatedTag, Blob, ContentHash, StateId},
840        store::{
841            CompressionConfig, FsStore, ObjectStore,
842            pack::{ObjectType as PackObjectType, PackBuilder, PackObjectId, PackReader},
843        },
844    };
845    use sley::ObjectFormat as GitObjectFormat;
846    use tempfile::TempDir;
847
848    use super::{
849        GitPackChunkState, GrowingPackChunkReader, MAX_RECEIVED_PACK_SIZE,
850        NativePackStreamingWriter, ObjectData, ObjectId, ObjectInfo, ObjectType, PackChunkSpool,
851        PackChunkState, PackFileChunkReader, build_native_pack, install_received_pack,
852        next_pack_chunk, receive_pack_chunk, receive_pack_chunk_with_limit,
853        reuse_native_pack_encoded_subset_in,
854    };
855
856    fn create_test_store() -> (TempDir, FsStore) {
857        let temp = TempDir::new().unwrap();
858        let store = FsStore::new(temp.path().join(".heddle"));
859        store.init().unwrap();
860        (temp, store)
861    }
862
863    fn hash(byte: u8) -> ContentHash {
864        ContentHash::from_bytes([byte; 32])
865    }
866
867    #[test]
868    fn encoded_snapshot_subset_is_wire_equivalent_without_local_artifacts_or_attachments() {
869        let source = TempDir::new().unwrap();
870        let spool = TempDir::new().unwrap();
871        let source_pack = source.path().join("snapshot.pack");
872        let source_index = source.path().join("snapshot.idx");
873        let blob = (
874            PackObjectId::Hash(hash(1)),
875            PackObjectType::Blob,
876            b"blob body".to_vec(),
877        );
878        let tree = (
879            PackObjectId::Hash(hash(2)),
880            PackObjectType::Tree,
881            b"tree body".to_vec(),
882        );
883        let state_id = StateId::from_bytes([3; 32]);
884        let state = (
885            PackObjectId::StateId(state_id),
886            PackObjectType::State,
887            b"state body".to_vec(),
888        );
889        let attachment_id = PackObjectId::Hash(hash(4));
890        let artifact_id = PackObjectId::Hash(hash(5));
891        let mut builder = PackBuilder::new(CompressionConfig {
892            max_delta_size: 0,
893            ..CompressionConfig::default()
894        });
895        for (id, kind, body) in [
896            blob.clone(),
897            tree.clone(),
898            state.clone(),
899            (
900                attachment_id,
901                PackObjectType::StateAttachment,
902                b"local attachment".to_vec(),
903            ),
904            (
905                artifact_id,
906                PackObjectType::SnapshotCommit,
907                b"local commit artifact".to_vec(),
908            ),
909        ] {
910            builder.add_id(id, kind, body);
911        }
912        let (pack, index, _) = builder.build().unwrap();
913        std::fs::write(&source_pack, pack).unwrap();
914        std::fs::write(&source_index, index).unwrap();
915
916        let wanted = vec![
917            ObjectInfo {
918                id: ObjectId::Hash(hash(1)),
919                obj_type: ObjectType::Blob,
920                size: blob.2.len() as u64,
921                delta_base: None,
922            },
923            ObjectInfo {
924                id: ObjectId::Hash(hash(2)),
925                obj_type: ObjectType::Tree,
926                size: tree.2.len() as u64,
927                delta_base: None,
928            },
929            ObjectInfo {
930                id: ObjectId::StateId(state_id),
931                obj_type: ObjectType::State,
932                size: state.2.len() as u64,
933                delta_base: None,
934            },
935        ];
936        let (bundle, stats) =
937            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &wanted)
938                .unwrap()
939                .expect("authoritative non-delta subset must be reusable");
940
941        assert_eq!(stats.object_count, wanted.len());
942        assert!(stats.encoded_bytes_copied > 0);
943        let reused =
944            PackReader::open(&bundle.pack_path, &bundle.index_path, &std::env::temp_dir()).unwrap();
945        let mut reused_ids = reused.list_ids().unwrap();
946        reused_ids.sort();
947        let mut wanted_ids = vec![blob.0, tree.0, state.0];
948        wanted_ids.sort();
949        assert_eq!(reused_ids, wanted_ids);
950        assert!(!reused.has_object(&attachment_id).unwrap());
951        assert!(!reused.has_object(&artifact_id).unwrap());
952
953        for path in [&bundle.pack_path, &bundle.index_path] {
954            let expected_wire_bytes = std::fs::read(path).unwrap();
955            let mut chunk_reader = PackFileChunkReader::open(path, 7).unwrap();
956            let mut wire_bytes = Vec::new();
957            while let Some((offset, chunk_index, data, is_final)) =
958                chunk_reader.next_chunk().unwrap()
959            {
960                assert_eq!(offset as usize, wire_bytes.len());
961                assert_eq!(chunk_index as usize, wire_bytes.len() / 7);
962                wire_bytes.extend_from_slice(&data);
963                assert_eq!(is_final, wire_bytes.len() == expected_wire_bytes.len());
964            }
965            assert_eq!(wire_bytes, expected_wire_bytes);
966        }
967        for (id, _, expected) in [blob, tree, state] {
968            assert_eq!(reused.get_object(&id).unwrap().unwrap().1, expected);
969        }
970    }
971
972    #[test]
973    fn encoded_snapshot_subset_falls_back_for_mismatch_delta_or_attachment_request() {
974        let source = TempDir::new().unwrap();
975        let spool = TempDir::new().unwrap();
976        let source_pack = source.path().join("snapshot.pack");
977        let source_index = source.path().join("snapshot.idx");
978        let first = b"This is the base content. ".repeat(100);
979        let second = b"This is modified content. ".repeat(100);
980        let mut builder = PackBuilder::new(CompressionConfig::default());
981        builder.add(hash(10), PackObjectType::Blob, first.clone());
982        builder.add(hash(11), PackObjectType::Blob, second.clone());
983        let (pack, index, stats) = builder.build().unwrap();
984        assert!(stats.delta_count > 0, "fixture must contain a delta");
985        std::fs::write(&source_pack, pack).unwrap();
986        std::fs::write(&source_index, index).unwrap();
987        let delta_wants = [ObjectInfo {
988            id: ObjectId::Hash(hash(11)),
989            obj_type: ObjectType::Blob,
990            size: second.len() as u64,
991            delta_base: None,
992        }];
993        assert!(
994            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &delta_wants)
995                .unwrap()
996                .is_none()
997        );
998
999        let missing_wants = [ObjectInfo {
1000            id: ObjectId::Hash(hash(12)),
1001            obj_type: ObjectType::Blob,
1002            size: 1,
1003            delta_base: None,
1004        }];
1005        assert!(
1006            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &missing_wants)
1007                .unwrap()
1008                .is_none()
1009        );
1010
1011        let attachment_wants = [ObjectInfo {
1012            id: ObjectId::StateAttachment {
1013                state: StateId::from_bytes([13; 32]),
1014                id: objects::object::StateAttachmentId::from_hash(hash(14)),
1015                kind: objects::object::StateAttachmentKind::SemanticIndex,
1016            },
1017            obj_type: ObjectType::StateAttachment,
1018            size: 1,
1019            delta_base: None,
1020        }];
1021        assert!(
1022            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &attachment_wants)
1023                .unwrap()
1024                .is_none()
1025        );
1026    }
1027
1028    #[test]
1029    fn receive_pack_chunk_rejects_cumulative_size_over_limit_before_buffering() {
1030        let mut state = PackChunkState::default();
1031
1032        receive_pack_chunk_with_limit(&mut state, false, 0, 0, false, b"abcd", false, 8).unwrap();
1033        receive_pack_chunk_with_limit(&mut state, false, 4, 1, false, b"efgh", false, 8).unwrap();
1034
1035        let error = receive_pack_chunk_with_limit(&mut state, false, 8, 2, false, b"i", false, 8)
1036            .unwrap_err();
1037
1038        assert_eq!(state.pack_data, b"abcdefgh");
1039        assert!(
1040            error
1041                .to_string()
1042                .contains("native pack body exceeds receive size limit")
1043        );
1044        assert!(error.to_string().contains("9 bytes (max 8)"));
1045    }
1046
1047    #[test]
1048    fn receive_pack_chunk_checks_production_limit_before_extending_buffer() {
1049        let mut state = PackChunkState {
1050            pack_progress: (MAX_RECEIVED_PACK_SIZE - 1, 0),
1051            ..PackChunkState::default()
1052        };
1053
1054        let error = receive_pack_chunk(
1055            &mut state,
1056            false,
1057            MAX_RECEIVED_PACK_SIZE - 1,
1058            0,
1059            false,
1060            b"xx",
1061            false,
1062        )
1063        .unwrap_err();
1064
1065        assert!(state.pack_data.is_empty());
1066        assert!(
1067            error
1068                .to_string()
1069                .contains("native pack body exceeds receive size limit")
1070        );
1071    }
1072
1073    #[test]
1074    fn receive_pack_chunk_rejects_resume_offset_mismatch_before_buffering() {
1075        let mut state = PackChunkState::default();
1076
1077        let error =
1078            receive_pack_chunk(&mut state, false, 1, 0, false, b"late chunk", false).unwrap_err();
1079
1080        assert!(state.pack_data.is_empty());
1081        assert!(
1082            error
1083                .to_string()
1084                .contains("native pack chunk resume offset mismatch: expected 0, got 1")
1085        );
1086    }
1087
1088    #[test]
1089    fn receive_pack_chunk_rejects_chunk_index_mismatch_before_buffering() {
1090        let mut state = PackChunkState::default();
1091
1092        receive_pack_chunk(&mut state, false, 0, 0, false, b"abc", false).unwrap();
1093        let error = receive_pack_chunk(&mut state, false, 3, 2, false, b"def", false).unwrap_err();
1094
1095        assert_eq!(state.pack_data, b"abc");
1096        assert!(
1097            error
1098                .to_string()
1099                .contains("native pack chunk index mismatch: expected 1, got 2")
1100        );
1101    }
1102
1103    #[test]
1104    fn git_pack_chunk_state_requires_ordered_chunks_and_final_size() {
1105        let mut state = GitPackChunkState::default();
1106
1107        assert!(
1108            state
1109                .receive_chunk("git-pack:test", 0, 0, false, 8, b"abcd")
1110                .unwrap()
1111                .is_none()
1112        );
1113        let error = state
1114            .receive_chunk("git-pack:test", 4, 2, true, 8, b"efgh")
1115            .unwrap_err();
1116
1117        assert!(
1118            error
1119                .to_string()
1120                .contains("Git pack chunk index mismatch: expected 1, got 2")
1121        );
1122        assert!(state.ensure_idle().is_err());
1123
1124        let mut state = GitPackChunkState::default();
1125        state
1126            .receive_chunk("git-pack:test", 0, 0, false, 8, b"abcd")
1127            .unwrap();
1128        let complete = state
1129            .receive_chunk("git-pack:test", 4, 1, true, 8, b"efgh")
1130            .unwrap()
1131            .unwrap();
1132
1133        assert_eq!(complete, b"abcdefgh");
1134        assert!(state.ensure_idle().is_ok());
1135    }
1136
1137    #[test]
1138    fn receive_pack_chunk_accepts_completion_flags_for_pack_and_index() {
1139        let mut state = PackChunkState::default();
1140
1141        receive_pack_chunk(&mut state, false, 0, 0, true, b"pack-body", false).unwrap();
1142        assert!(!state.is_complete());
1143        receive_pack_chunk(&mut state, true, 0, 0, false, b"pack-index", true).unwrap();
1144
1145        assert!(state.is_complete());
1146        assert_eq!(state.pack_data, b"pack-body");
1147        assert_eq!(state.index_data, b"pack-index");
1148    }
1149
1150    #[test]
1151    fn normal_size_native_pack_receives_and_installs() {
1152        let (_source_temp, source_store) = create_test_store();
1153        let (_dest_temp, dest_store) = create_test_store();
1154        let blob = Blob::from("native pack receive regression");
1155        let hash = source_store.put_blob(&blob).unwrap();
1156        let bundle = build_native_pack(
1157            &source_store,
1158            &[ObjectInfo {
1159                id: ObjectId::Hash(hash),
1160                obj_type: ObjectType::Blob,
1161                size: blob.size() as u64,
1162                delta_base: None,
1163            }],
1164        )
1165        .unwrap();
1166
1167        let mut state = PackChunkState::default();
1168        let mut chunk_index = 0usize;
1169        while let Some((start, data, is_final)) = next_pack_chunk(&bundle.pack_data, 7, chunk_index)
1170        {
1171            receive_pack_chunk(
1172                &mut state,
1173                false,
1174                start as u64,
1175                chunk_index as u32,
1176                is_final,
1177                &data,
1178                is_final,
1179            )
1180            .unwrap();
1181            chunk_index += 1;
1182        }
1183
1184        let mut index_chunk = 0usize;
1185        while let Some((start, data, is_final)) =
1186            next_pack_chunk(&bundle.index_data, 5, index_chunk)
1187        {
1188            receive_pack_chunk(
1189                &mut state,
1190                true,
1191                start as u64,
1192                index_chunk as u32,
1193                is_final,
1194                &data,
1195                is_final,
1196            )
1197            .unwrap();
1198            index_chunk += 1;
1199        }
1200
1201        assert!(state.is_complete());
1202        assert_eq!(state.pack_data, bundle.pack_data);
1203        assert_eq!(state.index_data, bundle.index_data);
1204
1205        let installed_ids =
1206            install_received_pack(&dest_store, &state.pack_data, &state.index_data).unwrap();
1207
1208        assert_eq!(installed_ids, vec![PackObjectId::Hash(hash)]);
1209        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1210        assert_eq!(installed_blob.content(), blob.content());
1211    }
1212
1213    #[test]
1214    fn native_pack_transfers_first_class_annotated_tag() {
1215        let (_source_temp, source_store) = create_test_store();
1216        let (_dest_temp, dest_store) = create_test_store();
1217        let tag = AnnotatedTag::new(
1218            GitObjectFormat::Sha1,
1219            b"object 1111111111111111111111111111111111111111\ntype commit\ntag v1\ntagger Test <test@example.com> 1700000000 +0100\n\nrelease\n".to_vec(),
1220            None,
1221            None,
1222        )
1223        .unwrap();
1224        let hash = source_store.put_annotated_tag(&tag).unwrap();
1225        let bundle = build_native_pack(
1226            &source_store,
1227            &[ObjectInfo {
1228                id: ObjectId::Hash(hash),
1229                obj_type: ObjectType::AnnotatedTag,
1230                size: tag.encode_current_msgpack().len() as u64,
1231                delta_base: None,
1232            }],
1233        )
1234        .unwrap();
1235
1236        let installed = install_received_pack(&dest_store, &bundle.pack_data, &bundle.index_data)
1237            .expect("install annotated-tag native pack");
1238
1239        assert_eq!(installed, vec![PackObjectId::AnnotatedTag(hash)]);
1240        assert_eq!(dest_store.get_annotated_tag(&hash).unwrap(), Some(tag));
1241    }
1242
1243    #[test]
1244    fn normal_size_native_pack_spools_and_installs() {
1245        let (_source_temp, source_store) = create_test_store();
1246        let (dest_temp, dest_store) = create_test_store();
1247        let blob = Blob::from("native pack spooled receive regression");
1248        let hash = source_store.put_blob(&blob).unwrap();
1249        let bundle = build_native_pack(
1250            &source_store,
1251            &[ObjectInfo {
1252                id: ObjectId::Hash(hash),
1253                obj_type: ObjectType::Blob,
1254                size: blob.size() as u64,
1255                delta_base: None,
1256            }],
1257        )
1258        .unwrap();
1259
1260        let mut spool = PackChunkSpool::new_in(dest_temp.path()).unwrap();
1261        let mut chunk_index = 0usize;
1262        while let Some((start, data, is_final)) = next_pack_chunk(&bundle.pack_data, 7, chunk_index)
1263        {
1264            spool
1265                .receive_chunk(
1266                    false,
1267                    start as u64,
1268                    chunk_index as u32,
1269                    is_final,
1270                    &data,
1271                    is_final,
1272                )
1273                .unwrap();
1274            chunk_index += 1;
1275        }
1276
1277        let mut index_chunk = 0usize;
1278        while let Some((start, data, is_final)) =
1279            next_pack_chunk(&bundle.index_data, 5, index_chunk)
1280        {
1281            spool
1282                .receive_chunk(
1283                    true,
1284                    start as u64,
1285                    index_chunk as u32,
1286                    is_final,
1287                    &data,
1288                    is_final,
1289                )
1290                .unwrap();
1291            index_chunk += 1;
1292        }
1293
1294        assert!(spool.is_complete());
1295        let installed_ids = spool
1296            .install_into(&dest_store)
1297            .unwrap()
1298            .ids()
1299            .collect::<objects::store::Result<Vec<_>>>()
1300            .unwrap();
1301
1302        assert_eq!(installed_ids, vec![PackObjectId::Hash(hash)]);
1303        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1304        assert_eq!(installed_blob.content(), blob.content());
1305    }
1306
1307    #[test]
1308    fn native_pack_streaming_writer_drains_growing_pack_and_installs() {
1309        let (source_temp, source_store) = create_test_store();
1310        let (dest_temp, dest_store) = create_test_store();
1311        let blob = Blob::from("native pack growing stream regression");
1312        let hash = source_store.put_blob(&blob).unwrap();
1313        let large_blob = Blob::from_slice(&vec![b'z'; 4096]);
1314        let large_hash = source_store.put_blob(&large_blob).unwrap();
1315
1316        let mut writer = NativePackStreamingWriter::new_in(source_temp.path(), 2).unwrap();
1317        let mut pack_reader = GrowingPackChunkReader::open(writer.pack_path(), 31).unwrap();
1318        let mut spool = PackChunkSpool::new_in(dest_temp.path()).unwrap();
1319        let mut saw_interleaved_pack_chunk = false;
1320
1321        for (id, obj_type, data) in [
1322            (
1323                ObjectId::Hash(hash),
1324                ObjectType::Blob,
1325                blob.content().to_vec(),
1326            ),
1327            (
1328                ObjectId::Hash(large_hash),
1329                ObjectType::Blob,
1330                large_blob.content().to_vec(),
1331            ),
1332        ] {
1333            writer
1334                .add_object_data(ObjectData {
1335                    id,
1336                    obj_type,
1337                    data,
1338                    is_delta: false,
1339                })
1340                .unwrap();
1341            writer.flush_pack().unwrap();
1342            while let Some((offset, chunk_index, data, is_final)) =
1343                pack_reader.next_available_chunk(false).unwrap()
1344            {
1345                assert!(
1346                    !is_final,
1347                    "pre-final growing pack drain must not mark chunks final"
1348                );
1349                saw_interleaved_pack_chunk = true;
1350                spool
1351                    .receive_chunk(false, offset, chunk_index, false, &data, false)
1352                    .unwrap();
1353            }
1354        }
1355
1356        let bundle = writer.finish().unwrap();
1357        let mut saw_final_pack_chunk = false;
1358        while let Some((offset, chunk_index, data, is_final)) =
1359            pack_reader.next_available_chunk(true).unwrap()
1360        {
1361            saw_final_pack_chunk |= is_final;
1362            spool
1363                .receive_chunk(false, offset, chunk_index, is_final, &data, is_final)
1364                .unwrap();
1365        }
1366
1367        let mut index_reader = PackFileChunkReader::open(&bundle.index_path, 17).unwrap();
1368        while let Some((offset, chunk_index, data, is_final)) = index_reader.next_chunk().unwrap() {
1369            spool
1370                .receive_chunk(true, offset, chunk_index, is_final, &data, is_final)
1371                .unwrap();
1372        }
1373
1374        assert!(
1375            saw_interleaved_pack_chunk,
1376            "expected at least one pack chunk before finalize"
1377        );
1378        assert!(
1379            saw_final_pack_chunk,
1380            "expected final pack chunk after finish"
1381        );
1382        assert!(spool.is_complete());
1383        let mut installed_ids = spool
1384            .install_into(&dest_store)
1385            .unwrap()
1386            .ids()
1387            .collect::<objects::store::Result<Vec<_>>>()
1388            .unwrap();
1389        let mut expected_ids = vec![PackObjectId::Hash(hash), PackObjectId::Hash(large_hash)];
1390        installed_ids.sort();
1391        expected_ids.sort();
1392
1393        assert_eq!(installed_ids, expected_ids);
1394        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1395        assert_eq!(installed_blob.content(), blob.content());
1396        let installed_large_blob = dest_store.get_blob(&large_hash).unwrap().unwrap();
1397        assert_eq!(installed_large_blob.content(), large_blob.content());
1398    }
1399}