Skip to main content

wire/
provider_pack.rs

1// SPDX-License-Identifier: Apache-2.0
2use std::{
3    collections::HashSet,
4    fs::{self, File, OpenOptions},
5    io::Write,
6    path::{Path, PathBuf},
7    sync::{Arc, Mutex},
8};
9
10use objects::store::{
11    ObjectStore,
12    pack::{PackContainerSpec, PackIndex, PackObjectId, PackReader, verify_container},
13};
14
15use crate::{
16    MAX_RECEIVED_PACK_SIZE, NativePackBundle, ProtocolError, Result, native_pack::unique_spool_dir,
17};
18
19const PACK_HEADER_LEN: usize = 16;
20const PACK_TRAILER_LEN: usize = 32;
21const PACK_SPEC: PackContainerSpec = PackContainerSpec {
22    magic: b"LMPK",
23    version: 3,
24};
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub struct ProviderPackIndexEntry {
28    pub id: PackObjectId,
29    pub output_offset: u64,
30}
31
32#[derive(Debug, Clone, PartialEq, Eq)]
33pub struct ProviderPackExtent {
34    pub output_offset: u64,
35    pub length: u64,
36    pub digest: [u8; 32],
37    pub objects: Vec<ProviderPackIndexEntry>,
38}
39
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct ProviderPackManifest {
42    pub header: [u8; PACK_HEADER_LEN],
43    pub output_pack_length: u64,
44    pub extents: Vec<ProviderPackExtent>,
45}
46
47#[derive(Debug)]
48pub struct ProviderPackBundle {
49    pub pack: NativePackBundle,
50    pub trailer_digest: [u8; 32],
51}
52
53/// A bounded-memory positional writer for one validated provider pack plan.
54#[derive(Clone, Debug)]
55pub struct ProviderPackWriter {
56    file: Arc<File>,
57    ranges: Arc<Vec<(u64, u64)>>,
58    verified: Arc<Mutex<Vec<bool>>>,
59}
60
61/// A pre-sized provider pack spool that owns all partial transfer state.
62#[derive(Debug)]
63pub struct ProviderPackSpool {
64    dir: Option<PathBuf>,
65    pack_path: PathBuf,
66    index_path: PathBuf,
67    file: Arc<File>,
68    manifest: ProviderPackManifest,
69    verified: Arc<Mutex<Vec<bool>>>,
70}
71
72/// A completely covered and validated provider pack ready for atomic install.
73#[derive(Debug)]
74pub struct CompletedProviderPack {
75    dir: PathBuf,
76    pack_path: PathBuf,
77    index_path: PathBuf,
78    pub trailer_digest: [u8; 32],
79}
80
81impl ProviderPackSpool {
82    /// Create an exact-length spool and write the validated virtual pack header.
83    pub fn new_in(root: &Path, manifest: ProviderPackManifest) -> Result<Self> {
84        validate_manifest(&manifest)?;
85        let base = root.join("transfer-spool");
86        fs::create_dir_all(&base)?;
87        let dir = unique_spool_dir(&base)?;
88        let pack_path = dir.join("provider.pack");
89        let index_path = dir.join("provider.idx");
90        let create_result = (|| -> Result<File> {
91            let file = OpenOptions::new()
92                .read(true)
93                .write(true)
94                .create_new(true)
95                .open(&pack_path)?;
96            file.set_len(manifest.output_pack_length)?;
97            write_all_at(&file, &manifest.header, 0)?;
98            Ok(file)
99        })();
100        let file = match create_result {
101            Ok(file) => file,
102            Err(error) => {
103                let _ = fs::remove_dir_all(&dir);
104                return Err(error);
105            }
106        };
107        let verified = Arc::new(Mutex::new(vec![false; manifest.extents.len()]));
108        Ok(Self {
109            dir: Some(dir),
110            pack_path,
111            index_path,
112            file: Arc::new(file),
113            manifest,
114            verified,
115        })
116    }
117
118    /// Obtain a cloneable positional writer for concurrent extent streams.
119    pub fn writer(&self) -> ProviderPackWriter {
120        ProviderPackWriter {
121            file: Arc::clone(&self.file),
122            ranges: Arc::new(
123                self.manifest
124                    .extents
125                    .iter()
126                    .map(|extent| (extent.output_offset, extent.length))
127                    .collect(),
128            ),
129            verified: Arc::clone(&self.verified),
130        }
131    }
132
133    /// Finalize the trailer and index after every extent has verified exactly once.
134    pub fn finish(mut self) -> Result<CompletedProviderPack> {
135        let verified = self.verified.lock().map_err(|_| {
136            ProtocolError::InvalidState("provider spool verification lock poisoned".to_string())
137        })?;
138        if verified.iter().any(|complete| !complete) {
139            return Err(ProtocolError::InvalidState(
140                "provider spool does not have complete verified coverage".to_string(),
141            ));
142        }
143        drop(verified);
144
145        let body_end = self.manifest.output_pack_length - PACK_TRAILER_LEN as u64;
146        let trailer_digest = hash_file_prefix(&self.file, body_end)?;
147        write_all_at(&self.file, &trailer_digest, body_end)?;
148        self.file.sync_all()?;
149        write_provider_index(&self.index_path, &self.manifest)?;
150
151        let reader = PackReader::open(&self.pack_path, &self.index_path)?;
152        let expected_objects = self
153            .manifest
154            .extents
155            .iter()
156            .map(|extent| extent.objects.len())
157            .sum::<usize>();
158        if reader.list_ids().len() != expected_objects {
159            return Err(ProtocolError::InvalidState(
160                "provider pack index does not match its manifest".to_string(),
161            ));
162        }
163        drop(reader);
164
165        let dir = self.dir.take().ok_or_else(|| {
166            ProtocolError::InvalidState("provider spool directory is missing".to_string())
167        })?;
168        Ok(CompletedProviderPack {
169            dir,
170            pack_path: self.pack_path.clone(),
171            index_path: self.index_path.clone(),
172            trailer_digest,
173        })
174    }
175}
176
177impl ProviderPackWriter {
178    /// Write one chunk at its manifest-assigned extent-relative position.
179    pub fn write_extent_chunk(
180        &self,
181        extent_index: usize,
182        relative_offset: u64,
183        data: &[u8],
184    ) -> Result<()> {
185        let (output_offset, extent_len) =
186            self.ranges.get(extent_index).copied().ok_or_else(|| {
187                ProtocolError::InvalidState("provider extent index is out of range".to_string())
188            })?;
189        let data_len = u64::try_from(data.len()).map_err(|_| {
190            ProtocolError::InvalidState("provider chunk length exceeds u64".to_string())
191        })?;
192        let relative_end = relative_offset.checked_add(data_len).ok_or_else(|| {
193            ProtocolError::InvalidState("provider extent write offset overflows".to_string())
194        })?;
195        if relative_end > extent_len {
196            return Err(ProtocolError::InvalidState(
197                "provider extent write exceeds its planned range".to_string(),
198            ));
199        }
200        let absolute_offset = output_offset.checked_add(relative_offset).ok_or_else(|| {
201            ProtocolError::InvalidState("provider spool write offset overflows".to_string())
202        })?;
203        write_all_at(&self.file, data, absolute_offset)
204    }
205
206    /// Rehash a retained prefix without retaining the extent body in memory.
207    pub fn hash_extent_prefix(
208        &self,
209        extent_index: usize,
210        prefix_len: u64,
211        hasher: &mut blake3::Hasher,
212    ) -> Result<()> {
213        let (output_offset, extent_len) =
214            self.ranges.get(extent_index).copied().ok_or_else(|| {
215                ProtocolError::InvalidState("provider extent index is out of range".to_string())
216            })?;
217        if prefix_len > extent_len {
218            return Err(ProtocolError::InvalidState(
219                "provider retained prefix exceeds its planned extent".to_string(),
220            ));
221        }
222        let mut buffer = [0_u8; 64 * 1024];
223        let mut read = 0_u64;
224        while read < prefix_len {
225            let remaining = prefix_len - read;
226            let length = usize::try_from(remaining.min(buffer.len() as u64)).map_err(|_| {
227                ProtocolError::InvalidState("provider prefix length exceeds usize".to_string())
228            })?;
229            let offset = output_offset.checked_add(read).ok_or_else(|| {
230                ProtocolError::InvalidState("provider prefix read offset overflows".to_string())
231            })?;
232            read_exact_at(&self.file, &mut buffer[..length], offset)?;
233            hasher.update(&buffer[..length]);
234            read += length as u64;
235        }
236        Ok(())
237    }
238
239    /// Mark one fully length- and digest-verified extent complete exactly once.
240    pub fn mark_verified(&self, extent_index: usize) -> Result<()> {
241        let mut verified = self.verified.lock().map_err(|_| {
242            ProtocolError::InvalidState("provider spool verification lock poisoned".to_string())
243        })?;
244        let complete = verified.get_mut(extent_index).ok_or_else(|| {
245            ProtocolError::InvalidState("provider extent index is out of range".to_string())
246        })?;
247        if *complete {
248            return Err(ProtocolError::InvalidState(
249                "provider extent completed more than once".to_string(),
250            ));
251        }
252        *complete = true;
253        Ok(())
254    }
255}
256
257impl CompletedProviderPack {
258    /// Atomically install a fully validated provider pack into the object store.
259    pub fn install_into(&mut self, store: &impl ObjectStore) -> Result<Vec<PackObjectId>> {
260        store
261            .install_pack_streaming(&self.pack_path, &self.index_path)
262            .map_err(ProtocolError::from)
263    }
264}
265
266impl Drop for ProviderPackSpool {
267    fn drop(&mut self) {
268        if let Some(dir) = self.dir.as_ref() {
269            let _ = fs::remove_dir_all(dir);
270        }
271    }
272}
273
274impl Drop for CompletedProviderPack {
275    fn drop(&mut self) {
276        let _ = fs::remove_dir_all(&self.dir);
277    }
278}
279
280fn write_provider_index(path: &Path, manifest: &ProviderPackManifest) -> Result<()> {
281    let mut index = PackIndex::new();
282    for extent in &manifest.extents {
283        for object in &extent.objects {
284            index.add(object.id, object.output_offset);
285        }
286    }
287    index.sort();
288    let mut file = OpenOptions::new().write(true).create_new(true).open(path)?;
289    file.write_all(&index.to_bytes())?;
290    file.flush()?;
291    file.sync_all()?;
292    Ok(())
293}
294
295fn hash_file_prefix(file: &File, length: u64) -> Result<[u8; 32]> {
296    let mut hasher = blake3::Hasher::new();
297    let mut buffer = [0_u8; 64 * 1024];
298    let mut offset = 0_u64;
299    while offset < length {
300        let remaining = length - offset;
301        let read_len = usize::try_from(remaining.min(buffer.len() as u64)).map_err(|_| {
302            ProtocolError::InvalidState("provider pack hash length exceeds usize".to_string())
303        })?;
304        read_exact_at(file, &mut buffer[..read_len], offset)?;
305        hasher.update(&buffer[..read_len]);
306        offset += read_len as u64;
307    }
308    Ok(*hasher.finalize().as_bytes())
309}
310
311#[cfg(unix)]
312fn write_all_at(file: &File, mut data: &[u8], mut offset: u64) -> Result<()> {
313    use std::os::unix::fs::FileExt;
314
315    while !data.is_empty() {
316        let written = file.write_at(data, offset)?;
317        if written == 0 {
318            return Err(ProtocolError::Io(std::io::Error::new(
319                std::io::ErrorKind::WriteZero,
320                "provider positional spool write returned zero",
321            )));
322        }
323        data = &data[written..];
324        offset = offset.checked_add(written as u64).ok_or_else(|| {
325            ProtocolError::InvalidState("provider positional write offset overflows".to_string())
326        })?;
327    }
328    Ok(())
329}
330
331#[cfg(windows)]
332fn write_all_at(file: &File, mut data: &[u8], mut offset: u64) -> Result<()> {
333    use std::os::windows::fs::FileExt;
334
335    while !data.is_empty() {
336        let written = file.seek_write(data, offset)?;
337        if written == 0 {
338            return Err(ProtocolError::Io(std::io::Error::new(
339                std::io::ErrorKind::WriteZero,
340                "provider positional spool write returned zero",
341            )));
342        }
343        data = &data[written..];
344        offset = offset.checked_add(written as u64).ok_or_else(|| {
345            ProtocolError::InvalidState("provider positional write offset overflows".to_string())
346        })?;
347    }
348    Ok(())
349}
350
351#[cfg(unix)]
352fn read_exact_at(file: &File, mut data: &mut [u8], mut offset: u64) -> Result<()> {
353    use std::os::unix::fs::FileExt;
354
355    while !data.is_empty() {
356        let read = file.read_at(data, offset)?;
357        if read == 0 {
358            return Err(ProtocolError::Io(std::io::Error::new(
359                std::io::ErrorKind::UnexpectedEof,
360                "provider positional spool read ended early",
361            )));
362        }
363        data = &mut data[read..];
364        offset = offset.checked_add(read as u64).ok_or_else(|| {
365            ProtocolError::InvalidState("provider positional read offset overflows".to_string())
366        })?;
367    }
368    Ok(())
369}
370
371#[cfg(windows)]
372fn read_exact_at(file: &File, mut data: &mut [u8], mut offset: u64) -> Result<()> {
373    use std::os::windows::fs::FileExt;
374
375    while !data.is_empty() {
376        let read = file.seek_read(data, offset)?;
377        if read == 0 {
378            return Err(ProtocolError::Io(std::io::Error::new(
379                std::io::ErrorKind::UnexpectedEof,
380                "provider positional spool read ended early",
381            )));
382        }
383        data = &mut data[read..];
384        offset = offset.checked_add(read as u64).ok_or_else(|| {
385            ProtocolError::InvalidState("provider positional read offset overflows".to_string())
386        })?;
387    }
388    Ok(())
389}
390
391/// Assemble and verify one virtual native pack from provider extent bodies.
392///
393/// `extent_bodies` uses the same order as `manifest.extents`. The manifest may
394/// arrive in any order, but its output layout must cover every byte between the
395/// 16-byte header and 32-byte trailer exactly once.
396pub fn assemble_provider_pack(
397    manifest: &ProviderPackManifest,
398    extent_bodies: &[Vec<u8>],
399) -> Result<ProviderPackBundle> {
400    validate_manifest(manifest)?;
401    if extent_bodies.len() != manifest.extents.len() {
402        return Err(ProtocolError::InvalidState(format!(
403            "provider extent count mismatch: expected {}, got {}",
404            manifest.extents.len(),
405            extent_bodies.len()
406        )));
407    }
408
409    let mut order = (0..manifest.extents.len()).collect::<Vec<_>>();
410    order.sort_unstable_by_key(|index| manifest.extents[*index].output_offset);
411    let output_len = usize::try_from(manifest.output_pack_length).map_err(|_| {
412        ProtocolError::InvalidState("provider output pack exceeds this platform".to_string())
413    })?;
414    let mut pack_data = Vec::with_capacity(output_len);
415    pack_data.extend_from_slice(&manifest.header);
416    for index in order {
417        let extent = &manifest.extents[index];
418        let body = &extent_bodies[index];
419        let expected_len = usize::try_from(extent.length).map_err(|_| {
420            ProtocolError::InvalidState("provider extent exceeds this platform".to_string())
421        })?;
422        if body.len() != expected_len {
423            return Err(ProtocolError::InvalidState(format!(
424                "provider extent length mismatch: expected {}, got {}",
425                extent.length,
426                body.len()
427            )));
428        }
429        if blake3::hash(body).as_bytes() != &extent.digest {
430            return Err(ProtocolError::InvalidState(
431                "provider extent digest mismatch".to_string(),
432            ));
433        }
434        pack_data.extend_from_slice(body);
435    }
436
437    let trailer_digest = *blake3::hash(&pack_data).as_bytes();
438    pack_data.extend_from_slice(&trailer_digest);
439    if pack_data.len() != output_len {
440        return Err(ProtocolError::InvalidState(format!(
441            "provider output pack length mismatch: expected {}, got {}",
442            manifest.output_pack_length,
443            pack_data.len()
444        )));
445    }
446    verify_container(&pack_data, PACK_SPEC).map_err(ProtocolError::from)?;
447
448    let mut index = PackIndex::new();
449    for extent in &manifest.extents {
450        for object in &extent.objects {
451            index.add(object.id, object.output_offset);
452        }
453    }
454    index.sort();
455
456    Ok(ProviderPackBundle {
457        pack: NativePackBundle {
458            pack_data,
459            index_data: index.to_bytes(),
460        },
461        trailer_digest,
462    })
463}
464
465fn validate_manifest(manifest: &ProviderPackManifest) -> Result<()> {
466    if manifest.output_pack_length > MAX_RECEIVED_PACK_SIZE
467        || manifest.output_pack_length < (PACK_HEADER_LEN + PACK_TRAILER_LEN) as u64
468    {
469        return Err(ProtocolError::InvalidState(
470            "provider output pack length is invalid".to_string(),
471        ));
472    }
473    if &manifest.header[..4] != PACK_SPEC.magic
474        || u32::from_be_bytes(manifest.header[4..8].try_into().map_err(|_| {
475            ProtocolError::InvalidState("provider pack header is truncated".to_string())
476        })?) != PACK_SPEC.version
477    {
478        return Err(ProtocolError::InvalidState(
479            "provider pack header has invalid magic or version".to_string(),
480        ));
481    }
482    let object_count = u64::from_be_bytes(manifest.header[8..16].try_into().map_err(|_| {
483        ProtocolError::InvalidState("provider pack header is truncated".to_string())
484    })?);
485    let expected_body_end = manifest.output_pack_length - PACK_TRAILER_LEN as u64;
486    let mut order = manifest.extents.iter().collect::<Vec<_>>();
487    order.sort_unstable_by_key(|extent| extent.output_offset);
488    let mut next_offset = PACK_HEADER_LEN as u64;
489    let mut ids = HashSet::new();
490    let mut object_offsets = HashSet::new();
491    let mut actual_object_count = 0_u64;
492
493    for extent in order {
494        if extent.length == 0 || extent.output_offset != next_offset {
495            return Err(ProtocolError::InvalidState(
496                "provider extents do not exactly cover the virtual pack body".to_string(),
497            ));
498        }
499        next_offset = extent
500            .output_offset
501            .checked_add(extent.length)
502            .ok_or_else(|| {
503                ProtocolError::InvalidState("provider extent output range overflows".to_string())
504            })?;
505        if next_offset > expected_body_end {
506            return Err(ProtocolError::InvalidState(
507                "provider extent exceeds the virtual pack body".to_string(),
508            ));
509        }
510        for object in &extent.objects {
511            if object.output_offset < extent.output_offset
512                || object.output_offset >= next_offset
513                || !ids.insert(object.id)
514                || !object_offsets.insert(object.output_offset)
515            {
516                return Err(ProtocolError::InvalidState(
517                    "provider object index is outside its extent or duplicated".to_string(),
518                ));
519            }
520            actual_object_count = actual_object_count.checked_add(1).ok_or_else(|| {
521                ProtocolError::InvalidState("provider object count overflows".to_string())
522            })?;
523        }
524    }
525    if next_offset != expected_body_end || actual_object_count != object_count {
526        return Err(ProtocolError::InvalidState(
527            "provider manifest body or object count is incomplete".to_string(),
528        ));
529    }
530    Ok(())
531}
532
533#[cfg(test)]
534mod tests {
535    use objects::{
536        object::ContentHash,
537        store::{
538            CompressionConfig,
539            pack::{ObjectType, PackBuilder, PackIndex, PackObjectId},
540        },
541    };
542
543    use super::*;
544
545    fn source_pack() -> (Vec<u8>, Vec<u8>, Vec<PackObjectId>) {
546        let ids = vec![
547            PackObjectId::Hash(ContentHash::from_bytes([1; 32])),
548            PackObjectId::Hash(ContentHash::from_bytes([2; 32])),
549        ];
550        let mut builder = PackBuilder::new(CompressionConfig {
551            enabled: false,
552            ..CompressionConfig::default()
553        });
554        builder.add_id(ids[0], ObjectType::Blob, b"provider-one".to_vec());
555        builder.add_id(ids[1], ObjectType::Blob, b"provider-two".to_vec());
556        let (pack, index, _) = builder.build().unwrap();
557        (pack, index, ids)
558    }
559
560    fn split_manifest() -> (ProviderPackManifest, Vec<Vec<u8>>, Vec<u8>, Vec<u8>) {
561        let (pack, index, ids) = source_pack();
562        let parsed_index = PackIndex::from_bytes(&index).unwrap();
563        let first = parsed_index.find(&ids[0]).unwrap();
564        let second = parsed_index.find(&ids[1]).unwrap();
565        let (first_id, first_offset, second_id, second_offset) = if first < second {
566            (ids[0], first, ids[1], second)
567        } else {
568            (ids[1], second, ids[0], first)
569        };
570        let body_end = pack.len() - PACK_TRAILER_LEN;
571        let first_body = pack[first_offset as usize..second_offset as usize].to_vec();
572        let second_body = pack[second_offset as usize..body_end].to_vec();
573        let manifest = ProviderPackManifest {
574            header: pack[..PACK_HEADER_LEN].try_into().unwrap(),
575            output_pack_length: pack.len() as u64,
576            extents: vec![
577                ProviderPackExtent {
578                    output_offset: first_offset,
579                    length: first_body.len() as u64,
580                    digest: *blake3::hash(&first_body).as_bytes(),
581                    objects: vec![ProviderPackIndexEntry {
582                        id: first_id,
583                        output_offset: first_offset,
584                    }],
585                },
586                ProviderPackExtent {
587                    output_offset: second_offset,
588                    length: second_body.len() as u64,
589                    digest: *blake3::hash(&second_body).as_bytes(),
590                    objects: vec![ProviderPackIndexEntry {
591                        id: second_id,
592                        output_offset: second_offset,
593                    }],
594                },
595            ],
596        };
597        (manifest, vec![first_body, second_body], pack, index)
598    }
599
600    #[test]
601    fn provider_and_ordinary_pack_results_are_byte_identical() {
602        let (manifest, bodies, source_pack, source_index) = split_manifest();
603
604        let assembled = assemble_provider_pack(&manifest, &bodies).unwrap();
605
606        assert_eq!(assembled.pack.pack_data, source_pack);
607        assert_eq!(assembled.pack.index_data, source_index);
608        let source_digest = blake3::Hash::from_bytes(
609            source_pack[source_pack.len() - PACK_TRAILER_LEN..]
610                .try_into()
611                .unwrap(),
612        );
613        println!(
614            "byte_identical provider_digest={} weft_digest={} pack_bytes={} index_bytes={} identical=true",
615            blake3::Hash::from_bytes(assembled.trailer_digest),
616            source_digest,
617            source_pack.len(),
618            source_index.len(),
619        );
620    }
621
622    #[test]
623    fn digest_mismatch_never_produces_an_installable_pack() {
624        let (manifest, mut bodies, _, _) = split_manifest();
625        bodies[1][0] ^= 0xff;
626
627        let error = assemble_provider_pack(&manifest, &bodies).unwrap_err();
628
629        assert!(error.to_string().contains("digest mismatch"));
630    }
631
632    #[test]
633    fn manifest_gaps_and_invalid_or_duplicate_index_entries_fail_closed() {
634        let (mut manifest, bodies, _, _) = split_manifest();
635        manifest.extents[1].output_offset += 1;
636        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
637
638        let (mut manifest, bodies, _, _) = split_manifest();
639        manifest.extents[1].objects[0].output_offset = manifest.extents[0].output_offset;
640        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
641
642        let (mut manifest, bodies, _, _) = split_manifest();
643        manifest.extents[1].objects[0].id = manifest.extents[0].objects[0].id;
644        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
645    }
646
647    #[test]
648    fn positional_spool_accepts_out_of_order_extents_and_is_byte_identical() {
649        let (manifest, bodies, source_pack, source_index) = split_manifest();
650        let root = tempfile::tempdir().unwrap();
651        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
652        assert_eq!(
653            spool.file.metadata().unwrap().len(),
654            source_pack.len() as u64,
655            "the sparse spool must be pre-sized to the exact virtual pack length"
656        );
657        let writer = spool.writer();
658
659        writer.write_extent_chunk(1, 0, &bodies[1]).unwrap();
660        writer.mark_verified(1).unwrap();
661        writer.write_extent_chunk(0, 0, &bodies[0]).unwrap();
662        writer.mark_verified(0).unwrap();
663        drop(writer);
664
665        let completed = spool.finish().unwrap();
666        assert_eq!(fs::read(&completed.pack_path).unwrap(), source_pack);
667        assert_eq!(fs::read(&completed.index_path).unwrap(), source_index);
668    }
669
670    #[test]
671    fn positional_spool_rejects_range_overrun_duplicate_and_incomplete_coverage() {
672        let (manifest, bodies, _, _) = split_manifest();
673        let root = tempfile::tempdir().unwrap();
674        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
675        let spool_dir = spool.dir.clone().unwrap();
676        let writer = spool.writer();
677
678        assert!(
679            writer
680                .write_extent_chunk(0, bodies[0].len() as u64, &[1])
681                .is_err()
682        );
683        writer.write_extent_chunk(0, 0, &bodies[0]).unwrap();
684        writer.mark_verified(0).unwrap();
685        assert!(writer.mark_verified(0).is_err());
686        drop(writer);
687
688        assert!(spool.finish().is_err());
689        assert!(
690            !spool_dir.exists(),
691            "failed finalization must remove partial spool state"
692        );
693    }
694
695    #[test]
696    fn dropping_partial_spool_removes_all_transfer_state() {
697        let (manifest, bodies, _, _) = split_manifest();
698        let root = tempfile::tempdir().unwrap();
699        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
700        let spool_dir = spool.dir.clone().unwrap();
701        let writer = spool.writer();
702        writer
703            .write_extent_chunk(0, 0, &bodies[0][..bodies[0].len() / 2])
704            .unwrap();
705
706        drop(writer);
707        drop(spool);
708
709        assert!(!spool_dir.exists());
710    }
711}