1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum Kind {
14 Blob,
16 ChunkedFile,
18}
19
20#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct InspectObject {
23 pub id: Hash,
25 pub size: u64,
27 pub kind: Kind,
29}
30
31#[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#[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#[must_use]
65pub fn limit_error() -> ServerError {
66 ServerError::invalid_argument("object index limit exceeded")
67}
68
69impl InspectionSet {
70 #[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 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 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 #[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 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 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 #[must_use]
183 pub fn finalize(self) -> Vec<InspectObject> {
184 self.objects.into_values().collect()
185 }
186}
187
188pub(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
289pub(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)] 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)] 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 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)] 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)] 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}