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: 4,
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        objects::fs_atomic::sync_file(&self.file, &self.pack_path)?;
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    /// Validated native pack and index paths, retained by this completed spool.
259    pub fn artifact_paths(&self) -> (&Path, &Path) {
260        (&self.pack_path, &self.index_path)
261    }
262
263    /// Atomically install a fully validated provider pack into the object store.
264    pub fn install_into(&mut self, store: &impl ObjectStore) -> Result<Vec<PackObjectId>> {
265        store
266            .install_pack_streaming(&self.pack_path, &self.index_path)
267            .map_err(ProtocolError::from)
268    }
269}
270
271impl Drop for ProviderPackSpool {
272    fn drop(&mut self) {
273        if let Some(dir) = self.dir.as_ref() {
274            let _ = fs::remove_dir_all(dir);
275        }
276    }
277}
278
279impl Drop for CompletedProviderPack {
280    fn drop(&mut self) {
281        let _ = fs::remove_dir_all(&self.dir);
282    }
283}
284
285fn write_provider_index(path: &Path, manifest: &ProviderPackManifest) -> Result<()> {
286    let mut index = PackIndex::new();
287    for extent in &manifest.extents {
288        for object in &extent.objects {
289            index.add(object.id, object.output_offset);
290        }
291    }
292    index.sort();
293    let mut file = OpenOptions::new().write(true).create_new(true).open(path)?;
294    file.write_all(&index.to_bytes())?;
295    file.flush()?;
296    objects::fs_atomic::sync_file(&file, path)?;
297    Ok(())
298}
299
300fn hash_file_prefix(file: &File, length: u64) -> Result<[u8; 32]> {
301    let mut hasher = blake3::Hasher::new();
302    let mut buffer = [0_u8; 64 * 1024];
303    let mut offset = 0_u64;
304    while offset < length {
305        let remaining = length - offset;
306        let read_len = usize::try_from(remaining.min(buffer.len() as u64)).map_err(|_| {
307            ProtocolError::InvalidState("provider pack hash length exceeds usize".to_string())
308        })?;
309        read_exact_at(file, &mut buffer[..read_len], offset)?;
310        hasher.update(&buffer[..read_len]);
311        offset += read_len as u64;
312    }
313    Ok(*hasher.finalize().as_bytes())
314}
315
316#[cfg(unix)]
317fn write_all_at(file: &File, mut data: &[u8], mut offset: u64) -> Result<()> {
318    use std::os::unix::fs::FileExt;
319
320    while !data.is_empty() {
321        let written = file.write_at(data, offset)?;
322        if written == 0 {
323            return Err(ProtocolError::Io(std::io::Error::new(
324                std::io::ErrorKind::WriteZero,
325                "provider positional spool write returned zero",
326            )));
327        }
328        data = &data[written..];
329        offset = offset.checked_add(written as u64).ok_or_else(|| {
330            ProtocolError::InvalidState("provider positional write offset overflows".to_string())
331        })?;
332    }
333    Ok(())
334}
335
336#[cfg(windows)]
337fn write_all_at(file: &File, mut data: &[u8], mut offset: u64) -> Result<()> {
338    use std::os::windows::fs::FileExt;
339
340    while !data.is_empty() {
341        let written = file.seek_write(data, offset)?;
342        if written == 0 {
343            return Err(ProtocolError::Io(std::io::Error::new(
344                std::io::ErrorKind::WriteZero,
345                "provider positional spool write returned zero",
346            )));
347        }
348        data = &data[written..];
349        offset = offset.checked_add(written as u64).ok_or_else(|| {
350            ProtocolError::InvalidState("provider positional write offset overflows".to_string())
351        })?;
352    }
353    Ok(())
354}
355
356#[cfg(unix)]
357fn read_exact_at(file: &File, mut data: &mut [u8], mut offset: u64) -> Result<()> {
358    use std::os::unix::fs::FileExt;
359
360    while !data.is_empty() {
361        let read = file.read_at(data, offset)?;
362        if read == 0 {
363            return Err(ProtocolError::Io(std::io::Error::new(
364                std::io::ErrorKind::UnexpectedEof,
365                "provider positional spool read ended early",
366            )));
367        }
368        data = &mut data[read..];
369        offset = offset.checked_add(read as u64).ok_or_else(|| {
370            ProtocolError::InvalidState("provider positional read offset overflows".to_string())
371        })?;
372    }
373    Ok(())
374}
375
376#[cfg(windows)]
377fn read_exact_at(file: &File, mut data: &mut [u8], mut offset: u64) -> Result<()> {
378    use std::os::windows::fs::FileExt;
379
380    while !data.is_empty() {
381        let read = file.seek_read(data, offset)?;
382        if read == 0 {
383            return Err(ProtocolError::Io(std::io::Error::new(
384                std::io::ErrorKind::UnexpectedEof,
385                "provider positional spool read ended early",
386            )));
387        }
388        data = &mut data[read..];
389        offset = offset.checked_add(read as u64).ok_or_else(|| {
390            ProtocolError::InvalidState("provider positional read offset overflows".to_string())
391        })?;
392    }
393    Ok(())
394}
395
396/// Assemble and verify one virtual native pack from provider extent bodies.
397///
398/// `extent_bodies` uses the same order as `manifest.extents`. The manifest may
399/// arrive in any order, but its output layout must cover every byte between the
400/// 16-byte header and 32-byte trailer exactly once.
401pub fn assemble_provider_pack(
402    manifest: &ProviderPackManifest,
403    extent_bodies: &[Vec<u8>],
404) -> Result<ProviderPackBundle> {
405    validate_manifest(manifest)?;
406    if extent_bodies.len() != manifest.extents.len() {
407        return Err(ProtocolError::InvalidState(format!(
408            "provider extent count mismatch: expected {}, got {}",
409            manifest.extents.len(),
410            extent_bodies.len()
411        )));
412    }
413
414    let mut order = (0..manifest.extents.len()).collect::<Vec<_>>();
415    order.sort_unstable_by_key(|index| manifest.extents[*index].output_offset);
416    let output_len = usize::try_from(manifest.output_pack_length).map_err(|_| {
417        ProtocolError::InvalidState("provider output pack exceeds this platform".to_string())
418    })?;
419    let mut pack_data = Vec::with_capacity(output_len);
420    pack_data.extend_from_slice(&manifest.header);
421    for index in order {
422        let extent = &manifest.extents[index];
423        let body = &extent_bodies[index];
424        let expected_len = usize::try_from(extent.length).map_err(|_| {
425            ProtocolError::InvalidState("provider extent exceeds this platform".to_string())
426        })?;
427        if body.len() != expected_len {
428            return Err(ProtocolError::InvalidState(format!(
429                "provider extent length mismatch: expected {}, got {}",
430                extent.length,
431                body.len()
432            )));
433        }
434        if blake3::hash(body).as_bytes() != &extent.digest {
435            return Err(ProtocolError::InvalidState(
436                "provider extent digest mismatch".to_string(),
437            ));
438        }
439        pack_data.extend_from_slice(body);
440    }
441
442    let trailer_digest = *blake3::hash(&pack_data).as_bytes();
443    pack_data.extend_from_slice(&trailer_digest);
444    if pack_data.len() != output_len {
445        return Err(ProtocolError::InvalidState(format!(
446            "provider output pack length mismatch: expected {}, got {}",
447            manifest.output_pack_length,
448            pack_data.len()
449        )));
450    }
451    verify_container(&pack_data, PACK_SPEC).map_err(ProtocolError::from)?;
452
453    let mut index = PackIndex::new();
454    for extent in &manifest.extents {
455        for object in &extent.objects {
456            index.add(object.id, object.output_offset);
457        }
458    }
459    index.sort();
460
461    Ok(ProviderPackBundle {
462        pack: NativePackBundle {
463            pack_data,
464            index_data: index.to_bytes(),
465        },
466        trailer_digest,
467    })
468}
469
470fn validate_manifest(manifest: &ProviderPackManifest) -> Result<()> {
471    if manifest.output_pack_length > MAX_RECEIVED_PACK_SIZE
472        || manifest.output_pack_length < (PACK_HEADER_LEN + PACK_TRAILER_LEN) as u64
473    {
474        return Err(ProtocolError::InvalidState(
475            "provider output pack length is invalid".to_string(),
476        ));
477    }
478    if &manifest.header[..4] != PACK_SPEC.magic
479        || u32::from_be_bytes(manifest.header[4..8].try_into().map_err(|_| {
480            ProtocolError::InvalidState("provider pack header is truncated".to_string())
481        })?) != PACK_SPEC.version
482    {
483        return Err(ProtocolError::InvalidState(
484            "provider pack header has invalid magic or version".to_string(),
485        ));
486    }
487    let object_count = u64::from_be_bytes(manifest.header[8..16].try_into().map_err(|_| {
488        ProtocolError::InvalidState("provider pack header is truncated".to_string())
489    })?);
490    let expected_body_end = manifest.output_pack_length - PACK_TRAILER_LEN as u64;
491    let mut order = manifest.extents.iter().collect::<Vec<_>>();
492    order.sort_unstable_by_key(|extent| extent.output_offset);
493    let mut next_offset = PACK_HEADER_LEN as u64;
494    let mut ids = HashSet::new();
495    let mut object_offsets = HashSet::new();
496    let mut actual_object_count = 0_u64;
497
498    for extent in order {
499        if extent.length == 0 || extent.output_offset != next_offset {
500            return Err(ProtocolError::InvalidState(
501                "provider extents do not exactly cover the virtual pack body".to_string(),
502            ));
503        }
504        next_offset = extent
505            .output_offset
506            .checked_add(extent.length)
507            .ok_or_else(|| {
508                ProtocolError::InvalidState("provider extent output range overflows".to_string())
509            })?;
510        if next_offset > expected_body_end {
511            return Err(ProtocolError::InvalidState(
512                "provider extent exceeds the virtual pack body".to_string(),
513            ));
514        }
515        for object in &extent.objects {
516            if object.output_offset < extent.output_offset
517                || object.output_offset >= next_offset
518                || !ids.insert(object.id)
519                || !object_offsets.insert(object.output_offset)
520            {
521                return Err(ProtocolError::InvalidState(
522                    "provider object index is outside its extent or duplicated".to_string(),
523                ));
524            }
525            actual_object_count = actual_object_count.checked_add(1).ok_or_else(|| {
526                ProtocolError::InvalidState("provider object count overflows".to_string())
527            })?;
528        }
529    }
530    if next_offset != expected_body_end || actual_object_count != object_count {
531        return Err(ProtocolError::InvalidState(
532            "provider manifest body or object count is incomplete".to_string(),
533        ));
534    }
535    Ok(())
536}
537
538#[cfg(test)]
539mod tests {
540    use objects::{
541        object::ContentHash,
542        store::{
543            CompressionConfig,
544            pack::{ObjectType, PackBuilder, PackIndex, PackObjectId},
545        },
546    };
547
548    use super::*;
549
550    fn source_pack() -> (Vec<u8>, Vec<u8>, Vec<PackObjectId>) {
551        let ids = vec![
552            PackObjectId::Hash(ContentHash::from_bytes([1; 32])),
553            PackObjectId::Hash(ContentHash::from_bytes([2; 32])),
554        ];
555        let mut builder = PackBuilder::new(CompressionConfig {
556            enabled: false,
557            ..CompressionConfig::default()
558        });
559        builder.add_id(ids[0], ObjectType::Blob, b"provider-one".to_vec());
560        builder.add_id(ids[1], ObjectType::Blob, b"provider-two".to_vec());
561        let (pack, index, _) = builder.build().unwrap();
562        (pack, index, ids)
563    }
564
565    fn split_manifest() -> (ProviderPackManifest, Vec<Vec<u8>>, Vec<u8>, Vec<u8>) {
566        let (pack, index, ids) = source_pack();
567        let parsed_index = PackIndex::from_bytes(&index).unwrap();
568        let first = parsed_index
569            .find(&ids[0])
570            .unwrap()
571            .expect("first fixture object indexed");
572        let second = parsed_index
573            .find(&ids[1])
574            .unwrap()
575            .expect("second fixture object indexed");
576        let (first_id, first_offset, second_id, second_offset) = if first < second {
577            (ids[0], first, ids[1], second)
578        } else {
579            (ids[1], second, ids[0], first)
580        };
581        let body_end = pack.len() - PACK_TRAILER_LEN;
582        let first_body = pack[first_offset as usize..second_offset as usize].to_vec();
583        let second_body = pack[second_offset as usize..body_end].to_vec();
584        let manifest = ProviderPackManifest {
585            header: pack[..PACK_HEADER_LEN].try_into().unwrap(),
586            output_pack_length: pack.len() as u64,
587            extents: vec![
588                ProviderPackExtent {
589                    output_offset: first_offset,
590                    length: first_body.len() as u64,
591                    digest: *blake3::hash(&first_body).as_bytes(),
592                    objects: vec![ProviderPackIndexEntry {
593                        id: first_id,
594                        output_offset: first_offset,
595                    }],
596                },
597                ProviderPackExtent {
598                    output_offset: second_offset,
599                    length: second_body.len() as u64,
600                    digest: *blake3::hash(&second_body).as_bytes(),
601                    objects: vec![ProviderPackIndexEntry {
602                        id: second_id,
603                        output_offset: second_offset,
604                    }],
605                },
606            ],
607        };
608        (manifest, vec![first_body, second_body], pack, index)
609    }
610
611    #[test]
612    fn provider_and_ordinary_pack_results_are_byte_identical() {
613        let (manifest, bodies, source_pack, source_index) = split_manifest();
614
615        let assembled = assemble_provider_pack(&manifest, &bodies).unwrap();
616
617        assert_eq!(assembled.pack.pack_data, source_pack);
618        assert_eq!(assembled.pack.index_data, source_index);
619        let source_digest = blake3::Hash::from_bytes(
620            source_pack[source_pack.len() - PACK_TRAILER_LEN..]
621                .try_into()
622                .unwrap(),
623        );
624        println!(
625            "byte_identical provider_digest={} weft_digest={} pack_bytes={} index_bytes={} identical=true",
626            blake3::Hash::from_bytes(assembled.trailer_digest),
627            source_digest,
628            source_pack.len(),
629            source_index.len(),
630        );
631    }
632
633    #[test]
634    fn digest_mismatch_never_produces_an_installable_pack() {
635        let (manifest, mut bodies, _, _) = split_manifest();
636        bodies[1][0] ^= 0xff;
637
638        let error = assemble_provider_pack(&manifest, &bodies).unwrap_err();
639
640        assert!(error.to_string().contains("digest mismatch"));
641    }
642
643    #[test]
644    fn manifest_gaps_and_invalid_or_duplicate_index_entries_fail_closed() {
645        let (mut manifest, bodies, _, _) = split_manifest();
646        manifest.extents[1].output_offset += 1;
647        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
648
649        let (mut manifest, bodies, _, _) = split_manifest();
650        manifest.extents[1].objects[0].output_offset = manifest.extents[0].output_offset;
651        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
652
653        let (mut manifest, bodies, _, _) = split_manifest();
654        manifest.extents[1].objects[0].id = manifest.extents[0].objects[0].id;
655        assert!(assemble_provider_pack(&manifest, &bodies).is_err());
656    }
657
658    #[test]
659    fn positional_spool_accepts_out_of_order_extents_and_is_byte_identical() {
660        let (manifest, bodies, source_pack, source_index) = split_manifest();
661        let root = tempfile::tempdir().unwrap();
662        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
663        assert_eq!(
664            spool.file.metadata().unwrap().len(),
665            source_pack.len() as u64,
666            "the sparse spool must be pre-sized to the exact virtual pack length"
667        );
668        let writer = spool.writer();
669
670        writer.write_extent_chunk(1, 0, &bodies[1]).unwrap();
671        writer.mark_verified(1).unwrap();
672        writer.write_extent_chunk(0, 0, &bodies[0]).unwrap();
673        writer.mark_verified(0).unwrap();
674        drop(writer);
675
676        let completed = spool.finish().unwrap();
677        assert_eq!(fs::read(&completed.pack_path).unwrap(), source_pack);
678        assert_eq!(fs::read(&completed.index_path).unwrap(), source_index);
679    }
680
681    #[test]
682    fn positional_spool_rejects_range_overrun_duplicate_and_incomplete_coverage() {
683        let (manifest, bodies, _, _) = split_manifest();
684        let root = tempfile::tempdir().unwrap();
685        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
686        let spool_dir = spool.dir.clone().unwrap();
687        let writer = spool.writer();
688
689        assert!(
690            writer
691                .write_extent_chunk(0, bodies[0].len() as u64, &[1])
692                .is_err()
693        );
694        writer.write_extent_chunk(0, 0, &bodies[0]).unwrap();
695        writer.mark_verified(0).unwrap();
696        assert!(writer.mark_verified(0).is_err());
697        drop(writer);
698
699        assert!(spool.finish().is_err());
700        assert!(
701            !spool_dir.exists(),
702            "failed finalization must remove partial spool state"
703        );
704    }
705
706    #[test]
707    fn dropping_partial_spool_removes_all_transfer_state() {
708        let (manifest, bodies, _, _) = split_manifest();
709        let root = tempfile::tempdir().unwrap();
710        let spool = ProviderPackSpool::new_in(root.path(), manifest).unwrap();
711        let spool_dir = spool.dir.clone().unwrap();
712        let writer = spool.writer();
713        writer
714            .write_extent_chunk(0, 0, &bodies[0][..bodies[0].len() / 2])
715            .unwrap();
716
717        drop(writer);
718        drop(spool);
719
720        assert!(!spool_dir.exists());
721    }
722}