1use super::{
20 IndexedConfig,
21 budget::{Budgeted, PackWindows, SliceBudget, Window, WindowError, is_exhausted},
22 checkpoint::{
23 self, BaseRow, FrameRow, Kind, Outcome, Phase, VerifyJobV1, decode_base, decode_frame,
24 encode_base, encode_frame, encode_job, parse_reference,
25 },
26 classify::{self, UploadType},
27 resolve::{self, MemberCache, ResolveFailure},
28 state::{self, VerificationV1},
29};
30use crate::pipeline::{LeaseParams, ShardMap, renew_for_relay};
31use crate::relay::{commit_relay_rows, relay_delivered_through};
32use crate::repo::RepoId;
33use crate::rt::{BoxFuture, Clock};
34use crate::store::{
35 Batch, BatchOutcome, BlobStore, Cursor, Key, NamespaceStore, Partition, Precondition,
36 StoreError, Value, Write,
37 codec::TicketV1,
38 codec::{self, decode_ticket},
39 index::{self, IndexEntry, IndexValue, LocatedObject},
40 keys,
41};
42use crate::telemetry::Metrics;
43use crate::timers::{
44 DueTimer, Fired, TimerCtx,
45 registry::{TimerHandler, TimerKind, kinds},
46};
47use mkit_core::hash::{Hash, hash};
48use mkit_core::object::Object;
49use mkit_core::ops::graph::{ClosureMode, children};
50use mkit_core::pack::window::{Step, WindowCursor, WindowReader};
51use mkit_core::pack::{
52 DecodeLimits, DeltaBaseSource, PackEntry, PackError, decode_entry_with, decode_frame_with,
53};
54use mkit_core::sign::verify_object_signature;
55use mkit_core::transfer::decode_packlist;
56use std::collections::{BTreeMap, BTreeSet, VecDeque};
57use std::sync::Arc;
58
59mod extraction;
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub struct SliceLimits {
64 pub window_bytes: u64,
66 pub resident_bytes: u64,
71 pub max_subrequests: u32,
73 pub max_entries: u32,
75}
76
77impl Default for SliceLimits {
78 fn default() -> Self {
79 Self {
80 window_bytes: super::geometry::FRAME_PAYLOAD_BYTES,
81 resident_bytes: super::geometry::RESIDENT_BYTES,
82 max_subrequests: 256,
83 max_entries: checkpoint::DEFAULT_ENTRY_CAP,
84 }
85 }
86}
87
88#[cfg(feature = "pack-ruzstd")]
92const DECODER_SCRATCH_BYTES: u64 = 3 * (8 << 20) + (4 << 20);
93#[cfg(feature = "pack-ruzstd")]
96const CACHE_BYTES: u64 = super::geometry::RESIDENT_BYTES
97 - super::geometry::FRAME_PAYLOAD_BYTES
98 - DECODER_SCRATCH_BYTES
99 - 2 * super::geometry::DELTA_STREAM_BYTES
100 - (1 << 20);
101#[cfg(not(feature = "pack-ruzstd"))]
102const CACHE_BYTES: u64 = super::geometry::ENTRY_CACHE_BYTES;
103const ATTEMPTS_PER_CAP: u32 = 3;
105const ENTRY_RESERVE: u32 = 64;
107const CLOSURE_CHUNK: u32 = 1;
109const EMIT_PAGE: u32 = 64;
111const CLEANUP_PAGE: u32 = 90;
113const CLEANUP_ROUNDS: usize = 16;
115const WRITE_BATCH: usize = 90;
117const LAG_BACKOFF_MS: u64 = 15_000;
119const WATCH_POLL_MS: u64 = 3_600_000;
122const MAX_SATISFYING: usize = index::MAX_LOOKUP_IDS;
124
125pub trait SliceExtension: crate::MaybeSend + crate::MaybeSync {
128 fn needs_extraction(&self, object: &Object, cfg: &IndexedConfig) -> bool;
131
132 fn extraction_enabled(&self) -> bool {
134 false
135 }
136
137 fn begin_object<'a>(
139 &'a self,
140 _key: crate::BlobKey,
141 _plan: &'a mkit_core::upload_parts::PartPlan,
142 _root: Hash,
143 _cvs: &'a [Hash],
144 _operation: Hash,
145 _budget: &'a SliceBudget,
146 ) -> BoxFuture<'a, Result<Option<Vec<u8>>, StoreError>> {
147 Box::pin(async { Ok(None) })
148 }
149
150 fn put_object_part<'a>(
152 &'a self,
153 _key: crate::BlobKey,
154 _session: &'a [u8],
155 _plan: &'a mkit_core::upload_parts::PartPlan,
156 _index: u32,
157 _cv: Hash,
158 _bytes: Vec<u8>,
159 _budget: &'a SliceBudget,
160 ) -> BoxFuture<'a, Result<Option<Vec<u8>>, StoreError>> {
161 Box::pin(async { Ok(None) })
162 }
163
164 fn complete_object<'a>(
166 &'a self,
167 _key: crate::BlobKey,
168 _session: &'a [u8],
169 _plan: &'a mkit_core::upload_parts::PartPlan,
170 _parts: Vec<crate::PartRef>,
171 _root: Hash,
172 _budget: &'a SliceBudget,
173 ) -> BoxFuture<'a, Result<Option<crate::CommitOutcome>, StoreError>> {
174 Box::pin(async { Ok(None) })
175 }
176
177 fn abort_object<'a>(
179 &'a self,
180 _key: crate::BlobKey,
181 _session: &'a [u8],
182 _plan: &'a mkit_core::upload_parts::PartPlan,
183 _budget: &'a SliceBudget,
184 ) -> BoxFuture<'a, Result<(), StoreError>> {
185 Box::pin(async { Ok(()) })
186 }
187}
188
189#[derive(Debug, Clone, Copy, Default)]
192pub struct FailClosedExtraction;
193
194impl SliceExtension for FailClosedExtraction {
195 fn needs_extraction(&self, object: &Object, cfg: &IndexedConfig) -> bool {
196 match object {
197 Object::ChunkedBlob(_) => true,
198 Object::Blob(blob) => blob.data.len() as u64 >= cfg.extract_min_bytes,
199 _ => false,
200 }
201 }
202}
203
204pub struct VerifyTimer<R, B, W, X = FailClosedExtraction> {
208 pub remote: R,
210 pub blobs: B,
212 pub windows: W,
214 pub shards: Arc<dyn ShardMap>,
216 pub cfg: IndexedConfig,
218 pub limits: SliceLimits,
220 pub lease: LeaseParams,
222 pub clock: Arc<dyn Clock>,
224 pub metrics: Arc<dyn Metrics>,
226 pub extension: X,
228}
229
230impl<R, B, W, X> core::fmt::Debug for VerifyTimer<R, B, W, X> {
231 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
232 f.debug_struct("VerifyTimer").finish_non_exhaustive()
233 }
234}
235
236impl<S, R, B, W, X> TimerHandler<S> for VerifyTimer<R, B, W, X>
237where
238 S: NamespaceStore,
239 R: NamespaceStore,
240 B: BlobStore,
241 W: PackWindows,
242 X: SliceExtension,
243{
244 fn kind(&self) -> TimerKind {
245 kinds::VERIFY
246 }
247
248 fn max_per_tick(&self) -> Option<u32> {
250 Some(1)
251 }
252
253 fn fire<'a>(
254 &'a self,
255 ctx: &'a TimerCtx<'a, S>,
256 timer: &'a DueTimer,
257 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
258 Box::pin(async move {
259 let Some((name, pack)) = parse_reference(&timer.reference) else {
260 tracing::warn!("malformed verification timer reference dropped");
261 return Ok(Fired::Done(Batch::new()));
262 };
263 let namespace = match ctx.partition {
264 Partition::Namespace(ns) | Partition::Ref { ns, .. } => ns.clone(),
265 _ => {
266 return Err(StoreError::Corrupt(
267 "verification timer in wrong partition".into(),
268 ));
269 }
270 };
271 let budget = SliceBudget::new(self.limits.max_subrequests);
272 let remote = Budgeted::new(&self.remote, &budget);
273 let blobs = Budgeted::new(&self.blobs, &budget);
274 let run = Run {
275 h: self,
276 local: ctx.store,
277 source: ctx.partition,
278 repo: RepoId { namespace, name },
279 pack,
280 budget: &budget,
281 remote: &remote,
282 blobs: &blobs,
283 now: ctx.now_ms,
284 };
285 run.slice(timer).await
286 })
287 }
288}
289
290enum Stop {
292 Store(StoreError),
294 Reject(&'static str),
296 Outcome(Outcome),
298 Wait(u64),
300 Yield(u64),
302 Restart,
304}
305
306impl From<StoreError> for Stop {
307 fn from(error: StoreError) -> Self {
308 Self::Store(error)
309 }
310}
311
312fn unavailable(reason: &'static str) -> Stop {
313 Stop::Store(StoreError::Unavailable(reason.into()))
314}
315
316#[derive(Default)]
319struct Lru {
320 map: BTreeMap<Hash, Arc<Vec<u8>>>,
321 order: VecDeque<Hash>,
322 bytes: u64,
323}
324
325impl Lru {
326 fn insert(&mut self, id: Hash, bytes: Arc<Vec<u8>>) {
327 let len = bytes.len() as u64;
328 if self.map.insert(id, bytes).is_none() {
329 self.order.push_back(id);
330 self.bytes += len;
331 }
332 while self.bytes > CACHE_BYTES && self.order.len() > 1 {
333 if let Some(old) = self.order.pop_front()
334 && let Some(gone) = self.map.remove(&old)
335 {
336 self.bytes -= gone.len() as u64;
337 }
338 }
339 }
340}
341
342struct CacheBases<'a>(&'a Lru);
343
344impl DeltaBaseSource for CacheBases<'_> {
345 const VERIFIED: bool = false;
346 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
347 Ok(self.0.map.get(id).map(|bytes| bytes.as_ref().clone()))
348 }
349}
350
351#[derive(Default)]
353struct SliceState {
354 guards: Vec<Precondition>,
355 job_guard: Option<Value>,
356 cache: Lru,
357 memo: MemberCache,
358 visiting: BTreeSet<(Hash, Hash, u64)>,
359 frames: BTreeMap<Hash, FrameRow>,
360 bases: BTreeMap<Hash, u32>,
361 charged: BTreeSet<(Hash, Hash, u64)>,
362 writes: Vec<Write>,
363 settled: Vec<Write>,
364 entry_idx: u64,
365}
366
367struct Run<'a, S, R, B, W, X> {
368 h: &'a VerifyTimer<R, B, W, X>,
369 local: &'a S,
370 source: &'a Partition,
371 repo: RepoId,
372 pack: Hash,
373 budget: &'a SliceBudget,
374 remote: &'a Budgeted<'a, R>,
375 blobs: &'a Budgeted<'a, B>,
376 now: u64,
377}
378
379fn now_ms(clock: &dyn Clock) -> u64 {
380 u64::try_from(clock.now_ms()).unwrap_or(0)
381}
382
383impl<S, R, B, W, X> Run<'_, S, R, B, W, X>
384where
385 S: NamespaceStore,
386 R: NamespaceStore,
387 B: BlobStore,
388 W: PackWindows,
389 X: SliceExtension,
390{
391 fn decode_limits(&self) -> DecodeLimits {
396 let limits = self.h.limits;
397 super::geometry::decode_limits(limits.resident_bytes, limits.window_bytes)
398 }
399
400 fn deadline(&self) -> u64 {
401 now_ms(self.h.clock.as_ref()).saturating_add(10_000)
402 }
403
404 fn job_key(&self) -> Key {
405 keys::verify_job(&self.repo.name, &self.pack)
406 }
407
408 fn row(&self, sub: u8, id: &Hash) -> Key {
409 keys::verify_row(&self.repo.name, &self.pack, sub, Some(id))
410 }
411
412 fn due(&self, timer: &DueTimer, delay_ms: u64) -> u64 {
414 now_ms(self.h.clock.as_ref())
415 .max(self.now)
416 .saturating_add(delay_ms)
417 .max(timer.due_at_ms.saturating_add(1))
418 }
419
420 async fn slice(&self, timer: &DueTimer) -> Result<Fired, StoreError> {
421 let fired = self.slice_inner(timer).await;
422 if let Err(error) = &fired {
423 tracing::warn!(%error, pack = %mkit_core::hash::to_hex(&self.pack), "verification slice failed");
424 }
425 self.h.metrics.gauge(
426 crate::telemetry::METRIC_INDEX_SLICE_SUBREQUESTS,
427 &[],
428 f64::from(self.budget.used()),
429 );
430 fired
431 }
432
433 #[allow(clippy::too_many_lines)] async fn slice_inner(&self, timer: &DueTimer) -> Result<Fired, StoreError> {
435 let (job, state) =
436 checkpoint::read_job(self.local, self.source, &self.repo.name, &self.pack).await?;
437 let Some((mut job, mut raw)) = job else {
438 return self.cleanup(timer, None, None).await;
439 };
440 if job.gone {
441 return self.cleanup(timer, Some(raw), None).await;
442 }
443 let ticket = match self
444 .local
445 .get(self.source, &keys::ticket(&job.ticket_id))
446 .await?
447 {
448 Some(value) => Some(decode_ticket(&value)?),
449 None if job.phase == Phase::Extract
450 && job.extraction.as_ref().is_some_and(|x| x.object.is_some()) =>
451 {
452 None
453 }
454 None => return self.cleanup(timer, Some(raw), Some(job.ticket_id)).await,
455 };
456 if job.phase == Phase::Watch {
457 return self.watch(
458 timer,
459 job,
460 &raw,
461 state.is_none(),
462 ticket
463 .as_ref()
464 .ok_or_else(|| StoreError::Corrupt("missing watch ticket".into()))?,
465 );
466 }
467 if matches!(
468 job.phase,
469 Phase::Decode | Phase::ClosureResolve | Phase::Recheck
470 ) {
471 match self.begin_attempt(&mut job, &raw).await? {
472 Some(next) => raw = next,
473 None => return Err(StoreError::Unavailable("verification job contended".into())),
474 }
475 }
476 let mut st = SliceState {
477 job_guard: Some(raw.clone()),
478 ..SliceState::default()
479 };
480 let mut held = None;
481 let start = job.clone();
482 let mut ran = job.phase;
483 let mut result = self
484 .step(&mut st, &mut job, state.as_ref(), &mut held)
485 .await;
486 let mut chained = 0;
489 while matches!(result, Ok(0))
490 && ran != Phase::Decode
491 && job.phase != ran
492 && job.phase != Phase::Watch
493 && !matches!(job.phase, Phase::ClosureResolve | Phase::Recheck | Phase::Extract)
495 && chained < 6
496 && self.budget.remaining() >= ENTRY_RESERVE
497 {
498 chained += 1;
499 ran = job.phase;
500 result = self
501 .step(&mut st, &mut job, state.as_ref(), &mut held)
502 .await;
503 }
504 let delay = match result {
505 Ok(delay) | Err(Stop::Yield(delay)) => delay,
506 Err(Stop::Store(error)) => {
507 if is_exhausted(&error) {
508 tracing::warn!(pack = %mkit_core::hash::to_hex(&self.pack), phase = ?job.phase, "verification slice spent its subrequest budget");
509 }
510 return Err(error);
511 }
512 Err(Stop::Wait(delay)) => {
514 job = start;
515 delay
516 }
517 Err(Stop::Restart) => {
518 if job.restarts >= 3 {
519 return Err(StoreError::Unavailable("pack source keeps changing".into()));
520 }
521 job.restart();
522 0
523 }
524 Err(Stop::Outcome(outcome)) => {
525 job.outcome = Some(outcome);
526 job.phase = Phase::Watch;
527 0
528 }
529 Err(Stop::Reject(message)) => {
530 held = Some(self.reject(state.as_ref(), held.as_ref(), message).await?);
531 job.phase = Phase::Watch;
532 0
533 }
534 };
535 job.attempts = 0;
537 self.flush(&mut st).await?;
538 let mut batch = Batch::new()
539 .require(Precondition::NotAfter(self.deadline()))
540 .require(Precondition::Equals(self.job_key(), raw.clone()));
541 let vs = keys::verification(&self.repo.name, &self.pack);
542 batch = batch.require(match held.or_else(|| state.map(|(_, raw)| raw)) {
543 Some(raw) => Precondition::Equals(vs, raw),
544 None => Precondition::Absent(vs),
545 });
546 batch.writes.extend(st.settled);
549 batch.preconditions.extend(st.guards);
550 let batch =
551 checkpoint::write_job(batch, &mut job, Some(&raw), &self.repo.name, &self.pack)?;
552 Ok(self.reschedule(timer, delay, batch))
553 }
554
555 fn watch(
559 &self,
560 timer: &DueTimer,
561 mut job: VerifyJobV1,
562 raw: &Value,
563 vs_missing: bool,
564 ticket: &TicketV1,
565 ) -> Result<Fired, StoreError> {
566 let resume = job.closure_retry();
567 if resume || (vs_missing && job.outcome.is_none()) {
568 if resume {
569 job.phase = Phase::Extract;
570 } else {
571 job.restart();
572 }
573 let batch = Batch::new()
574 .require(Precondition::NotAfter(self.deadline()))
575 .require(Precondition::Equals(self.job_key(), raw.clone()));
576 let batch =
577 checkpoint::write_job(batch, &mut job, Some(raw), &self.repo.name, &self.pack)?;
578 return Ok(self.reschedule(timer, 0, batch));
579 }
580 let delay = ticket
581 .expires_at_ms
582 .saturating_add(1)
583 .saturating_sub(self.now)
584 .min(WATCH_POLL_MS);
585 Ok(self.reschedule(timer, delay, Batch::new()))
586 }
587
588 fn reschedule(&self, timer: &DueTimer, delay_ms: u64, batch: Batch) -> Fired {
589 Fired::Reschedule {
590 due_at_ms: self.due(timer, delay_ms),
591 value: timer.value.clone(),
592 batch,
593 }
594 }
595
596 async fn begin_attempt(
600 &self,
601 job: &mut VerifyJobV1,
602 raw: &Value,
603 ) -> Result<Option<Value>, StoreError> {
604 let (cap, outcome) = if job.phase == Phase::Decode {
605 job.entry_cap = job.entry_cap.min(self.h.limits.max_entries).max(1);
606 (&mut job.entry_cap, Outcome::DecodeBudget)
607 } else {
608 job.closure_cap = job.closure_cap.clamp(1, 4);
609 (&mut job.closure_cap, Outcome::ClosureCapped)
610 };
611 if job.attempts >= ATTEMPTS_PER_CAP {
612 job.attempts = 0;
613 if *cap <= 1 {
614 job.outcome = Some(outcome);
615 job.phase = Phase::Watch;
616 } else {
617 *cap = (*cap / 2).max(1);
618 }
619 }
620 job.attempts += 1;
621 let batch = Batch::new()
622 .require(Precondition::NotAfter(self.deadline()))
623 .require(Precondition::Equals(self.job_key(), raw.clone()));
624 let batch = checkpoint::write_job(batch, job, Some(raw), &self.repo.name, &self.pack)?;
625 let next = encode_job(job);
626 Ok(matches!(
627 self.local.apply(self.source, batch).await?,
628 BatchOutcome::Committed
629 )
630 .then_some(next))
631 }
632
633 async fn reject(
634 &self,
635 state: Option<&(VerificationV1, Value)>,
636 held: Option<&Value>,
637 message: &'static str,
638 ) -> Result<Value, StoreError> {
639 if matches!(state, Some((VerificationV1::Verified { .. }, _))) {
640 return Err(StoreError::Unavailable(
641 "verified pack content unavailable".into(),
642 ));
643 }
644 let rejected = VerificationV1::Rejected {
645 code: "invalid_argument".into(),
646 message: message.into(),
647 };
648 let written = state::write(
649 self.local,
650 self.source,
651 &self.repo.name,
652 &self.pack,
653 held.or_else(|| state.map(|(_, raw)| raw)),
654 &rejected,
655 self.deadline(),
656 )
657 .await?;
658 if !written {
659 tracing::error!(pack = %mkit_core::hash::to_hex(&self.pack), "failed to persist rejected verification state");
660 self.h
661 .metrics
662 .incr(crate::telemetry::METRIC_INDEX_REJECTED_WRITE_FAILED, &[], 1);
663 return Err(StoreError::Unavailable(
664 "rejected state not persisted".into(),
665 ));
666 }
667 Ok(state::encode(&rejected))
668 }
669
670 async fn flush(&self, st: &mut SliceState) -> Result<(), StoreError> {
672 for chunk in std::mem::take(&mut st.writes).chunks(WRITE_BATCH) {
673 let mut batch = Batch::new().require(Precondition::NotAfter(self.deadline()));
674 if let Some(raw) = &st.job_guard {
675 batch = batch.require(Precondition::Equals(self.job_key(), raw.clone()));
676 }
677 batch.writes.extend_from_slice(chunk);
678 if !matches!(
679 self.local.apply(self.source, batch).await?,
680 BatchOutcome::Committed
681 ) {
682 return Err(StoreError::Unavailable(
683 "verification rows contended".into(),
684 ));
685 }
686 }
687 Ok(())
688 }
689
690 #[allow(clippy::too_many_lines)] async fn cleanup(
695 &self,
696 timer: &DueTimer,
697 mut job: Option<Value>,
698 ticket: Option<Hash>,
699 ) -> Result<Fired, StoreError> {
700 let mut peer_guards = Vec::new();
701 if let Some(raw) = &job {
702 let current = checkpoint::decode_job(raw)?;
703 for member in ¤t.extraction_group {
704 if member.pack == self.pack {
705 continue;
706 }
707 let key = keys::verify_job(&self.repo.name, &member.pack);
708 if let Some(raw) = self.local.get(self.source, &key).await? {
709 let peer = checkpoint::decode_job(&raw)?;
710 if peer.extraction_group == current.extraction_group
711 && !peer.gone
712 && !peer.usable()
713 && (peer.outcome.is_none() || peer.closure_retry())
714 {
715 let ticket_key = keys::ticket(&peer.ticket_id);
716 let ticket_raw = self.local.get(self.source, &ticket_key).await?;
717 let live = ticket_raw
718 .as_ref()
719 .map(decode_ticket)
720 .transpose()?
721 .is_some_and(|t| t.expires_at_ms > self.now);
722 if live || peer.extraction.as_ref().is_some_and(|x| x.object.is_some()) {
723 return Ok(self.reschedule(timer, 1_000, Batch::new()));
724 }
725 peer_guards.push(match ticket_raw {
726 Some(raw) => Precondition::Equals(ticket_key, raw),
727 None => Precondition::Absent(ticket_key),
728 });
729 }
730 peer_guards.push(Precondition::Equals(key, raw));
731 } else {
732 peer_guards.push(Precondition::Absent(key));
733 }
734 }
735 }
736 if let Some(raw) = &job
737 && !checkpoint::decode_job(raw)?.gone
738 {
739 let mut gone = VerifyJobV1 {
740 gone: true,
741 members_loaded: true,
742 ..VerifyJobV1::default()
743 };
744 let mut batch = Batch::new()
745 .require(Precondition::NotAfter(self.deadline()))
746 .require(Precondition::Equals(self.job_key(), raw.clone()));
747 batch.preconditions.extend(peer_guards);
748 if let Some(id) = ticket {
749 batch = batch.require(Precondition::Absent(keys::ticket(&id)));
750 }
751 let batch =
752 checkpoint::write_job(batch, &mut gone, Some(raw), &self.repo.name, &self.pack)?;
753 if !matches!(
754 self.local.apply(self.source, batch).await?,
755 BatchOutcome::Committed
756 ) {
757 return Err(StoreError::Unavailable("cleanup header contended".into()));
758 }
759 job = Some(encode_job(&gone));
760 }
761 let (start, end) = keys::verify_range(&self.repo.name, &self.pack, None);
762 for _ in 0..CLEANUP_ROUNDS {
763 let page = self
764 .local
765 .scan(self.source, &start, &end, None, CLEANUP_PAGE)
766 .await?;
767 let mut batch = Batch::new().require(Precondition::NotAfter(self.deadline()));
768 batch = batch.require(match job.as_ref() {
769 Some(raw) => Precondition::Equals(self.job_key(), raw.clone()),
770 None => Precondition::Absent(self.job_key()),
771 });
772 if let Some(id) = ticket {
773 batch = batch.require(Precondition::Absent(keys::ticket(&id)));
774 }
775 for (key, _) in &page.entries {
776 if *key != self.job_key() {
777 batch = batch.delete(key.clone());
778 }
779 }
780 if page.next.is_none() {
781 let member = self
782 .local
783 .has(self.source, &keys::membership(&self.repo.name, &self.pack))
784 .await?;
785 let key = keys::verification(&self.repo.name, &self.pack);
786 if !member && let Some(raw) = self.local.get(self.source, &key).await? {
787 batch = batch
788 .require(Precondition::Absent(keys::membership(
789 &self.repo.name,
790 &self.pack,
791 )))
792 .require(Precondition::Equals(key.clone(), raw))
793 .delete(key);
794 }
795 return Ok(Fired::Done(batch));
796 }
797 if !matches!(
798 self.local.apply(self.source, batch).await?,
799 BatchOutcome::Committed
800 ) {
801 return Err(StoreError::Unavailable(
802 "verification cleanup contended".into(),
803 ));
804 }
805 }
806 Ok(self.reschedule(timer, 0, Batch::new()))
807 }
808}
809
810impl<S, R, B, W, X> Run<'_, S, R, B, W, X>
811where
812 S: NamespaceStore,
813 R: NamespaceStore,
814 B: BlobStore,
815 W: PackWindows,
816 X: SliceExtension,
817{
818 async fn step(
819 &self,
820 st: &mut SliceState,
821 job: &mut VerifyJobV1,
822 state: Option<&(VerificationV1, Value)>,
823 held: &mut Option<Value>,
824 ) -> Result<u64, Stop> {
825 match job.phase {
826 Phase::Decode => self.decode(st, job, state, held).await,
827 Phase::ClosureResolve => {
828 if self.closure(st, job).await? {
829 job.phase = Phase::EmitIndex;
830 job.scan.clear();
831 }
832 Ok(0)
833 }
834 Phase::EmitIndex => self.emit(job).await,
835 Phase::AwaitDelivery => self.await_delivery(job).await,
836 Phase::Extract => {
837 if self.h.extension.extraction_enabled() {
838 return self.extraction(st, job).await;
839 }
840 if job.extract_needed {
841 return Err(Stop::Outcome(Outcome::ExtractionUnavailable));
842 }
843 job.phase = Phase::Verify;
844 Ok(0)
845 }
846 Phase::Verify => self.verify(job, state, held).await,
847 Phase::Recheck => self.recheck(st, job).await,
848 Phase::Watch => Ok(WATCH_POLL_MS),
849 }
850 }
851
852 async fn read(&self, job: &mut VerifyJobV1, offset: u64, len: u64) -> Result<Window, Stop> {
854 let etag = super::etag::resolve(
855 self.local,
856 self.source,
857 &self.repo.name,
858 &self.pack,
859 job.etag.as_deref(),
860 )
861 .await?;
862 self.budget.charge()?;
863 match self
864 .h
865 .windows
866 .read(&self.pack, offset, len, etag.as_deref())
867 .await
868 {
869 Ok(window) => {
870 if job.etag.is_none() {
871 job.etag = Some(
872 super::etag::capture(
873 self.local,
874 self.source,
875 &self.repo.name,
876 &self.pack,
877 &window.etag,
878 self.deadline(),
879 )
880 .await?,
881 );
882 }
883 Ok(window)
884 }
885 Err(WindowError::EtagChanged) => Err(Stop::Restart),
886 Err(WindowError::Missing | WindowError::Unavailable) => {
887 Err(unavailable("pack window read failed"))
888 }
889 }
890 }
891
892 fn reader_error(job: &VerifyJobV1, error: &PackError) -> Stop {
893 match error {
894 PackError::PackfileTooLarge => Stop::Outcome(Outcome::DecodeBudget),
895 _ if !job.cursor.is_empty() && job.restarts == 0 => Stop::Restart,
896 _ => Stop::Reject("object hash mismatch"),
897 }
898 }
899
900 fn packlist(job: &mut VerifyJobV1, window: &Window, pack: &Hash) -> Result<u64, Stop> {
904 job.kind = Kind::Packlist;
905 if window.bytes.len() as u64 != job.pack_len {
906 return Err(Stop::Outcome(Outcome::ClosureCapped));
907 }
908 if hash(&window.bytes) != *pack {
909 return Err(Stop::Reject("object hash mismatch"));
910 }
911 let list =
912 decode_packlist(&window.bytes).map_err(|_| Stop::Reject("object hash mismatch"))?;
913 if list.packs.len() > index::MAX_LOOKUP_IDS + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
914 {
915 return Err(Stop::Outcome(Outcome::ClosureCapped));
916 }
917 job.packlist_prev = list.prev;
918 job.packlist = list.packs;
919 job.phase = Phase::Verify;
920 Ok(0)
921 }
922
923 #[allow(clippy::too_many_lines)] async fn decode(
925 &self,
926 st: &mut SliceState,
927 job: &mut VerifyJobV1,
928 state: Option<&(VerificationV1, Value)>,
929 held: &mut Option<Value>,
930 ) -> Result<u64, Stop> {
931 if job.pack_len > self.h.cfg.max_pack_bytes {
934 return Err(Stop::Reject(super::PACK_CAP_MESSAGE));
935 }
936 if job.kind == Kind::Unknown {
939 let start = keys::verify_row(&self.repo.name, &self.pack, keys::VC_FRAME, None);
940 let (_, end) = keys::verify_range(&self.repo.name, &self.pack, None);
941 let page = self
942 .local
943 .scan(self.source, &start, &end, None, CLEANUP_PAGE)
944 .await?;
945 st.settled
946 .extend(page.entries.into_iter().map(|(key, _)| Write::Delete(key)));
947 if !st.settled.is_empty() {
948 return Ok(0);
949 }
950 }
951 match state {
952 Some((VerificationV1::Rejected { .. }, _)) => {
953 job.phase = Phase::Watch;
954 return Ok(0);
955 }
956 Some((VerificationV1::Verified { pack_len, .. }, _)) if *pack_len != job.pack_len => {
957 return Err(Stop::Store(StoreError::Corrupt(
958 "verified pack length changed".into(),
959 )));
960 }
961 Some((VerificationV1::Verified { .. }, _)) => {}
963 other => {
964 let pending = VerificationV1::Pending {
965 lease_until_ms: now_ms(self.h.clock.as_ref())
966 .saturating_add(state::VERIFICATION_LEASE_MS),
967 };
968 let prior = other.map(|(_, raw)| raw);
969 if !state::write(
970 self.local,
971 self.source,
972 &self.repo.name,
973 &self.pack,
974 prior,
975 &pending,
976 self.deadline(),
977 )
978 .await?
979 {
980 return Err(unavailable("verification state contended"));
981 }
982 *held = Some(state::encode(&pending));
983 }
984 }
985 let window_bytes = self.h.limits.window_bytes;
986 let mut preloaded = None;
987 if job.kind == Kind::Unknown {
988 let window = self.read(job, 0, job.pack_len.min(window_bytes)).await?;
989 match classify::classify(&window.bytes) {
990 Ok(UploadType::Packlist) => return Self::packlist(job, &window, &self.pack),
991 Ok(UploadType::Pack) => {
992 job.kind = Kind::Pack;
993 job.version = window
994 .bytes
995 .get(4..8)
996 .and_then(|v| v.try_into().ok())
997 .map_or(0, u32::from_le_bytes);
998 preloaded = Some(window);
999 }
1000 Err(_) => return Err(Stop::Reject("unknown upload type")),
1001 }
1002 }
1003 let limits = self.decode_limits();
1004 let mut reader = if job.cursor.is_empty() {
1005 WindowReader::new(job.pack_len, window_bytes, limits, Some(self.pack))
1006 } else {
1007 WindowCursor::from_bytes(&job.cursor)
1008 .and_then(|cursor| WindowReader::resume(&cursor, limits))
1009 }
1010 .map_err(|e| Self::reader_error(job, &e))?;
1011 let (mut fed, mut processed) = (0_u32, 0_u32);
1012 loop {
1013 match reader.step().map_err(|e| Self::reader_error(job, &e))? {
1014 Step::NeedWindow(request) => {
1015 let window = match preloaded.take() {
1016 Some(window)
1017 if request.offset == 0 && window.bytes.len() as u64 == request.len =>
1018 {
1019 window
1020 }
1021 _ => self.read(job, request.offset, request.len).await?,
1022 };
1023 reader
1024 .feed_owned(request.offset, window.bytes)
1025 .map_err(|e| Self::reader_error(job, &e))?;
1026 fed += 1;
1027 job.windows_done = job.windows_done.saturating_add(1);
1028 }
1029 Step::Entry(entry) => {
1030 let frame = reader
1031 .last_frame()
1032 .ok_or_else(|| unavailable("window reader lost its frame"))?;
1033 #[cfg(feature = "pack-ruzstd")]
1038 if matches!(entry, PackEntry::Delta { .. }) {
1039 let cursor = reader
1040 .checkpoint()
1041 .ok_or_else(|| unavailable("delta entry lost its boundary"))?;
1042 drop(reader);
1043 self.entry(st, job, frame, entry).await?;
1044 job.cursor = cursor.to_bytes();
1045 job.attempts = 0;
1046 return Ok(0);
1047 }
1048 self.entry(st, job, frame, entry).await?;
1049 processed += 1;
1050 if (fed >= 2
1053 || processed >= job.entry_cap
1054 || self.budget.remaining() < ENTRY_RESERVE)
1055 && let Some(cursor) = reader.checkpoint()
1056 {
1057 job.cursor = cursor.to_bytes();
1058 job.attempts = 0;
1059 return Ok(0);
1060 }
1061 }
1062 Step::Done(summary) => {
1063 if u64::from(summary.entry_count) != job.entries {
1064 return Err(Stop::Reject("object hash mismatch"));
1065 }
1066 if job.bad_signature {
1067 return Err(Stop::Reject("bad signature"));
1068 }
1069 job.cursor.clear();
1070 job.attempts = 0;
1071 job.scan.clear();
1072 job.owed = 0;
1073 job.phase = Phase::ClosureResolve;
1074 return Ok(0);
1075 }
1076 _ => return Err(unavailable("unexpected window reader step")),
1077 }
1078 }
1079 }
1080
1081 async fn frame_row(&self, st: &mut SliceState, id: &Hash) -> Result<Option<FrameRow>, Stop> {
1082 if let Some(row) = st.frames.get(id) {
1083 return Ok(Some(*row));
1084 }
1085 let Some(value) = self
1086 .local
1087 .get(self.source, &self.row(keys::VC_FRAME, id))
1088 .await?
1089 else {
1090 return Ok(None);
1091 };
1092 let row = decode_frame(id, &value)?;
1093 st.frames.insert(*id, row);
1094 Ok(Some(row))
1095 }
1096
1097 async fn base_depth(&self, st: &mut SliceState, id: &Hash) -> Result<u32, Stop> {
1098 if let Some(depth) = st.bases.get(id) {
1099 return Ok(*depth);
1100 }
1101 let depth = match self
1102 .local
1103 .get(self.source, &self.row(keys::VC_BASE, id))
1104 .await?
1105 {
1106 Some(value) => decode_base(&value)?.depth,
1107 None => 0,
1108 };
1109 st.bases.insert(*id, depth);
1110 Ok(depth)
1111 }
1112
1113 #[allow(clippy::too_many_lines)] async fn entry(
1116 &self,
1117 st: &mut SliceState,
1118 job: &mut VerifyJobV1,
1119 frame: mkit_core::pack::window::FrameInfo,
1120 entry: PackEntry<'static>,
1121 ) -> Result<(), Stop> {
1122 let cap = self.h.cfg.max_delta_chain_depth;
1123 st.entry_idx = job.entries;
1124 st.memo = MemberCache::default();
1125 if frame.length > super::geometry::FRAME_BYTES {
1126 return Err(Stop::Outcome(Outcome::DecodeBudget));
1127 }
1128 let base = match &entry {
1129 PackEntry::Delta { base, .. } => Some(*base),
1130 PackEntry::Raw { .. } => None,
1131 };
1132 let (hops, external) = match base {
1133 None => (0, None),
1134 Some(b) => match self
1135 .frame_row(st, &b)
1136 .await?
1137 .filter(|row| row.value.frame_offset < frame.offset)
1138 {
1139 Some(row) => (row.value.chain_depth.saturating_add(1), row.external),
1140 None => (1, Some(b)),
1141 },
1142 };
1143 if hops > cap {
1144 return Err(Stop::Reject("delta chain too deep"));
1145 }
1146 if let Some(b) = base {
1147 self.ensure_base(st, job, b, frame.offset).await?;
1148 }
1149 if let Some(x) = external
1150 && hops.saturating_add(self.base_depth(st, &x).await?) > cap
1151 {
1152 return Err(Stop::Outcome(Outcome::ExternalTooDeep));
1153 }
1154 let limits = self.decode_limits();
1155 let (id, bytes) =
1156 decode_entry_with(entry, &mut CacheBases(&st.cache), limits).map_err(|e| {
1157 if matches!(e, PackError::PackfileTooLarge) {
1158 Stop::Outcome(Outcome::DecodeBudget)
1159 } else {
1160 Stop::Reject("object hash mismatch")
1161 }
1162 })?;
1163 crate::takedown::denial::require_clear(self.remote, &id)
1164 .await
1165 .map_err(|error| {
1166 if error.code() == crate::Code::PermissionDenied {
1167 Stop::Outcome(Outcome::Blocked)
1168 } else {
1169 unavailable("decoded object unavailable")
1170 }
1171 })?;
1172 let object = mkit_core::serialize::deserialize(&bytes)
1173 .map_err(|_| Stop::Reject("object hash mismatch"))?;
1174 crate::takedown::inventory::stage(
1175 self.remote,
1176 &self.pack,
1177 job.pack_len,
1178 &id,
1179 &object,
1180 base,
1181 self.now,
1182 )
1183 .await?;
1184 let size = bytes.len() as u64;
1185 let existing = self.frame_row(st, &id).await?;
1186 if existing.is_none_or(|row| row.value.frame_offset == frame.offset) {
1189 job.in_pack_bytes = job.in_pack_bytes.saturating_add(size);
1190 let budget = self.h.cfg.decode_budget;
1191 if job.in_pack_bytes > budget {
1192 return Err(Stop::Reject("pack exceeds indexed decode budget"));
1193 }
1194 if job.in_pack_bytes.saturating_add(job.external_bytes) > budget {
1195 return Err(Stop::Outcome(Outcome::DecodeBudget));
1196 }
1197 let row = FrameRow {
1198 value: IndexValue {
1199 frame_offset: frame.offset,
1200 frame_length: frame.length,
1201 wire_type: frame.wire_type,
1202 decoded_size: size,
1203 chain_depth: hops,
1204 delta_base: base,
1205 },
1206 object_type: object.object_type() as u8,
1207 external,
1208 };
1209 let encoded =
1210 encode_frame(&id, &row).map_err(|_| Stop::Reject("object hash mismatch"))?;
1211 st.writes
1212 .push(Write::Put(self.row(keys::VC_FRAME, &id), encoded));
1213 st.frames.insert(id, row);
1214 let fact = super::selection::SelectionFact::from_object(&object);
1215 let projection = super::selection::Projection::from_fact(id, &fact);
1216 for (index, references) in fact
1217 .references()
1218 .chunks(super::selection::REFERENCES_PER_PAGE)
1219 .enumerate()
1220 {
1221 let index = u32::try_from(index).expect("decoded entry bounds page count");
1222 st.writes.push(Write::Put(
1223 self.row(keys::VC_CANDIDATE, &projection.page_id(index)),
1224 projection.encode_page(index, references),
1225 ));
1226 if st.writes.len() >= WRITE_BATCH {
1227 self.flush(st).await?;
1228 }
1229 }
1230 drop(fact);
1231 st.writes.push(Write::Put(
1235 self.row(keys::VC_CANDIDATE, &id),
1236 projection.encode(),
1237 ));
1238 if let Some(parents) = super::verify::history_parents(&object) {
1239 st.writes.push(Write::Put(
1240 self.row(keys::VC_HISTORY, &id),
1241 Value::new(parents.concat()),
1242 ));
1243 }
1244 for child in children(&object, ClosureMode::History) {
1245 st.writes.push(Write::Put(
1246 self.row(keys::VC_CHILD, &child),
1247 Value::default(),
1248 ));
1249 if st.writes.len() >= WRITE_BATCH {
1250 self.flush(st).await?;
1251 }
1252 }
1253 if verify_object_signature(&object).is_err() {
1254 job.bad_signature = true;
1255 }
1256 if self.h.extension.needs_extraction(&object, &self.h.cfg) {
1257 job.extract_needed = true;
1258 }
1259 }
1260 if size <= super::geometry::ENTRY_CACHE_BYTES {
1261 st.cache.insert(id, Arc::from(bytes));
1262 }
1263 job.entries += 1;
1264 Ok(())
1265 }
1266
1267 fn ensure_base<'x>(
1270 &'x self,
1271 st: &'x mut SliceState,
1272 job: &'x mut VerifyJobV1,
1273 base: Hash,
1274 before: u64,
1275 ) -> BoxFuture<'x, Result<(), Stop>> {
1276 Box::pin(async move {
1277 if st.cache.map.contains_key(&base) {
1278 return Ok(());
1279 }
1280 match self
1281 .frame_row(st, &base)
1282 .await?
1283 .filter(|row| row.value.frame_offset < before)
1284 {
1285 Some(row) => {
1286 let limits = self.decode_limits();
1287 if row.value.frame_length > super::geometry::FRAME_BYTES
1288 || row.value.decoded_size > limits.max_decoded_bytes
1289 || row.value.chain_depth > self.h.cfg.max_delta_chain_depth
1290 {
1291 return Err(StoreError::Corrupt("invalid source frame".into()).into());
1292 }
1293 if let Some(next) = row.value.delta_base {
1294 if self.frame_row(st, &next).await?.is_some_and(|p| {
1295 p.value.frame_offset < row.value.frame_offset
1296 && p.value.chain_depth >= row.value.chain_depth
1297 }) {
1298 return Err(StoreError::Corrupt("invalid source chain".into()).into());
1299 }
1300 self.ensure_base(st, job, next, row.value.frame_offset)
1301 .await?;
1302 }
1303 let window = self
1304 .read(job, row.value.frame_offset, row.value.frame_length)
1305 .await?;
1306 let (id, bytes) = decode_frame_with(
1307 &window.bytes,
1308 job.version,
1309 &mut CacheBases(&st.cache),
1310 limits,
1311 )
1312 .map_err(|_| Stop::Restart)?;
1313 if id != base {
1314 return Err(Stop::Restart);
1315 }
1316 st.cache.insert(base, Arc::from(bytes));
1317 Ok(())
1318 }
1319 None => self.resolve_external(st, job, base).await,
1320 }
1321 })
1322 }
1323
1324 fn missing(&self, job: &VerifyJobV1) -> Stop {
1325 if resolve::lagged(self.now, job.created_at_ms, self.h.cfg.relay_lag_bound_ms) {
1326 Stop::Wait(LAG_BACKOFF_MS)
1327 } else {
1328 Stop::Outcome(Outcome::BaseMissing)
1329 }
1330 }
1331
1332 async fn resolve_external(
1335 &self,
1336 st: &mut SliceState,
1337 job: &mut VerifyJobV1,
1338 base: Hash,
1339 ) -> Result<(), Stop> {
1340 if let Some((_, (bytes, _))) = st.memo.rows().find(|((id, ..), _)| *id == base) {
1342 st.cache.insert(base, Arc::from(bytes.to_vec()));
1343 return Ok(());
1344 }
1345 let found = resolve::locate_split(
1346 self.remote,
1347 self.h.shards.as_ref(),
1348 &self.repo,
1349 &[base],
1350 self.h.metrics.as_ref(),
1351 )
1352 .await
1353 .map_err(|_| unavailable("index lookup failed"))?;
1354 let located: LocatedObject = match found.get(&base) {
1355 Some(Ok(Some(located))) => *located,
1356 Some(Err(_)) => return Err(Stop::Outcome(Outcome::BaseCapped)),
1357 _ => return Err(self.missing(job)),
1358 };
1359 let cfg = &self.h.cfg;
1360 let memo_budget = self.decode_limits().max_decoded_bytes;
1364 let (canonical, _) = resolve::member_object(
1365 self.blobs,
1366 self.remote,
1367 self.h.shards.as_ref(),
1368 &self.repo,
1369 base,
1370 located,
1371 cfg.max_delta_chain_depth,
1372 memo_budget,
1373 &mut st.memo,
1374 &mut st.visiting,
1375 self.h.metrics.as_ref(),
1376 )
1377 .await
1378 .map_err(|failure| match failure {
1379 ResolveFailure::Missing => self.missing(job),
1380 ResolveFailure::Capped => Stop::Outcome(Outcome::BaseCapped),
1381 ResolveFailure::Corrupt(_) => unavailable("member content unavailable"),
1382 ResolveFailure::Other(error) => match error.public_message() {
1383 "pack exceeds indexed decode budget" => Stop::Outcome(Outcome::DecodeBudget),
1384 "delta chain too deep" => Stop::Outcome(Outcome::ExternalTooDeep),
1385 "object blocked" => Stop::Outcome(Outcome::Blocked),
1386 _ => unavailable("member content unavailable"),
1387 },
1388 })?;
1389 st.cache.insert(base, Arc::from(canonical.to_vec()));
1390 self.charge_bases(st, job).await?;
1391 st.memo = MemberCache::default();
1395 Ok(())
1396 }
1397
1398 async fn charge_bases(&self, st: &mut SliceState, job: &mut VerifyJobV1) -> Result<(), Stop> {
1402 let fresh: Vec<_> = st
1403 .memo
1404 .rows()
1405 .filter(|(location, _)| !st.charged.contains(*location))
1406 .map(|(location, (bytes, depth))| (*location, bytes.len() as u64, *depth))
1407 .collect();
1408 for ((id, pack, offset), size, depth) in fresh {
1409 crate::takedown::inventory::dependency(
1410 self.remote,
1411 &self.pack,
1412 job.pack_len,
1413 &id,
1414 self.now,
1415 )
1416 .await?;
1417 st.writes.push(Write::Put(
1418 self.row(keys::VC_DEPENDENCY, &pack),
1419 Value::default(),
1420 ));
1421 st.charged.insert((id, pack, offset));
1422 let location = mkit_core::hash::domain_digest(
1423 b"mkit:vc-base:v1\0",
1424 &[
1425 id.as_slice(),
1426 pack.as_slice(),
1427 offset.to_be_bytes().as_slice(),
1428 ]
1429 .concat(),
1430 );
1431 let prior = match self
1432 .local
1433 .get(self.source, &self.row(keys::VC_BASE, &location))
1434 .await?
1435 {
1436 Some(value) => Some(decode_base(&value)?),
1437 None => None,
1438 };
1439 if prior.is_none_or(|row| row.entry == st.entry_idx) {
1440 job.external_bytes = job.external_bytes.saturating_add(size);
1441 st.writes.push(Write::Put(
1442 self.row(keys::VC_BASE, &location),
1443 encode_base(&BaseRow {
1444 size,
1445 depth,
1446 entry: st.entry_idx,
1447 }),
1448 ));
1449 st.bases.insert(id, depth);
1450 } else if let Some(row) = prior {
1451 st.bases.insert(id, row.depth);
1452 }
1453 st.writes.push(Write::Put(
1456 self.row(keys::VC_BASE, &id),
1457 encode_base(&BaseRow {
1458 size: 0,
1459 depth,
1460 entry: st.entry_idx,
1461 }),
1462 ));
1463 }
1464 if job.in_pack_bytes.saturating_add(job.external_bytes) > self.h.cfg.decode_budget {
1465 return Err(Stop::Outcome(Outcome::DecodeBudget));
1466 }
1467 Ok(())
1468 }
1469
1470 async fn closure(&self, st: &mut SliceState, job: &mut VerifyJobV1) -> Result<bool, Stop> {
1476 let (start, end) = keys::verify_range(&self.repo.name, &self.pack, Some(keys::VC_CHILD));
1477 for _ in 0..job.closure_cap {
1478 if self.budget.remaining() < ENTRY_RESERVE
1479 || st.settled.len() + CLOSURE_CHUNK as usize > WRITE_BATCH
1480 {
1481 return Ok(false);
1482 }
1483 let cursor = (!job.scan.is_empty()).then(|| Cursor::new(job.scan.clone()));
1484 let page = self
1485 .local
1486 .scan(self.source, &start, &end, cursor.as_ref(), CLOSURE_CHUNK)
1487 .await?;
1488 let mut ids = Vec::new();
1489 for (key, _) in &page.entries {
1490 let Some(keys::ParsedKey::VerifyCursor { id: Some(id), .. }) = keys::parse(key)
1491 else {
1492 return Err(Stop::Store(StoreError::Corrupt(
1493 "bad owed child row".into(),
1494 )));
1495 };
1496 ids.push(id);
1497 }
1498 if !ids.is_empty() {
1499 let frame_keys: Vec<_> =
1500 ids.iter().map(|id| self.row(keys::VC_FRAME, id)).collect();
1501 let present = self.local.get_many(self.source, &frame_keys).await?;
1502 let mut wanted = Vec::new();
1503 for (id, row) in ids.iter().zip(present) {
1504 if row.is_some() {
1505 st.settled.push(Write::Delete(self.row(keys::VC_CHILD, id)));
1506 } else {
1507 wanted.push(*id);
1508 }
1509 }
1510 if !wanted.is_empty() {
1511 let reserve = self.h.limits.max_subrequests.min(
1514 u32::try_from(index::MAX_LOOKUP_PAGES + index::MAX_LOOKUP_MEMBERSHIP_READS)
1515 .unwrap_or(u32::MAX),
1516 );
1517 if self.budget.remaining() < reserve {
1518 return Ok(false);
1519 }
1520 let found = resolve::locate_split(
1521 self.remote,
1522 self.h.shards.as_ref(),
1523 &self.repo,
1524 &wanted,
1525 self.h.metrics.as_ref(),
1526 )
1527 .await
1528 .map_err(|_| {
1529 if self.budget.remaining() == 0 {
1530 Stop::Outcome(Outcome::ClosureCapped)
1531 } else {
1532 unavailable("index lookup failed")
1533 }
1534 })?;
1535 for id in wanted {
1536 match found.get(&id) {
1537 Some(Ok(Some(located))) => {
1538 if !job.satisfying.contains(&located.pack) {
1539 if job.satisfying.len() >= MAX_SATISFYING {
1540 return Err(Stop::Outcome(Outcome::ClosureCapped));
1541 }
1542 job.satisfying.push(located.pack);
1543 }
1544 st.settled
1545 .push(Write::Delete(self.row(keys::VC_CHILD, &id)));
1546 }
1547 Some(Err(_)) => return Err(Stop::Outcome(Outcome::ClosureCapped)),
1548 _ => job.owed += 1,
1549 }
1550 }
1551 }
1552 }
1553 let Some(next) = page.next else {
1554 job.scan.clear();
1555 return Ok(true);
1556 };
1557 job.scan = next.into_bytes().to_vec();
1558 }
1559 Ok(false)
1560 }
1561
1562 async fn emit(&self, job: &mut VerifyJobV1) -> Result<u64, Stop> {
1565 let (start, end) = keys::verify_range(&self.repo.name, &self.pack, Some(keys::VC_FRAME));
1566 let cursor = (!job.scan.is_empty()).then(|| Cursor::new(job.scan.clone()));
1567 let page = self
1568 .local
1569 .scan(self.source, &start, &end, cursor.as_ref(), EMIT_PAGE)
1570 .await?;
1571 let mut entries: Vec<IndexEntry> = Vec::with_capacity(page.entries.len());
1572 for (key, value) in &page.entries {
1573 let Some(keys::ParsedKey::VerifyCursor { id: Some(id), .. }) = keys::parse(key) else {
1574 return Err(Stop::Store(StoreError::Corrupt("bad frame row".into())));
1575 };
1576 entries.push(checkpoint::index_entry(id, &decode_frame(&id, value)?));
1577 }
1578 let clock = self.h.clock.as_ref();
1579 let plan = index::plan_index_rows(
1580 self.h.shards.as_ref(),
1581 &self.repo,
1582 self.source,
1583 &self.pack,
1584 &entries,
1585 now_ms(clock),
1586 )?;
1587 for direct in plan.direct {
1588 let mut batch = Batch::new().require(Precondition::NotAfter(self.deadline()));
1589 for (key, value) in direct.puts {
1590 if let Some(keys::ParsedKey::ObjectIndex { object, .. }) = keys::parse(&key) {
1591 crate::takedown::denial::require_clear(self.remote, &object)
1592 .await
1593 .map_err(|error| {
1594 if error.code() == crate::Code::PermissionDenied {
1595 Stop::Outcome(Outcome::Blocked)
1596 } else {
1597 unavailable("index object unavailable")
1598 }
1599 })?;
1600 }
1601 batch = batch.put(key, value);
1602 }
1603 if !matches!(
1604 self.local.apply(&direct.target, batch).await?,
1605 BatchOutcome::Committed
1606 ) {
1607 return Err(unavailable("index rows contended"));
1608 }
1609 }
1610 if !plan.relay.is_empty() {
1611 let lease = if matches!(self.source, Partition::Ref { .. }) {
1612 Some(
1613 renew_for_relay(
1614 self.local,
1615 self.remote,
1616 self.h.shards.as_ref(),
1617 clock,
1618 self.h.metrics.as_ref(),
1619 &self.repo,
1620 self.source,
1621 &self.h.lease,
1622 )
1623 .await
1624 .map_err(|_| unavailable("epoch lease renewal failed"))?,
1625 )
1626 } else {
1627 None
1628 };
1629 commit_relay_rows(
1630 self.local,
1631 self.source,
1632 &plan.relay,
1633 now_ms(clock),
1634 self.deadline(),
1635 lease.as_ref(),
1636 )
1637 .await?;
1638 job.last_relay_seq = self
1639 .local
1640 .get(self.source, &keys::outbox_sequence())
1641 .await?
1642 .as_ref()
1643 .map(codec::decode_u64)
1644 .transpose()?;
1645 }
1646 if let Some(next) = page.next {
1647 job.scan = next.into_bytes().to_vec();
1648 } else {
1649 job.scan.clear();
1650 job.phase = Phase::AwaitDelivery;
1651 }
1652 Ok(0)
1653 }
1654
1655 async fn await_delivery(&self, job: &mut VerifyJobV1) -> Result<u64, Stop> {
1657 let delivered = match job.last_relay_seq {
1658 Some(seq) => relay_delivered_through(self.local, self.source, seq).await?,
1659 None => true,
1660 };
1661 if delivered {
1662 job.phase = Phase::Extract;
1663 return Ok(0);
1664 }
1665 Ok(2_000)
1666 }
1667
1668 async fn verify(
1670 &self,
1671 job: &mut VerifyJobV1,
1672 state: Option<&(VerificationV1, Value)>,
1673 held: &mut Option<Value>,
1674 ) -> Result<u64, Stop> {
1675 let now = now_ms(self.h.clock.as_ref());
1676 if job.kind == Kind::Packlist {
1677 crate::takedown::inventory::stage_packlist(
1678 self.remote,
1679 &self.pack,
1680 job.pack_len,
1681 job.packlist_prev,
1682 &job.packlist,
1683 now,
1684 )
1685 .await?;
1686 let mut offset = if job.scan.is_empty() {
1688 0
1689 } else {
1690 codec::decode_u64(&Value::new(job.scan.clone()))?
1691 };
1692 let start =
1693 usize::try_from(offset).map_err(|_| unavailable("invalid inventory cursor"))?;
1694 let children = job
1695 .packlist
1696 .get(start..)
1697 .ok_or_else(|| unavailable("invalid inventory cursor"))?;
1698 for child in children {
1699 if self.budget.remaining() < ENTRY_RESERVE {
1700 return Ok(1);
1701 }
1702 crate::takedown::inventory::dependency(
1703 self.remote,
1704 &self.pack,
1705 job.pack_len,
1706 child,
1707 now,
1708 )
1709 .await?;
1710 offset += 1;
1711 job.scan = offset.to_be_bytes().to_vec();
1712 }
1713 }
1714 crate::takedown::inventory::complete(self.remote, &self.pack, job.pack_len, now).await?;
1715 job.scan.clear();
1716 match state {
1717 Some((VerificationV1::Rejected { .. }, _)) => {
1718 job.phase = Phase::Watch;
1719 return Ok(0);
1720 }
1721 Some((VerificationV1::Verified { pack_len, .. }, _)) if *pack_len == job.pack_len => {}
1722 Some((VerificationV1::Verified { .. }, _)) => {
1723 return Err(Stop::Store(StoreError::Corrupt(
1724 "verified pack length changed".into(),
1725 )));
1726 }
1727 other => {
1728 let verified = VerificationV1::Verified {
1729 pack_len: job.pack_len,
1730 verified_at_ms: now,
1731 publication: None,
1732 };
1733 let prior = other.map(|(_, raw)| raw);
1734 if !state::write(
1735 self.local,
1736 self.source,
1737 &self.repo.name,
1738 &self.pack,
1739 prior,
1740 &verified,
1741 self.deadline(),
1742 )
1743 .await?
1744 {
1745 return Err(unavailable("verification state contended"));
1746 }
1747 *held = Some(state::encode(&verified));
1748 }
1749 }
1750 if job.kind == Kind::Pack && job.owed > 0 {
1751 job.phase = Phase::Recheck;
1752 } else {
1753 job.closure_final_at_ms = Some(now);
1754 job.phase = Phase::Watch;
1755 }
1756 Ok(0)
1757 }
1758
1759 async fn recheck(&self, st: &mut SliceState, job: &mut VerifyJobV1) -> Result<u64, Stop> {
1762 let now = now_ms(self.h.clock.as_ref());
1763 if !job.final_pass {
1764 let end = job
1765 .created_at_ms
1766 .saturating_add(self.h.cfg.relay_lag_bound_ms);
1767 if now < end {
1768 return Ok(end - now + 1);
1769 }
1770 job.final_pass = true;
1771 job.owed = 0;
1772 job.scan.clear();
1773 }
1774 if self.closure(st, job).await? {
1775 job.closure_final_at_ms = Some(now);
1776 job.phase = Phase::Watch;
1777 }
1778 Ok(0)
1779 }
1780}