1use 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
18pub const MAX_RECEIVED_PACK_SIZE: u64 = 2 * 1024 * 1024 * 1024;
29
30pub const MAX_RECEIVED_PACK_INDEX_SIZE: u64 = 256 * 1024 * 1024;
36
37pub 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
66pub 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 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}