Skip to main content

mkit_server/indexed/
inspection.rs

1//! Bounded, metadata-only enumeration of added-pack files for launch inspection.
2use std::collections::{BTreeMap, BTreeSet};
3
4use mkit_core::hash::Hash;
5use mkit_core::object::ObjectType;
6
7use crate::ServerError;
8use crate::store::{BlobBody, BlobKey, BlobStore, ByteRange, codec::TicketV1};
9use futures::StreamExt as _;
10
11/// Launch inspection kinds come directly from verified object types.
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum Kind {
14    /// Any blob, including one used only as a chunk.
15    Blob,
16    /// A `ChunkedBlob` manifest.
17    ChunkedFile,
18}
19
20/// One inspected object's metadata, without its bytes.
21#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct InspectObject {
23    /// Verified content identity.
24    pub id: Hash,
25    /// Canonical decoded object length.
26    pub size: u64,
27    /// Blob or chunked file.
28    pub kind: Kind,
29}
30
31/// The sorted, deduplicated file-typed entries of the advance's added packs.
32#[derive(Debug, Clone, PartialEq, Eq)]
33pub struct InspectionSet {
34    limit: usize,
35    objects: BTreeMap<Hash, InspectObject>,
36    added_entries: u64,
37    raw_packs: BTreeSet<Hash>,
38    pending: Option<Added>,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub(super) struct NativeEntry {
43    pub id: Hash,
44    pub size: u64,
45    pub object_type: u8,
46}
47
48#[derive(Debug, Clone, PartialEq, Eq)]
49enum Added {
50    Native(Vec<NativeEntry>),
51    Scheduled(crate::Partition, Vec<ScheduledPack>),
52}
53
54/// Request-local proof of the exact verified job accepted during ticket checking.
55#[derive(Debug, Clone, PartialEq, Eq)]
56pub(super) struct ScheduledPack {
57    pub pack: Hash,
58    pub job: crate::Value,
59    pub verification: crate::Value,
60    pub decoded_bytes: u64,
61}
62
63/// The established bounded-index refusal; request size is never unavailability.
64#[must_use]
65pub fn limit_error() -> ServerError {
66    ServerError::invalid_argument("object index limit exceeded")
67}
68
69impl InspectionSet {
70    /// Create a whole-advance collector with its configured object limit.
71    #[must_use]
72    pub fn new(limit: usize) -> Self {
73        Self {
74            limit,
75            objects: BTreeMap::new(),
76            added_entries: 0,
77            raw_packs: BTreeSet::new(),
78            pending: None,
79        }
80    }
81
82    /// Refuse the conservative count before decoding or enumeration.
83    ///
84    /// # Errors
85    /// The existing index-limit error if the bound exceeds the launch limit.
86    pub fn preflight(&self, count: u64) -> Result<(), ServerError> {
87        if count > self.limit as u64 {
88            return Err(limit_error());
89        }
90        Ok(())
91    }
92
93    /// Reserve the conservative count of all added-pack entries.
94    ///
95    /// # Errors
96    /// The existing index-limit refusal when the header/job bound is oversized.
97    pub fn reserve_added_count(&mut self, count: u64) -> Result<(), ServerError> {
98        self.preflight(count)?;
99        self.added_entries = count;
100        Ok(())
101    }
102
103    pub(super) fn add_raw_pack(&mut self, id: Hash) {
104        self.raw_packs.insert(id);
105    }
106
107    /// Verified raw additions; excludes packlist nodes without new byte reads.
108    #[cfg(feature = "remote-hooks")]
109    pub(crate) fn raw_packs(&self) -> &BTreeSet<Hash> {
110        &self.raw_packs
111    }
112
113    pub(super) fn defer_native(&mut self, entries: Vec<NativeEntry>) {
114        self.pending = Some(Added::Native(entries));
115    }
116
117    pub(super) fn defer_scheduled(&mut self, packs: Vec<ScheduledPack>, source: crate::Partition) {
118        self.raw_packs.extend(packs.iter().map(|pack| pack.pack));
119        self.pending = Some(Added::Scheduled(source, packs));
120    }
121
122    /// Collect added entries after header/job preflight, without reading bytes.
123    ///
124    /// # Errors
125    /// Existing indexed storage, metadata decoding and bounded-enumeration refusals.
126    pub async fn complete_added<S: crate::NamespaceStore>(
127        &mut self,
128        store: &S,
129        repo: &crate::RepoId,
130    ) -> Result<(), ServerError> {
131        self.preflight(self.added_entries)?;
132        match self.pending.take() {
133            Some(Added::Native(entries)) => {
134                for entry in entries {
135                    self.entry(entry.id, entry.size, entry.object_type)?;
136                }
137            }
138            Some(Added::Scheduled(source, packs)) => {
139                let added =
140                    scheduled_entries(store, repo, &source, &packs, self.limit, self.added_entries)
141                        .await?;
142                self.objects = added.objects;
143            }
144            None => {}
145        }
146        Ok(())
147    }
148
149    /// Insert verified metadata, ignoring non-file object types.
150    ///
151    /// # Errors
152    /// Index-limit refusal or inconsistent immutable metadata.
153    pub fn entry(&mut self, id: Hash, size: u64, object_type: u8) -> Result<(), ServerError> {
154        if !(ObjectType::Blob as u8..=ObjectType::Tag as u8).contains(&object_type) {
155            return Err(ServerError::unavailable(
156                "verified object metadata inconsistency",
157            ));
158        }
159        let kind = if object_type == ObjectType::Blob as u8 {
160            Kind::Blob
161        } else if object_type == ObjectType::ChunkedBlob as u8 {
162            Kind::ChunkedFile
163        } else {
164            return Ok(());
165        };
166        if let Some(old) = self.objects.get(&id) {
167            if old.size != size || old.kind != kind {
168                return Err(ServerError::unavailable(
169                    "verified object metadata inconsistency",
170                ));
171            }
172            return Ok(());
173        }
174        if self.objects.len() >= self.limit {
175            return Err(limit_error());
176        }
177        self.objects.insert(id, InspectObject { id, size, kind });
178        Ok(())
179    }
180
181    /// Final metadata in stable object-id order.
182    #[must_use]
183    pub fn finalize(self) -> Vec<InspectObject> {
184        self.objects.into_values().collect()
185    }
186}
187
188/// Read only fixed-size headers before native whole-pack allocation/decoding.
189pub(super) async fn preflight_native<B: BlobStore>(
190    blobs: &B,
191    tickets: &[TicketV1],
192    limit: usize,
193) -> Result<u64, ServerError> {
194    let mut count = 0_u64;
195    let mut packs = BTreeSet::new();
196    for ticket in tickets {
197        if !packs.insert(ticket.pack_id) {
198            continue;
199        }
200        let body = blobs
201            .get(
202                &BlobKey::pack(ticket.pack_id),
203                Some(ByteRange {
204                    start: 0,
205                    end_inclusive: 11,
206                }),
207            )
208            .await
209            .map_err(|_| ServerError::unavailable("object storage request failed"))?
210            .ok_or_else(|| ServerError::unavailable("object storage request failed"))?;
211        let mut prefix = Vec::with_capacity(12);
212        match body {
213            BlobBody::Bytes(bytes) => {
214                if bytes.len() != 12 {
215                    return Err(ServerError::invalid_argument("object hash mismatch"));
216                }
217                prefix.extend_from_slice(&bytes);
218            }
219            BlobBody::Stream { len, mut stream } => {
220                if len != 12 {
221                    return Err(ServerError::unavailable("object storage request failed"));
222                }
223                while let Some(bytes) = stream.next().await {
224                    let bytes = bytes
225                        .map_err(|_| ServerError::unavailable("object storage request failed"))?;
226                    if prefix.len().saturating_add(bytes.len()) > 12 {
227                        return Err(ServerError::unavailable("object storage request failed"));
228                    }
229                    prefix.extend_from_slice(&bytes);
230                }
231            }
232        }
233        if prefix.len() != 12 {
234            return Err(ServerError::invalid_argument("object hash mismatch"));
235        }
236        if super::classify::classify(&prefix)? == super::classify::UploadType::Pack {
237            let entries = u32::from_le_bytes(prefix[8..12].try_into().map_err(|_| limit_error())?);
238            count = count.saturating_add(u64::from(entries));
239            if count > limit as u64 {
240                return Err(limit_error());
241            }
242        }
243    }
244    Ok(count)
245}
246
247fn metadata_error(error: &crate::StoreError) -> ServerError {
248    if super::budget::is_exhausted(error) {
249        limit_error()
250    } else {
251        ServerError::unavailable("object storage request failed")
252    }
253}
254
255async fn guard_jobs<S: crate::NamespaceStore>(
256    store: &S,
257    repo: &crate::RepoId,
258    source: &crate::Partition,
259    packs: &[ScheduledPack],
260) -> Result<(), ServerError> {
261    if packs.is_empty() {
262        return Ok(());
263    }
264    let keys = packs
265        .iter()
266        .flat_map(|pack| {
267            [
268                crate::store::keys::verify_job(&repo.name, &pack.pack),
269                crate::store::keys::verification(&repo.name, &pack.pack),
270            ]
271        })
272        .collect::<Vec<_>>();
273    let current = store
274        .get_many(source, &keys)
275        .await
276        .map_err(|error| metadata_error(&error))?;
277    if current.len() != keys.len() {
278        return Err(ServerError::unavailable("object storage request failed"));
279    }
280    if packs.iter().enumerate().any(|(i, pack)| {
281        current[2 * i].as_ref() != Some(&pack.job)
282            || current[2 * i + 1].as_ref() != Some(&pack.verification)
283    }) {
284        return Err(super::pending(1000));
285    }
286    Ok(())
287}
288
289/// Page verified first-occurrence frame rows, without reference scans or R2 reads.
290pub(super) async fn scheduled_entries<S: crate::NamespaceStore>(
291    store: &S,
292    repo: &crate::RepoId,
293    source: &crate::Partition,
294    additions: &[ScheduledPack],
295    limit: usize,
296    added_count: u64,
297) -> Result<InspectionSet, ServerError> {
298    use crate::store::keys;
299    let failed = || ServerError::unavailable("object storage request failed");
300    let mut set = InspectionSet::new(limit);
301    set.reserve_added_count(added_count)?;
302    guard_jobs(store, repo, source, additions).await?;
303    let mut entries = 0_u64;
304    for pack in additions {
305        let mut decoded_bytes = 0_u64;
306        let (start, end) = keys::verify_range(&repo.name, &pack.pack, Some(keys::VC_FRAME));
307        let mut after = None;
308        loop {
309            let page = store
310                .scan(source, &start, &end, after.as_ref(), 1000)
311                .await
312                .map_err(|error| metadata_error(&error))?;
313            if page.entries.len() > 1000 {
314                return Err(failed());
315            }
316            for (key, raw) in page.entries {
317                let Some(keys::ParsedKey::VerifyCursor {
318                    repo: found,
319                    pack_id,
320                    sub,
321                    id: Some(id),
322                }) = keys::parse(&key)
323                else {
324                    return Err(failed());
325                };
326                if found != repo.name || pack_id != pack.pack || sub != keys::VC_FRAME {
327                    return Err(failed());
328                }
329                entries = entries.saturating_add(1);
330                if entries > added_count {
331                    return Err(failed());
332                }
333                let frame = super::checkpoint::decode_frame(&id, &raw).map_err(|_| failed())?;
334                decoded_bytes = decoded_bytes
335                    .checked_add(frame.value.decoded_size)
336                    .ok_or_else(failed)?;
337                set.entry(id, frame.value.decoded_size, frame.object_type)?;
338            }
339            match page.next {
340                Some(next) if after.as_ref() != Some(&next) => after = Some(next),
341                Some(_) => return Err(failed()),
342                None => break,
343            }
344        }
345        if decoded_bytes != pack.decoded_bytes {
346            return Err(super::pending(1000));
347        }
348    }
349    guard_jobs(store, repo, source, additions).await?;
350    Ok(set)
351}
352
353#[cfg(all(test, feature = "memory"))]
354mod tests {
355    use super::*;
356    use crate::indexed::checkpoint::{FrameRow, encode_frame};
357    use crate::indexed::tests::{NOW, repo, source, ticket, upload};
358    use crate::memory::{MemoryBlobStore, MemoryKv};
359    use crate::pipeline::SinglePartition;
360    use crate::rt::ManualClock;
361    use crate::store::{Batch, NamespaceStore, index::IndexValue, keys};
362    use crate::telemetry::NoopMetrics;
363    use futures_executor::block_on;
364    use mkit_core::hash::hash;
365    use mkit_core::object::{
366        Blob, ChunkedBlob, Commit, EntryMode, Identity, Object, Tree, TreeEntry,
367    };
368    use mkit_core::pack::{DecodeLimits, NoExternalBases, PackWriter, decode_entries_with};
369    use mkit_core::serialize::serialize;
370    use mkit_core::sign::{KeyPair, sign_commit};
371    use std::sync::Arc;
372
373    fn verified_pack(
374        store: &MemoryKv,
375        repo: &crate::RepoId,
376        pack: Hash,
377        entries: u64,
378        decoded_bytes: u64,
379    ) -> ScheduledPack {
380        let mut job = super::super::checkpoint::VerifyJobV1::new(hash(&pack), NOW as u64, 1, 4096);
381        job.kind = super::super::checkpoint::Kind::Pack;
382        job.phase = super::super::checkpoint::Phase::Watch;
383        job.entries = entries;
384        job.in_pack_bytes = decoded_bytes;
385        let snapshot = ScheduledPack {
386            pack,
387            job: super::super::checkpoint::encode_job(&job),
388            verification: super::super::state::encode(
389                &super::super::state::VerificationV1::Verified {
390                    pack_len: 1,
391                    verified_at_ms: NOW as u64,
392                    publication: None,
393                },
394            ),
395            decoded_bytes,
396        };
397        block_on(
398            store.apply(
399                &source(repo),
400                Batch::new()
401                    .put(keys::verify_job(&repo.name, &pack), snapshot.job.clone())
402                    .put(
403                        keys::verification(&repo.name, &pack),
404                        snapshot.verification.clone(),
405                    ),
406            ),
407        )
408        .expect("seed deterministic verified-pack metadata");
409        snapshot
410    }
411
412    #[allow(clippy::unwrap_used)] // Deterministic valid objects used only by these tests.
413    fn fixture() -> (Vec<u8>, Hash, Vec<(Hash, Kind)>) {
414        let chunk = Object::Blob(Blob {
415            data: b"one".to_vec(),
416        });
417        let dual = Object::Blob(Blob {
418            data: b"two".to_vec(),
419        });
420        let extra = Object::Blob(Blob {
421            data: b"surplus".to_vec(),
422        });
423        let manifest = Object::ChunkedBlob(ChunkedBlob {
424            total_size: 6,
425            chunk_size: 0,
426            chunks: vec![chunk.id().unwrap(), dual.id().unwrap()],
427        });
428        let tree = Object::Tree(Tree {
429            entries: vec![
430                TreeEntry {
431                    name: b"direct".to_vec(),
432                    mode: EntryMode::Blob,
433                    object_hash: dual.id().unwrap(),
434                },
435                TreeEntry {
436                    name: b"manifest".to_vec(),
437                    mode: EntryMode::Blob,
438                    object_hash: manifest.id().unwrap(),
439                },
440            ],
441        });
442        let key = KeyPair::from_seed([7; 32]);
443        let mut commit = Commit::new_unannotated(
444            tree.id().unwrap(),
445            vec![],
446            Identity::ed25519(key.public.0),
447            key.public.0,
448            b"inspection".to_vec(),
449            42,
450            [0; 64],
451        );
452        commit.signature = sign_commit(&commit, &key).unwrap().0;
453        let commit = Object::Commit(commit);
454        let mut writer = PackWriter::new_raw_only();
455        for object in [&chunk, &dual, &extra, &manifest, &tree, &commit, &chunk] {
456            writer
457                .push_raw(object.id().unwrap(), &serialize(object).unwrap())
458                .unwrap();
459        }
460        let mut expected = vec![
461            (chunk.id().unwrap(), Kind::Blob),
462            (dual.id().unwrap(), Kind::Blob),
463            (extra.id().unwrap(), Kind::Blob),
464            (manifest.id().unwrap(), Kind::ChunkedFile),
465        ];
466        expected.sort_by_key(|e| e.0);
467        (writer.finish().unwrap(), commit.id().unwrap(), expected)
468    }
469
470    #[test]
471    #[allow(clippy::too_many_lines)] // One native/scheduled parity scenario and its preflight failure.
472    fn native_and_worker_enumerate_surplus_manifests_chunks_dual_uses_and_duplicates() {
473        let (pack, head, expected) = fixture();
474        let repo = repo("inspection-union");
475        let blobs = MemoryBlobStore::default();
476        upload(&blobs, &pack);
477        let clock = Arc::new(ManualClock::new(NOW));
478        let store = MemoryKv::with_clock(clock.clone());
479        let ticket = ticket(&repo, &pack, NOW as u64);
480        let mut native = block_on(super::super::verify::verify_ticketed_inspected(
481            &blobs,
482            &store,
483            &SinglePartition,
484            &repo,
485            &source(&repo),
486            std::slice::from_ref(&ticket),
487            &[hash(&ticket.pack_id)],
488            head,
489            super::super::IndexedConfig::default(),
490            clock.as_ref(),
491            &NoopMetrics,
492            10,
493        ))
494        .unwrap()
495        .inspection
496        .unwrap();
497        assert!(native.objects.is_empty());
498        block_on(native.complete_added(&store, &repo)).unwrap();
499        let native = native.finalize();
500        assert_eq!(
501            native.iter().map(|e| (e.id, e.kind)).collect::<Vec<_>>(),
502            expected
503        );
504        let mut frames = BTreeMap::new();
505        decode_entries_with(
506            &pack,
507            &mut NoExternalBases,
508            DecodeLimits::default(),
509            |entry| {
510                let row = FrameRow {
511                    object_type: entry.object.object_type() as u8,
512                    external: None,
513                    value: IndexValue {
514                        frame_offset: entry.frame_offset,
515                        frame_length: entry.frame_length,
516                        wire_type: entry.wire_type,
517                        decoded_size: entry.bytes.len() as u64,
518                        chain_depth: 0,
519                        delta_base: None,
520                    },
521                };
522                frames
523                    .entry(entry.id)
524                    .or_insert_with(|| encode_frame(&entry.id, &row).unwrap());
525                Ok(())
526            },
527        )
528        .unwrap();
529        let in_pack_bytes = frames
530            .iter()
531            .map(|(id, raw)| {
532                super::super::checkpoint::decode_frame(id, raw)
533                    .unwrap()
534                    .value
535                    .decoded_size
536            })
537            .sum();
538        let mut batch = Batch::new();
539        for (id, raw) in frames {
540            batch = batch.put(
541                keys::verify_row(&repo.name, &ticket.pack_id, keys::VC_FRAME, Some(&id)),
542                raw,
543            );
544        }
545        block_on(store.apply(&source(&repo), batch)).unwrap();
546        let mut job = super::super::checkpoint::VerifyJobV1::new(
547            hash(&ticket.pack_id),
548            NOW as u64,
549            pack.len() as u64,
550            4096,
551        );
552        job.kind = super::super::checkpoint::Kind::Pack;
553        job.phase = super::super::checkpoint::Phase::Watch;
554        job.entries = 7;
555        job.in_pack_bytes = in_pack_bytes;
556        let ready = Batch::new()
557            .put(
558                keys::verify_job(&repo.name, &ticket.pack_id),
559                super::super::checkpoint::encode_job(&job),
560            )
561            .put(
562                keys::verification(&repo.name, &ticket.pack_id),
563                super::super::state::encode(&super::super::state::VerificationV1::Verified {
564                    pack_len: pack.len() as u64,
565                    verified_at_ms: NOW as u64,
566                    publication: None,
567                }),
568            );
569        block_on(store.apply(&source(&repo), ready)).unwrap();
570        let accepted = ScheduledPack {
571            pack: ticket.pack_id,
572            job: super::super::checkpoint::encode_job(&job),
573            verification: super::super::state::encode(
574                &super::super::state::VerificationV1::Verified {
575                    pack_len: pack.len() as u64,
576                    verified_at_ms: NOW as u64,
577                    publication: None,
578                },
579            ),
580            decoded_bytes: job.in_pack_bytes,
581        };
582        let worker = block_on(scheduled_entries(
583            &store,
584            &repo,
585            &source(&repo),
586            &[accepted],
587            10,
588            7,
589        ))
590        .unwrap()
591        .finalize();
592        assert_eq!(native, worker);
593        let mut checked = block_on(super::super::scheduled::check_inspected(
594            &blobs,
595            &store,
596            &SinglePartition,
597            &repo,
598            &source(&repo),
599            std::slice::from_ref(&ticket),
600            &[job.ticket_id],
601            head,
602            super::super::IndexedConfig::default(),
603            clock.as_ref(),
604            &NoopMetrics,
605            10,
606        ))
607        .unwrap()
608        .inspection
609        .unwrap();
610        assert!(checked.objects.is_empty());
611        block_on(checked.complete_added(&store, &repo)).unwrap();
612        assert_eq!(checked.finalize(), native);
613        let error = block_on(super::super::scheduled::check_inspected(
614            &blobs,
615            &store,
616            &SinglePartition,
617            &repo,
618            &source(&repo),
619            &[ticket],
620            &[job.ticket_id],
621            head,
622            super::super::IndexedConfig::default(),
623            clock.as_ref(),
624            &NoopMetrics,
625            6,
626        ))
627        .unwrap_err();
628        assert_eq!(error.public_message(), "object index limit exceeded");
629    }
630
631    #[test]
632    fn native_preflight_counts_repeated_pack_once_at_the_entry_limit() {
633        let (pack, _, _) = fixture();
634        let repo = repo("inspection-repeated-pack");
635        let blobs = MemoryBlobStore::default();
636        upload(&blobs, &pack);
637        let mut first = ticket(&repo, &pack, NOW as u64);
638        first.reservation_id = "s:duplicate-first".into();
639        let mut second = first.clone();
640        second.reservation_id = "s:duplicate-second".into();
641        assert_ne!(
642            crate::store::tickets::ticket_id(&first.reservation_id),
643            crate::store::tickets::ticket_id(&second.reservation_id),
644        );
645        assert_eq!(first.pack_id, second.pack_id);
646        let calls = super::super::budget::SliceBudget::new(4);
647        let counted = super::super::budget::Budgeted::new(&blobs, &calls);
648        assert_eq!(
649            block_on(preflight_native(&counted, &[first, second], 7))
650                .map_err(|error| (error.code(), error.public_message().to_owned())),
651            Ok(7),
652        );
653        // One header range reserves both metadata and byte requests.
654        assert_eq!(calls.used(), 2);
655    }
656
657    #[test]
658    fn oversize_native_header_refuses_before_verification_writes() {
659        let (pack, head, _) = fixture();
660        let repo = repo("inspection-size");
661        let blobs = MemoryBlobStore::default();
662        upload(&blobs, &pack);
663        let clock = Arc::new(ManualClock::new(NOW));
664        let store = MemoryKv::with_clock(clock.clone());
665        let ticket = ticket(&repo, &pack, NOW as u64);
666        let error = block_on(super::super::verify::verify_ticketed_inspected(
667            &blobs,
668            &store,
669            &SinglePartition,
670            &repo,
671            &source(&repo),
672            &[ticket],
673            &[[1; 32]],
674            head,
675            super::super::IndexedConfig::default(),
676            clock.as_ref(),
677            &NoopMetrics,
678            6,
679        ))
680        .unwrap_err();
681        assert_eq!(error.code(), crate::Code::InvalidArgument);
682        assert_eq!(error.public_message(), "object index limit exceeded");
683        assert_eq!(block_on(store.stats(&source(&repo))).unwrap().keys, Some(0));
684    }
685
686    #[test]
687    fn frame_scan_budget_exhaustion_is_an_index_limit_refusal() {
688        let repo = repo("inspection-scan-budget");
689        let store = MemoryKv::default();
690        let accepted = verified_pack(&store, &repo, [1; 32], 1, 11);
691        let budget = super::super::budget::SliceBudget::new(1);
692        let bounded = super::super::budget::Budgeted::new(&store, &budget);
693        let error = block_on(scheduled_entries(
694            &bounded,
695            &repo,
696            &source(&repo),
697            &[accepted],
698            10_000,
699            1,
700        ))
701        .unwrap_err();
702        assert_eq!(error.code(), crate::Code::InvalidArgument);
703        assert_eq!(error.public_message(), "object index limit exceeded");
704        assert_eq!(budget.used(), 1);
705    }
706
707    #[test]
708    fn conservative_added_count_is_independent_of_deduplication() {
709        let mut set = InspectionSet::new(4);
710        set.reserve_added_count(4).unwrap();
711        set.entry([1; 32], 11, ObjectType::Blob as u8).unwrap();
712        set.entry([1; 32], 11, ObjectType::Blob as u8).unwrap();
713        assert!(set.reserve_added_count(5).is_err());
714        assert_eq!(set.finalize().len(), 1);
715    }
716
717    #[test]
718    fn unknown_checkpoint_object_types_fail_closed() {
719        let mut set = InspectionSet::new(4);
720        for tag in [0, 8, 255] {
721            assert_eq!(
722                set.entry([tag; 32], 11, tag).unwrap_err().code(),
723                crate::Code::Unavailable
724            );
725        }
726        assert!(set.finalize().is_empty());
727    }
728
729    #[test]
730    #[allow(clippy::too_many_lines)] // One real 400-entry pack compares native and Worker enumeration.
731    fn small_valid_manifest_pack_accepts_with_one_scan_two_guards_and_native_parity() {
732        let repo = repo("manifest-budget-reproduction");
733        let mut writer = PackWriter::new_raw_only();
734        for n in 0_u8..200 {
735            let blob = Object::Blob(Blob { data: vec![n] });
736            let manifest = Object::ChunkedBlob(ChunkedBlob {
737                total_size: 1,
738                chunk_size: 0,
739                chunks: vec![blob.id().unwrap()],
740            });
741            for object in [&blob, &manifest] {
742                writer
743                    .push_raw(object.id().unwrap(), &serialize(object).unwrap())
744                    .unwrap();
745            }
746        }
747        let pack = writer.finish().unwrap();
748        let pack_id = hash(&pack);
749        let blobs = MemoryBlobStore::default();
750        upload(&blobs, &pack);
751        let store = MemoryKv::default();
752        let mut frames = Vec::new();
753        let mut entries = Vec::new();
754        decode_entries_with(
755            &pack,
756            &mut NoExternalBases,
757            DecodeLimits::default(),
758            |entry| {
759                let object_type = entry.object.object_type() as u8;
760                let row = FrameRow {
761                    object_type,
762                    external: None,
763                    value: IndexValue {
764                        frame_offset: entry.frame_offset,
765                        frame_length: entry.frame_length,
766                        wire_type: entry.wire_type,
767                        decoded_size: entry.bytes.len() as u64,
768                        chain_depth: 0,
769                        delta_base: None,
770                    },
771                };
772                frames.push((
773                    keys::verify_row(&repo.name, &pack_id, keys::VC_FRAME, Some(&entry.id)),
774                    encode_frame(&entry.id, &row).unwrap(),
775                ));
776                entries.push(NativeEntry {
777                    id: entry.id,
778                    size: entry.bytes.len() as u64,
779                    object_type,
780                });
781                Ok(())
782            },
783        )
784        .unwrap();
785        assert_eq!(entries.len(), 400);
786        for rows in frames.chunks(90) {
787            let batch = rows.iter().fold(Batch::new(), |batch, (key, raw)| {
788                batch.put(key.clone(), raw.clone())
789            });
790            block_on(store.apply(&source(&repo), batch)).unwrap();
791        }
792        let mut native = InspectionSet::new(10_000);
793        native.reserve_added_count(400).unwrap();
794        native.defer_native(entries);
795        block_on(native.complete_added(&store, &repo)).unwrap();
796        let native = native.finalize();
797        assert_eq!(native.len(), 400);
798        let metadata = super::super::budget::SliceBudget::new(256);
799        let mut worker = InspectionSet::new(10_000);
800        worker.reserve_added_count(400).unwrap();
801        let accepted = verified_pack(
802            &store,
803            &repo,
804            pack_id,
805            400,
806            native.iter().map(|e| e.size).sum(),
807        );
808        worker.defer_scheduled(vec![accepted], source(&repo));
809        block_on(worker.complete_added(
810            &super::super::budget::Budgeted::new(&store, &metadata),
811            &repo,
812        ))
813        .unwrap();
814        assert_eq!(worker.finalize(), native);
815        assert_eq!(metadata.used(), 3);
816    }
817
818    #[test]
819    fn worker_enumeration_resumes_frame_pages_within_the_shared_call_budget() {
820        let repo = repo("inspection-pages");
821        let store = MemoryKv::default();
822        let pack = [9; 32];
823        let mut decoded_bytes = 0_u64;
824        for first in (0_u64..1500).step_by(90) {
825            let mut batch = Batch::new();
826            for n in first..(first + 90).min(1500) {
827                let bytes = serialize(&Object::Blob(Blob {
828                    data: n.to_le_bytes().to_vec(),
829                }))
830                .unwrap();
831                decoded_bytes += bytes.len() as u64;
832                let id = hash(&bytes);
833                let frame = FrameRow {
834                    object_type: ObjectType::Blob as u8,
835                    external: None,
836                    value: IndexValue {
837                        frame_offset: 12 + n * 23,
838                        frame_length: 23,
839                        wire_type: 0,
840                        decoded_size: bytes.len() as u64,
841                        chain_depth: 0,
842                        delta_base: None,
843                    },
844                };
845                batch = batch.put(
846                    keys::verify_row(&repo.name, &pack, keys::VC_FRAME, Some(&id)),
847                    encode_frame(&id, &frame).unwrap(),
848                );
849            }
850            block_on(store.apply(&source(&repo), batch)).unwrap();
851        }
852        let accepted = verified_pack(&store, &repo, pack, 1500, decoded_bytes);
853        let budget = super::super::budget::SliceBudget::new(256);
854        let bounded = super::super::budget::Budgeted::new(&store, &budget);
855        let entries = block_on(scheduled_entries(
856            &bounded,
857            &repo,
858            &source(&repo),
859            &[accepted],
860            1500,
861            1500,
862        ))
863        .unwrap()
864        .finalize();
865        assert_eq!(entries.len(), 1500);
866        assert_eq!(budget.used(), 4);
867        assert!(entries.windows(2).all(|w| w[0].id < w[1].id));
868    }
869    #[test]
870    #[allow(clippy::too_many_lines)] // Worst-case page rounding across seven consumed packs.
871    fn cap_across_seven_packs_uses_sixteen_scans_two_guards_and_matches_native() {
872        let repo = repo("inspection-seven-pack-cap");
873        let store = MemoryKv::default();
874        let counts = [1001_u64, 1001, 1001, 1001, 1001, 1001, 3994];
875        let mut packs = Vec::new();
876        let mut entries = Vec::new();
877        let mut sequence = 0_u64;
878        for (pack_number, count) in counts.into_iter().enumerate() {
879            let pack = [u8::try_from(pack_number + 1).unwrap(); 32];
880            let mut decoded_bytes = 0_u64;
881            for first in (0..count).step_by(90) {
882                let mut batch = Batch::new();
883                for n in first..(first + 90).min(count) {
884                    let bytes = serialize(&Object::Blob(Blob {
885                        data: sequence.to_le_bytes().to_vec(),
886                    }))
887                    .unwrap();
888                    sequence += 1;
889                    decoded_bytes += bytes.len() as u64;
890                    let id = hash(&bytes);
891                    let row = FrameRow {
892                        object_type: ObjectType::Blob as u8,
893                        external: None,
894                        value: IndexValue {
895                            frame_offset: 12 + n * 23,
896                            frame_length: 23,
897                            wire_type: 0,
898                            decoded_size: bytes.len() as u64,
899                            chain_depth: 0,
900                            delta_base: None,
901                        },
902                    };
903                    entries.push(NativeEntry {
904                        id,
905                        size: bytes.len() as u64,
906                        object_type: row.object_type,
907                    });
908                    batch = batch.put(
909                        keys::verify_row(&repo.name, &pack, keys::VC_FRAME, Some(&id)),
910                        encode_frame(&id, &row).unwrap(),
911                    );
912                }
913                block_on(store.apply(&source(&repo), batch)).unwrap();
914            }
915            packs.push(verified_pack(&store, &repo, pack, count, decoded_bytes));
916        }
917        assert_eq!(sequence, 10_000);
918        let metadata = super::super::budget::SliceBudget::new(18);
919        let bounded_store = super::super::budget::Budgeted::new(&store, &metadata);
920        let mut native = InspectionSet::new(10_000);
921        native.reserve_added_count(10_000).unwrap();
922        native.defer_native(entries);
923        block_on(native.complete_added(&bounded_store, &repo)).unwrap();
924        assert_eq!(metadata.used(), 0);
925        let mut worker = InspectionSet::new(10_000);
926        worker.reserve_added_count(10_000).unwrap();
927        worker.defer_scheduled(packs, source(&repo));
928        block_on(worker.complete_added(&bounded_store, &repo)).unwrap();
929        let native = native.finalize();
930        assert_eq!(native.len(), 10_000);
931        assert_eq!(worker.finalize(), native);
932        assert_eq!(metadata.used(), 18);
933    }
934    #[test]
935    fn missing_verified_frames_are_pending_before_inspection() {
936        let repo = repo("inspection-missing-frames");
937        let store = MemoryKv::default();
938        let pack = [7; 32];
939        let accepted = verified_pack(&store, &repo, pack, 1, 11);
940        let mut worker = InspectionSet::new(10_000);
941        worker.reserve_added_count(1).unwrap();
942        worker.defer_scheduled(vec![accepted], source(&repo));
943        let error = block_on(worker.complete_added(&store, &repo)).unwrap_err();
944        assert_eq!(error.public_message(), "pack verification pending");
945    }
946    struct ReplaceAfterScan<'a> {
947        inner: &'a MemoryKv,
948        job_key: crate::Key,
949        replacement: crate::Value,
950    }
951    impl NamespaceStore for ReplaceAfterScan<'_> {
952        fn capabilities(&self) -> crate::StoreCapabilities {
953            self.inner.capabilities()
954        }
955        async fn get(
956            &self,
957            p: &crate::Partition,
958            key: &crate::Key,
959        ) -> Result<Option<crate::Value>, crate::StoreError> {
960            self.inner.get(p, key).await
961        }
962        async fn get_many(
963            &self,
964            p: &crate::Partition,
965            keys: &[crate::Key],
966        ) -> Result<Vec<Option<crate::Value>>, crate::StoreError> {
967            self.inner.get_many(p, keys).await
968        }
969        async fn scan(
970            &self,
971            p: &crate::Partition,
972            start: &crate::Key,
973            end: &crate::Key,
974            after: Option<&crate::store::Cursor>,
975            limit: u32,
976        ) -> Result<crate::store::ScanPage, crate::StoreError> {
977            let page = self.inner.scan(p, start, end, after, limit).await?;
978            self.inner
979                .apply(
980                    p,
981                    Batch::new().put(self.job_key.clone(), self.replacement.clone()),
982                )
983                .await?;
984            Ok(page)
985        }
986        async fn apply(
987            &self,
988            p: &crate::Partition,
989            batch: Batch,
990        ) -> Result<crate::BatchOutcome, crate::StoreError> {
991            self.inner.apply(p, batch).await
992        }
993        async fn stats(
994            &self,
995            p: &crate::Partition,
996        ) -> Result<crate::PartitionStats, crate::StoreError> {
997            self.inner.stats(p).await
998        }
999        async fn probe(&self) -> Result<(), crate::StoreError> {
1000            self.inner.probe().await
1001        }
1002    }
1003
1004    #[test]
1005    fn replacement_verification_job_is_pending_before_and_during_enumeration() {
1006        for during_scan in [false, true] {
1007            let repo = repo("inspection-replaced-job");
1008            let store = MemoryKv::default();
1009            let pack = [7; 32];
1010            let accepted = verified_pack(&store, &repo, pack, 1, 11);
1011            let frame = FrameRow {
1012                object_type: ObjectType::Blob as u8,
1013                external: None,
1014                value: IndexValue {
1015                    frame_offset: 12,
1016                    frame_length: 23,
1017                    wire_type: 0,
1018                    decoded_size: 11,
1019                    chain_depth: 0,
1020                    delta_base: None,
1021                },
1022            };
1023            block_on(store.apply(
1024                &source(&repo),
1025                Batch::new().put(
1026                    keys::verify_row(&repo.name, &pack, keys::VC_FRAME, Some(&[9; 32])),
1027                    encode_frame(&[9; 32], &frame).unwrap(),
1028                ),
1029            ))
1030            .unwrap();
1031            let mut replacement = super::super::checkpoint::decode_job(&accepted.job).unwrap();
1032            replacement.ticket_id = [99; 32];
1033            replacement.phase = super::super::checkpoint::Phase::Decode;
1034            let replacement = super::super::checkpoint::encode_job(&replacement);
1035            let job_key = keys::verify_job(&repo.name, &pack);
1036            let budget = super::super::budget::SliceBudget::new(3);
1037            let mut worker = InspectionSet::new(10_000);
1038            worker.reserve_added_count(1).unwrap();
1039            worker.defer_scheduled(vec![accepted], source(&repo));
1040            let error =
1041                if during_scan {
1042                    let racing = ReplaceAfterScan {
1043                        inner: &store,
1044                        job_key,
1045                        replacement,
1046                    };
1047                    block_on(worker.complete_added(
1048                        &super::super::budget::Budgeted::new(&racing, &budget),
1049                        &repo,
1050                    ))
1051                    .unwrap_err()
1052                } else {
1053                    block_on(store.apply(&source(&repo), Batch::new().put(job_key, replacement)))
1054                        .unwrap();
1055                    block_on(worker.complete_added(
1056                        &super::super::budget::Budgeted::new(&store, &budget),
1057                        &repo,
1058                    ))
1059                    .unwrap_err()
1060                };
1061            assert_eq!(error.public_message(), "pack verification pending");
1062            assert_eq!(budget.used(), if during_scan { 3 } else { 1 });
1063        }
1064    }
1065}