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)?;
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(&mut self, store: &impl ObjectStore) -> Result<Vec<PackObjectId>> {
544        if !self.is_complete() {
545            return Err(ProtocolError::InvalidState(
546                "native pack spool is incomplete".to_string(),
547            ));
548        }
549        self.pack.close()?;
550        self.index.close()?;
551        store
552            .install_pack_streaming(&self.pack.path, &self.index.path)
553            .map_err(ProtocolError::from)
554    }
555}
556
557impl Drop for PackChunkSpool {
558    fn drop(&mut self) {
559        let _ = fs::remove_dir_all(&self.dir);
560    }
561}
562
563#[derive(Debug)]
564struct PackStreamSpool {
565    path: PathBuf,
566    file: Option<File>,
567    progress: (u64, u32),
568    complete: bool,
569}
570
571impl PackStreamSpool {
572    fn new(path: PathBuf) -> Result<Self> {
573        let file = File::create(&path)?;
574        Ok(Self {
575            path,
576            file: Some(file),
577            progress: (0, 0),
578            complete: false,
579        })
580    }
581
582    fn write_all(&mut self, data: &[u8]) -> Result<()> {
583        let Some(file) = self.file.as_mut() else {
584            return Err(ProtocolError::InvalidState(
585                "native pack spool stream is already closed".to_string(),
586            ));
587        };
588        file.write_all(data)?;
589        Ok(())
590    }
591
592    fn close(&mut self) -> Result<()> {
593        if let Some(mut file) = self.file.take() {
594            file.flush()?;
595            objects::fs_atomic::sync_file(&file, &self.path)?;
596        }
597        Ok(())
598    }
599}
600
601pub fn native_pack_excluded_object_types() -> &'static [ObjectType] {
602    &[
603        ObjectType::Redaction,
604        ObjectType::StateVisibility,
605        ObjectType::KeyBinding,
606    ]
607}
608
609pub fn is_native_packable_object_type(obj_type: ObjectType) -> bool {
610    obj_type.packable()
611}
612
613pub fn build_native_pack(
614    store: &impl ObjectStore,
615    objects: &[ObjectInfo],
616) -> Result<NativePackBundle> {
617    let mut builder = PackBuilder::new(sync_pack_compression());
618
619    for info in objects {
620        // Sidecar records (redaction + state-visibility) live outside
621        // `.heddle/objects/` so GC cannot touch them, and must not be
622        // folded into the content-addressed pack. They ship via the
623        // per-object transfer path instead; callers split them out before
624        // packing.
625        if !is_native_packable_object_type(info.obj_type) {
626            continue;
627        }
628        let object = load_object_data(store, &info.id, info.obj_type)?;
629        let pack_id = to_pack_object_id(&object.id, object.obj_type);
630        builder.add_id(pack_id, object.obj_type.pack_object_type()?, object.data);
631    }
632
633    let (pack_data, index_data, _) = builder.build()?;
634    Ok(NativePackBundle {
635        pack_data,
636        index_data,
637    })
638}
639
640fn sync_pack_compression() -> CompressionConfig {
641    CompressionConfig {
642        level: 1,
643        min_size: 1024,
644        max_delta_size: 0,
645        ..CompressionConfig::default()
646    }
647}
648
649pub fn install_received_pack(
650    store: &impl ObjectStore,
651    pack_data: &[u8],
652    index_data: &[u8],
653) -> Result<Vec<PackObjectId>> {
654    store
655        .install_pack(pack_data, index_data)
656        .map_err(ProtocolError::from)
657}
658
659pub fn next_pack_chunk(
660    data: &[u8],
661    chunk_size: usize,
662    chunk_index: usize,
663) -> Option<(usize, Vec<u8>, bool)> {
664    let (start, len) = crate::chunk_bounds(data.len(), chunk_size.max(1), chunk_index)?;
665    let is_final = start + len == data.len();
666    Some((start, data[start..start + len].to_vec(), is_final))
667}
668
669pub fn receive_pack_chunk(
670    state: &mut PackChunkState,
671    is_index: bool,
672    resume_offset: u64,
673    chunk_index: u32,
674    is_complete: bool,
675    data: &[u8],
676    is_final_chunk: bool,
677) -> Result<()> {
678    let max_bytes = if is_index {
679        MAX_RECEIVED_PACK_INDEX_SIZE
680    } else {
681        MAX_RECEIVED_PACK_SIZE
682    };
683    receive_pack_chunk_with_limit(
684        state,
685        is_index,
686        resume_offset,
687        chunk_index,
688        is_complete,
689        data,
690        is_final_chunk,
691        max_bytes,
692    )
693}
694
695#[allow(clippy::too_many_arguments)]
696fn receive_pack_chunk_with_limit(
697    state: &mut PackChunkState,
698    is_index: bool,
699    resume_offset: u64,
700    chunk_index: u32,
701    is_complete: bool,
702    data: &[u8],
703    is_final_chunk: bool,
704    max_bytes: u64,
705) -> Result<()> {
706    let (buffer, progress, complete) = if is_index {
707        (
708            &mut state.index_data,
709            &mut state.index_progress,
710            &mut state.index_complete,
711        )
712    } else {
713        (
714            &mut state.pack_data,
715            &mut state.pack_progress,
716            &mut state.pack_complete,
717        )
718    };
719
720    let next_progress = validate_pack_chunk(
721        *progress,
722        is_index,
723        resume_offset,
724        chunk_index,
725        data,
726        max_bytes,
727    )?;
728
729    buffer.extend_from_slice(data);
730    *progress = next_progress;
731    if is_final_chunk || is_complete {
732        *complete = true;
733    }
734    Ok(())
735}
736
737#[allow(clippy::too_many_arguments)]
738fn receive_pack_chunk_to_spool(
739    stream: &mut PackStreamSpool,
740    is_index: bool,
741    resume_offset: u64,
742    chunk_index: u32,
743    is_complete: bool,
744    data: &[u8],
745    is_final_chunk: bool,
746    max_bytes: u64,
747) -> Result<()> {
748    let next_progress = validate_pack_chunk(
749        stream.progress,
750        is_index,
751        resume_offset,
752        chunk_index,
753        data,
754        max_bytes,
755    )?;
756    stream.write_all(data)?;
757    stream.progress = next_progress;
758    if is_final_chunk || is_complete {
759        stream.complete = true;
760    }
761    Ok(())
762}
763
764fn validate_pack_chunk(
765    progress: (u64, u32),
766    is_index: bool,
767    resume_offset: u64,
768    chunk_index: u32,
769    data: &[u8],
770    max_bytes: u64,
771) -> Result<(u64, u32)> {
772    if resume_offset != progress.0 {
773        return Err(ProtocolError::InvalidState(format!(
774            "native pack chunk resume offset mismatch: expected {}, got {}",
775            progress.0, resume_offset
776        )));
777    }
778    if chunk_index != progress.1 {
779        return Err(ProtocolError::InvalidState(format!(
780            "native pack chunk index mismatch: expected {}, got {}",
781            progress.1, chunk_index
782        )));
783    }
784
785    let data_len = u64::try_from(data.len()).map_err(|_| {
786        ProtocolError::InvalidState("native pack chunk length does not fit in u64".to_string())
787    })?;
788    let next_offset = progress.0.checked_add(data_len).ok_or_else(|| {
789        ProtocolError::InvalidState("native pack chunk offset overflow".to_string())
790    })?;
791    if next_offset > max_bytes {
792        let stream_name = if is_index { "index" } else { "body" };
793        return Err(ProtocolError::InvalidState(format!(
794            "native pack {stream_name} exceeds receive size limit: {next_offset} bytes (max {max_bytes})"
795        )));
796    }
797    let next_chunk = progress.1.checked_add(1).ok_or_else(|| {
798        ProtocolError::InvalidState("native pack chunk index overflow".to_string())
799    })?;
800
801    Ok((next_offset, next_chunk))
802}
803
804pub(crate) fn unique_spool_dir(base: &Path) -> Result<PathBuf> {
805    let stamp = SystemTime::now()
806        .duration_since(UNIX_EPOCH)
807        .map_err(|err| {
808            ProtocolError::InvalidState(format!("system clock before UNIX epoch: {err}"))
809        })?
810        .as_nanos();
811    for attempt in 0..100u32 {
812        let dir = base.join(format!("pack-{}-{stamp}-{attempt}", std::process::id()));
813        match fs::create_dir(&dir) {
814            Ok(()) => return Ok(dir),
815            Err(err) if err.kind() == std::io::ErrorKind::AlreadyExists => continue,
816            Err(err) => return Err(ProtocolError::Io(err)),
817        }
818    }
819    Err(ProtocolError::InvalidState(
820        "failed to allocate native pack spool directory".to_string(),
821    ))
822}
823
824fn to_pack_object_id(id: &ObjectId, object_type: ObjectType) -> PackObjectId {
825    match (id, object_type) {
826        (ObjectId::Hash(hash), ObjectType::AnnotatedTag) => PackObjectId::AnnotatedTag(*hash),
827        (ObjectId::Hash(hash), _) => PackObjectId::Hash(*hash),
828        (ObjectId::StateId(state_id), _) => PackObjectId::StateId(*state_id),
829        (ObjectId::StateAttachment { id, .. }, _) => PackObjectId::Hash(*id.as_hash()),
830    }
831}
832
833#[cfg(test)]
834mod tests {
835    use objects::{
836        object::{AnnotatedTag, Blob, ContentHash, StateId},
837        store::{
838            CompressionConfig, FsStore, ObjectStore,
839            pack::{ObjectType as PackObjectType, PackBuilder, PackObjectId, PackReader},
840        },
841    };
842    use sley::ObjectFormat as GitObjectFormat;
843    use tempfile::TempDir;
844
845    use super::{
846        GitPackChunkState, GrowingPackChunkReader, MAX_RECEIVED_PACK_SIZE,
847        NativePackStreamingWriter, ObjectData, ObjectId, ObjectInfo, ObjectType, PackChunkSpool,
848        PackChunkState, PackFileChunkReader, build_native_pack, install_received_pack,
849        next_pack_chunk, receive_pack_chunk, receive_pack_chunk_with_limit,
850        reuse_native_pack_encoded_subset_in,
851    };
852
853    fn create_test_store() -> (TempDir, FsStore) {
854        let temp = TempDir::new().unwrap();
855        let store = FsStore::new(temp.path().join(".heddle"));
856        store.init().unwrap();
857        (temp, store)
858    }
859
860    fn hash(byte: u8) -> ContentHash {
861        ContentHash::from_bytes([byte; 32])
862    }
863
864    #[test]
865    fn encoded_snapshot_subset_is_wire_equivalent_without_local_artifacts_or_attachments() {
866        let source = TempDir::new().unwrap();
867        let spool = TempDir::new().unwrap();
868        let source_pack = source.path().join("snapshot.pack");
869        let source_index = source.path().join("snapshot.idx");
870        let blob = (
871            PackObjectId::Hash(hash(1)),
872            PackObjectType::Blob,
873            b"blob body".to_vec(),
874        );
875        let tree = (
876            PackObjectId::Hash(hash(2)),
877            PackObjectType::Tree,
878            b"tree body".to_vec(),
879        );
880        let state_id = StateId::from_bytes([3; 32]);
881        let state = (
882            PackObjectId::StateId(state_id),
883            PackObjectType::State,
884            b"state body".to_vec(),
885        );
886        let attachment_id = PackObjectId::Hash(hash(4));
887        let artifact_id = PackObjectId::Hash(hash(5));
888        let mut builder = PackBuilder::new(CompressionConfig {
889            max_delta_size: 0,
890            ..CompressionConfig::default()
891        });
892        for (id, kind, body) in [
893            blob.clone(),
894            tree.clone(),
895            state.clone(),
896            (
897                attachment_id,
898                PackObjectType::StateAttachment,
899                b"local attachment".to_vec(),
900            ),
901            (
902                artifact_id,
903                PackObjectType::SnapshotCommit,
904                b"local commit artifact".to_vec(),
905            ),
906        ] {
907            builder.add_id(id, kind, body);
908        }
909        let (pack, index, _) = builder.build().unwrap();
910        std::fs::write(&source_pack, pack).unwrap();
911        std::fs::write(&source_index, index).unwrap();
912
913        let wanted = vec![
914            ObjectInfo {
915                id: ObjectId::Hash(hash(1)),
916                obj_type: ObjectType::Blob,
917                size: blob.2.len() as u64,
918                delta_base: None,
919            },
920            ObjectInfo {
921                id: ObjectId::Hash(hash(2)),
922                obj_type: ObjectType::Tree,
923                size: tree.2.len() as u64,
924                delta_base: None,
925            },
926            ObjectInfo {
927                id: ObjectId::StateId(state_id),
928                obj_type: ObjectType::State,
929                size: state.2.len() as u64,
930                delta_base: None,
931            },
932        ];
933        let (bundle, stats) =
934            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &wanted)
935                .unwrap()
936                .expect("authoritative non-delta subset must be reusable");
937
938        assert_eq!(stats.object_count, wanted.len());
939        assert!(stats.encoded_bytes_copied > 0);
940        let reused = PackReader::open(&bundle.pack_path, &bundle.index_path).unwrap();
941        let mut reused_ids = reused.list_ids().unwrap();
942        reused_ids.sort();
943        let mut wanted_ids = vec![blob.0, tree.0, state.0];
944        wanted_ids.sort();
945        assert_eq!(reused_ids, wanted_ids);
946        assert!(!reused.has_object(&attachment_id).unwrap());
947        assert!(!reused.has_object(&artifact_id).unwrap());
948
949        for path in [&bundle.pack_path, &bundle.index_path] {
950            let expected_wire_bytes = std::fs::read(path).unwrap();
951            let mut chunk_reader = PackFileChunkReader::open(path, 7).unwrap();
952            let mut wire_bytes = Vec::new();
953            while let Some((offset, chunk_index, data, is_final)) =
954                chunk_reader.next_chunk().unwrap()
955            {
956                assert_eq!(offset as usize, wire_bytes.len());
957                assert_eq!(chunk_index as usize, wire_bytes.len() / 7);
958                wire_bytes.extend_from_slice(&data);
959                assert_eq!(is_final, wire_bytes.len() == expected_wire_bytes.len());
960            }
961            assert_eq!(wire_bytes, expected_wire_bytes);
962        }
963        for (id, _, expected) in [blob, tree, state] {
964            assert_eq!(reused.get_object(&id).unwrap().unwrap().1, expected);
965        }
966    }
967
968    #[test]
969    fn encoded_snapshot_subset_falls_back_for_mismatch_delta_or_attachment_request() {
970        let source = TempDir::new().unwrap();
971        let spool = TempDir::new().unwrap();
972        let source_pack = source.path().join("snapshot.pack");
973        let source_index = source.path().join("snapshot.idx");
974        let first = b"This is the base content. ".repeat(100);
975        let second = b"This is modified content. ".repeat(100);
976        let mut builder = PackBuilder::new(CompressionConfig::default());
977        builder.add(hash(10), PackObjectType::Blob, first.clone());
978        builder.add(hash(11), PackObjectType::Blob, second.clone());
979        let (pack, index, stats) = builder.build().unwrap();
980        assert!(stats.delta_count > 0, "fixture must contain a delta");
981        std::fs::write(&source_pack, pack).unwrap();
982        std::fs::write(&source_index, index).unwrap();
983        let delta_wants = [ObjectInfo {
984            id: ObjectId::Hash(hash(11)),
985            obj_type: ObjectType::Blob,
986            size: second.len() as u64,
987            delta_base: None,
988        }];
989        assert!(
990            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &delta_wants)
991                .unwrap()
992                .is_none()
993        );
994
995        let missing_wants = [ObjectInfo {
996            id: ObjectId::Hash(hash(12)),
997            obj_type: ObjectType::Blob,
998            size: 1,
999            delta_base: None,
1000        }];
1001        assert!(
1002            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &missing_wants)
1003                .unwrap()
1004                .is_none()
1005        );
1006
1007        let attachment_wants = [ObjectInfo {
1008            id: ObjectId::StateAttachment {
1009                state: StateId::from_bytes([13; 32]),
1010                id: objects::object::StateAttachmentId::from_hash(hash(14)),
1011                kind: objects::object::StateAttachmentKind::SemanticIndex,
1012            },
1013            obj_type: ObjectType::StateAttachment,
1014            size: 1,
1015            delta_base: None,
1016        }];
1017        assert!(
1018            reuse_native_pack_encoded_subset_in(spool.path(), &source_pack, &attachment_wants)
1019                .unwrap()
1020                .is_none()
1021        );
1022    }
1023
1024    #[test]
1025    fn receive_pack_chunk_rejects_cumulative_size_over_limit_before_buffering() {
1026        let mut state = PackChunkState::default();
1027
1028        receive_pack_chunk_with_limit(&mut state, false, 0, 0, false, b"abcd", false, 8).unwrap();
1029        receive_pack_chunk_with_limit(&mut state, false, 4, 1, false, b"efgh", false, 8).unwrap();
1030
1031        let error = receive_pack_chunk_with_limit(&mut state, false, 8, 2, false, b"i", false, 8)
1032            .unwrap_err();
1033
1034        assert_eq!(state.pack_data, b"abcdefgh");
1035        assert!(
1036            error
1037                .to_string()
1038                .contains("native pack body exceeds receive size limit")
1039        );
1040        assert!(error.to_string().contains("9 bytes (max 8)"));
1041    }
1042
1043    #[test]
1044    fn receive_pack_chunk_checks_production_limit_before_extending_buffer() {
1045        let mut state = PackChunkState {
1046            pack_progress: (MAX_RECEIVED_PACK_SIZE - 1, 0),
1047            ..PackChunkState::default()
1048        };
1049
1050        let error = receive_pack_chunk(
1051            &mut state,
1052            false,
1053            MAX_RECEIVED_PACK_SIZE - 1,
1054            0,
1055            false,
1056            b"xx",
1057            false,
1058        )
1059        .unwrap_err();
1060
1061        assert!(state.pack_data.is_empty());
1062        assert!(
1063            error
1064                .to_string()
1065                .contains("native pack body exceeds receive size limit")
1066        );
1067    }
1068
1069    #[test]
1070    fn receive_pack_chunk_rejects_resume_offset_mismatch_before_buffering() {
1071        let mut state = PackChunkState::default();
1072
1073        let error =
1074            receive_pack_chunk(&mut state, false, 1, 0, false, b"late chunk", false).unwrap_err();
1075
1076        assert!(state.pack_data.is_empty());
1077        assert!(
1078            error
1079                .to_string()
1080                .contains("native pack chunk resume offset mismatch: expected 0, got 1")
1081        );
1082    }
1083
1084    #[test]
1085    fn receive_pack_chunk_rejects_chunk_index_mismatch_before_buffering() {
1086        let mut state = PackChunkState::default();
1087
1088        receive_pack_chunk(&mut state, false, 0, 0, false, b"abc", false).unwrap();
1089        let error = receive_pack_chunk(&mut state, false, 3, 2, false, b"def", false).unwrap_err();
1090
1091        assert_eq!(state.pack_data, b"abc");
1092        assert!(
1093            error
1094                .to_string()
1095                .contains("native pack chunk index mismatch: expected 1, got 2")
1096        );
1097    }
1098
1099    #[test]
1100    fn git_pack_chunk_state_requires_ordered_chunks_and_final_size() {
1101        let mut state = GitPackChunkState::default();
1102
1103        assert!(
1104            state
1105                .receive_chunk("git-pack:test", 0, 0, false, 8, b"abcd")
1106                .unwrap()
1107                .is_none()
1108        );
1109        let error = state
1110            .receive_chunk("git-pack:test", 4, 2, true, 8, b"efgh")
1111            .unwrap_err();
1112
1113        assert!(
1114            error
1115                .to_string()
1116                .contains("Git pack chunk index mismatch: expected 1, got 2")
1117        );
1118        assert!(state.ensure_idle().is_err());
1119
1120        let mut state = GitPackChunkState::default();
1121        state
1122            .receive_chunk("git-pack:test", 0, 0, false, 8, b"abcd")
1123            .unwrap();
1124        let complete = state
1125            .receive_chunk("git-pack:test", 4, 1, true, 8, b"efgh")
1126            .unwrap()
1127            .unwrap();
1128
1129        assert_eq!(complete, b"abcdefgh");
1130        assert!(state.ensure_idle().is_ok());
1131    }
1132
1133    #[test]
1134    fn receive_pack_chunk_accepts_completion_flags_for_pack_and_index() {
1135        let mut state = PackChunkState::default();
1136
1137        receive_pack_chunk(&mut state, false, 0, 0, true, b"pack-body", false).unwrap();
1138        assert!(!state.is_complete());
1139        receive_pack_chunk(&mut state, true, 0, 0, false, b"pack-index", true).unwrap();
1140
1141        assert!(state.is_complete());
1142        assert_eq!(state.pack_data, b"pack-body");
1143        assert_eq!(state.index_data, b"pack-index");
1144    }
1145
1146    #[test]
1147    fn normal_size_native_pack_receives_and_installs() {
1148        let (_source_temp, source_store) = create_test_store();
1149        let (_dest_temp, dest_store) = create_test_store();
1150        let blob = Blob::from("native pack receive regression");
1151        let hash = source_store.put_blob(&blob).unwrap();
1152        let bundle = build_native_pack(
1153            &source_store,
1154            &[ObjectInfo {
1155                id: ObjectId::Hash(hash),
1156                obj_type: ObjectType::Blob,
1157                size: blob.size() as u64,
1158                delta_base: None,
1159            }],
1160        )
1161        .unwrap();
1162
1163        let mut state = PackChunkState::default();
1164        let mut chunk_index = 0usize;
1165        while let Some((start, data, is_final)) = next_pack_chunk(&bundle.pack_data, 7, chunk_index)
1166        {
1167            receive_pack_chunk(
1168                &mut state,
1169                false,
1170                start as u64,
1171                chunk_index as u32,
1172                is_final,
1173                &data,
1174                is_final,
1175            )
1176            .unwrap();
1177            chunk_index += 1;
1178        }
1179
1180        let mut index_chunk = 0usize;
1181        while let Some((start, data, is_final)) =
1182            next_pack_chunk(&bundle.index_data, 5, index_chunk)
1183        {
1184            receive_pack_chunk(
1185                &mut state,
1186                true,
1187                start as u64,
1188                index_chunk as u32,
1189                is_final,
1190                &data,
1191                is_final,
1192            )
1193            .unwrap();
1194            index_chunk += 1;
1195        }
1196
1197        assert!(state.is_complete());
1198        assert_eq!(state.pack_data, bundle.pack_data);
1199        assert_eq!(state.index_data, bundle.index_data);
1200
1201        let installed_ids =
1202            install_received_pack(&dest_store, &state.pack_data, &state.index_data).unwrap();
1203
1204        assert_eq!(installed_ids, vec![PackObjectId::Hash(hash)]);
1205        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1206        assert_eq!(installed_blob.content(), blob.content());
1207    }
1208
1209    #[test]
1210    fn native_pack_transfers_first_class_annotated_tag() {
1211        let (_source_temp, source_store) = create_test_store();
1212        let (_dest_temp, dest_store) = create_test_store();
1213        let tag = AnnotatedTag::new(
1214            GitObjectFormat::Sha1,
1215            b"object 1111111111111111111111111111111111111111\ntype commit\ntag v1\ntagger Test <test@example.com> 1700000000 +0100\n\nrelease\n".to_vec(),
1216            None,
1217            None,
1218        )
1219        .unwrap();
1220        let hash = source_store.put_annotated_tag(&tag).unwrap();
1221        let bundle = build_native_pack(
1222            &source_store,
1223            &[ObjectInfo {
1224                id: ObjectId::Hash(hash),
1225                obj_type: ObjectType::AnnotatedTag,
1226                size: tag.encode_current_msgpack().len() as u64,
1227                delta_base: None,
1228            }],
1229        )
1230        .unwrap();
1231
1232        let installed = install_received_pack(&dest_store, &bundle.pack_data, &bundle.index_data)
1233            .expect("install annotated-tag native pack");
1234
1235        assert_eq!(installed, vec![PackObjectId::AnnotatedTag(hash)]);
1236        assert_eq!(dest_store.get_annotated_tag(&hash).unwrap(), Some(tag));
1237    }
1238
1239    #[test]
1240    fn normal_size_native_pack_spools_and_installs() {
1241        let (_source_temp, source_store) = create_test_store();
1242        let (dest_temp, dest_store) = create_test_store();
1243        let blob = Blob::from("native pack spooled receive regression");
1244        let hash = source_store.put_blob(&blob).unwrap();
1245        let bundle = build_native_pack(
1246            &source_store,
1247            &[ObjectInfo {
1248                id: ObjectId::Hash(hash),
1249                obj_type: ObjectType::Blob,
1250                size: blob.size() as u64,
1251                delta_base: None,
1252            }],
1253        )
1254        .unwrap();
1255
1256        let mut spool = PackChunkSpool::new_in(dest_temp.path()).unwrap();
1257        let mut chunk_index = 0usize;
1258        while let Some((start, data, is_final)) = next_pack_chunk(&bundle.pack_data, 7, chunk_index)
1259        {
1260            spool
1261                .receive_chunk(
1262                    false,
1263                    start as u64,
1264                    chunk_index as u32,
1265                    is_final,
1266                    &data,
1267                    is_final,
1268                )
1269                .unwrap();
1270            chunk_index += 1;
1271        }
1272
1273        let mut index_chunk = 0usize;
1274        while let Some((start, data, is_final)) =
1275            next_pack_chunk(&bundle.index_data, 5, index_chunk)
1276        {
1277            spool
1278                .receive_chunk(
1279                    true,
1280                    start as u64,
1281                    index_chunk as u32,
1282                    is_final,
1283                    &data,
1284                    is_final,
1285                )
1286                .unwrap();
1287            index_chunk += 1;
1288        }
1289
1290        assert!(spool.is_complete());
1291        let installed_ids = spool.install_into(&dest_store).unwrap();
1292
1293        assert_eq!(installed_ids, vec![PackObjectId::Hash(hash)]);
1294        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1295        assert_eq!(installed_blob.content(), blob.content());
1296    }
1297
1298    #[test]
1299    fn native_pack_streaming_writer_drains_growing_pack_and_installs() {
1300        let (source_temp, source_store) = create_test_store();
1301        let (dest_temp, dest_store) = create_test_store();
1302        let blob = Blob::from("native pack growing stream regression");
1303        let hash = source_store.put_blob(&blob).unwrap();
1304        let large_blob = Blob::from_slice(&vec![b'z'; 4096]);
1305        let large_hash = source_store.put_blob(&large_blob).unwrap();
1306
1307        let mut writer = NativePackStreamingWriter::new_in(source_temp.path(), 2).unwrap();
1308        let mut pack_reader = GrowingPackChunkReader::open(writer.pack_path(), 31).unwrap();
1309        let mut spool = PackChunkSpool::new_in(dest_temp.path()).unwrap();
1310        let mut saw_interleaved_pack_chunk = false;
1311
1312        for (id, obj_type, data) in [
1313            (
1314                ObjectId::Hash(hash),
1315                ObjectType::Blob,
1316                blob.content().to_vec(),
1317            ),
1318            (
1319                ObjectId::Hash(large_hash),
1320                ObjectType::Blob,
1321                large_blob.content().to_vec(),
1322            ),
1323        ] {
1324            writer
1325                .add_object_data(ObjectData {
1326                    id,
1327                    obj_type,
1328                    data,
1329                    is_delta: false,
1330                })
1331                .unwrap();
1332            writer.flush_pack().unwrap();
1333            while let Some((offset, chunk_index, data, is_final)) =
1334                pack_reader.next_available_chunk(false).unwrap()
1335            {
1336                assert!(
1337                    !is_final,
1338                    "pre-final growing pack drain must not mark chunks final"
1339                );
1340                saw_interleaved_pack_chunk = true;
1341                spool
1342                    .receive_chunk(false, offset, chunk_index, false, &data, false)
1343                    .unwrap();
1344            }
1345        }
1346
1347        let bundle = writer.finish().unwrap();
1348        let mut saw_final_pack_chunk = false;
1349        while let Some((offset, chunk_index, data, is_final)) =
1350            pack_reader.next_available_chunk(true).unwrap()
1351        {
1352            saw_final_pack_chunk |= is_final;
1353            spool
1354                .receive_chunk(false, offset, chunk_index, is_final, &data, is_final)
1355                .unwrap();
1356        }
1357
1358        let mut index_reader = PackFileChunkReader::open(&bundle.index_path, 17).unwrap();
1359        while let Some((offset, chunk_index, data, is_final)) = index_reader.next_chunk().unwrap() {
1360            spool
1361                .receive_chunk(true, offset, chunk_index, is_final, &data, is_final)
1362                .unwrap();
1363        }
1364
1365        assert!(
1366            saw_interleaved_pack_chunk,
1367            "expected at least one pack chunk before finalize"
1368        );
1369        assert!(
1370            saw_final_pack_chunk,
1371            "expected final pack chunk after finish"
1372        );
1373        assert!(spool.is_complete());
1374        let mut installed_ids = spool.install_into(&dest_store).unwrap();
1375        let mut expected_ids = vec![PackObjectId::Hash(hash), PackObjectId::Hash(large_hash)];
1376        installed_ids.sort();
1377        expected_ids.sort();
1378
1379        assert_eq!(installed_ids, expected_ids);
1380        let installed_blob = dest_store.get_blob(&hash).unwrap().unwrap();
1381        assert_eq!(installed_blob.content(), blob.content());
1382        let installed_large_blob = dest_store.get_blob(&large_hash).unwrap().unwrap();
1383        assert_eq!(installed_large_blob.content(), large_blob.content());
1384    }
1385}