Skip to main content

mkit_server/indexed/
job.rs

1//! Scheduled verification: the kind-7 slice machine (WP-4.8, R-171).
2//!
3//! A ticketed pack verifies in checkpointed slices, one per alarm fire, with
4//! repository-isolated checks; Scheduled error priority and duplicate budget
5//! charging follow R-171's stricter deployment rules. Each fire runs one phase
6//! step of a [`VerifyJobV1`] and ends in `Fired::Reschedule`, whose batch
7//! carries the job put guarded by the job row (and `vs`) as it was read, so a
8//! duplicate fire loses. Rows a slice writes ahead of that batch (frames,
9//! owed children, charged bases) are pure functions of the pack, so a crash
10//! only replays them. Phases: `Decode`, `ClosureResolve`, `EmitIndex`,
11//! `AwaitDelivery`, `Extract`, `Verify`, `Recheck`, `Watch`.
12//!
13//! Index rows are emitted only after the decode reaches `Done`
14//! (SPEC-PACKFILE §11: entries are provisional until then), and `Verified` is
15//! written only after the relay delivered them (R-130). A failure caused by
16//! repository membership or the platform is a terminal [`Outcome`], never a
17//! persisted `Rejected` (R-148).
18
19use 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/// Fixed work units of one slice. A Worker fixes them for its plan.
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub struct SliceLimits {
64    /// One window of the resumable decoder.
65    pub window_bytes: u64,
66    /// Bytes one slice may hold: window, carried entry, decoded entry, the
67    /// in-pack cache and retained external bases. A single object over it
68    /// gets `pack exceeds indexed decode budget` (a documented deployment
69    /// limit, SPEC-SERVER §9.8).
70    pub resident_bytes: u64,
71    /// Calls that leave the object: R2 ranges and index or membership shards.
72    pub max_subrequests: u32,
73    /// Entries one slice may decode before its next checkpoint.
74    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/// R-203: old/new ring growth is at most three 8 MiB windows; four MiB
89/// covers bounded blocks and decoder tables. Decoder and object parsing do
90/// not overlap. Reserve this scratch while retaining one reader window.
91#[cfg(feature = "pack-ruzstd")]
92const DECODER_SCRATCH_BYTES: u64 = 3 * (8 << 20) + (4 << 20);
93/// A 16 MiB window, 28 MiB decoder, two delta streams and 1 MiB metadata
94/// leave an LRU under 1 MiB; always retain its required newest base.
95#[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;
103/// Slice failures on one cursor before its entry cap halves.
104const ATTEMPTS_PER_CAP: u32 = 3;
105/// Subrequests kept back when a slice decides to fetch its next entry.
106const ENTRY_RESERVE: u32 = 64;
107/// Each closure lookup has its own durable id boundary.
108const CLOSURE_CHUNK: u32 = 1;
109/// Frame rows per index emission slice.
110const EMIT_PAGE: u32 = 64;
111/// Rows per cleanup batch.
112const CLEANUP_PAGE: u32 = 90;
113/// Cleanup pages per fire.
114const CLEANUP_ROUNDS: usize = 16;
115/// Buffered idempotent rows per batch.
116const WRITE_BATCH: usize = 90;
117/// Retry wait while a base or child is inside the membership lag window.
118const LAG_BACKOFF_MS: u64 = 15_000;
119/// The longest a finished job waits before checking whether its ticket
120/// closed: an hour, and never past the ticket's own expiry.
121const WATCH_POLL_MS: u64 = 3_600_000;
122/// Distinct member packs one job may depend on.
123const MAX_SATISFYING: usize = index::MAX_LOOKUP_IDS;
124
125/// The existing extraction extension of the `Extract` phase. Enabled Workers
126/// use the internal upload callbacks; the disconnected default fails closed.
127pub trait SliceExtension: crate::MaybeSend + crate::MaybeSync {
128    /// Whether `object` must be extracted into the object store before its
129    /// pack is `Verified`.
130    fn needs_extraction(&self, object: &Object, cfg: &IndexedConfig) -> bool;
131
132    /// Whether this existing extension supplies the internal upload callbacks.
133    fn extraction_enabled(&self) -> bool {
134        false
135    }
136
137    /// Start a root-pinned private object session after all source CVs verify.
138    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    /// Verify and commit one bounded private part, returning its opaque receipt.
151    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    /// Publish only the exact root-pinned parts selected by these receipts.
165    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    /// Abort the private session; the default extension has no session.
178    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/// Every `ChunkedBlob` and every Blob of at least `extract_min_bytes` needs
190/// extraction, and there is none yet: the job ends `ExtractionUnavailable`.
191#[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
204/// The kind-7 handler. `remote` reaches other partitions (the Worker's
205/// namespace client, budgeted per fire); the fire's own store serves the ref
206/// shard's rows, timers and outbox.
207pub struct VerifyTimer<R, B, W, X = FailClosedExtraction> {
208    /// Cross-partition store: index shards, membership shards, coordinator.
209    pub remote: R,
210    /// Member packs, read for external bases.
211    pub blobs: B,
212    /// The ticketed pack's own windows.
213    pub windows: W,
214    /// Shard placement.
215    pub shards: Arc<dyn ShardMap>,
216    /// Indexed limits.
217    pub cfg: IndexedConfig,
218    /// Work units per slice.
219    pub limits: SliceLimits,
220    /// Epoch lease timing for the relay seam.
221    pub lease: LeaseParams,
222    /// Backend clock.
223    pub clock: Arc<dyn Clock>,
224    /// Metrics sink.
225    pub metrics: Arc<dyn Metrics>,
226    /// The extraction seam.
227    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    /// One slice per alarm: a slice spends most of an alarm's budget.
249    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
290/// How a phase step ends, other than by success.
291enum Stop {
292    /// A storage failure or spent budget: the fire fails and retries.
293    Store(StoreError),
294    /// Content-intrinsic failure, persisted as `Rejected` (§9.8).
295    Reject(&'static str),
296    /// A terminal result that is not persisted (R-148).
297    Outcome(Outcome),
298    /// Membership may still catch up: run again later.
299    Wait(u64),
300    /// Persist bounded reconstruction progress before the next alarm.
301    Yield(u64),
302    /// The source changed under the job: start over.
303    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/// Recently decoded in-pack objects, kept as delta bases. The newest entry
317/// stays even when it alone passes the cap.
318#[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/// Per-slice memory: nothing in it is authoritative, rows are.
352#[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    /// Preserve the original entry admission allowance. Two windows overlap
392    /// only in `feed`; decoder scratch runs after the acquired window drops.
393    /// Delta resolution releases the idle reader before decoding another frame.
394    /// Cache retention separately reserves R-203's decoder working memory.
395    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    /// The timer's next due time: never the timer's own, so the key moves.
413    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)] // Phase results and their guarded checkpoint commit together.
434    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        // The phases after the decode are cheap: run them back to back while
487        // each one finishes at once and the slice still has calls to spend.
488        let mut chained = 0;
489        while matches!(result, Ok(0))
490            && ran != Phase::Decode
491            && job.phase != ran
492            && job.phase != Phase::Watch
493            // Enter guarded lookup phases from a durable phase checkpoint.
494            && !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            // Nothing this slice counted is persisted: the cursor did not move.
513            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        // A slice that ended by itself is not a killed one.
536        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        // Closure deletions and the satisfying-pack list commit together.
547        // A crash cannot discard the row before recording its dependency.
548        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    /// A finished job waits for its ticket to close; the timer's next look is
556    /// at expiry, or sooner. If `vs` vanished under it (a concurrent ticket's
557    /// expiry) it verifies again rather than answer pending for ever.
558    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    /// Count the slice durably before working, so a slice the runtime kills
597    /// leaves a mark: repeated kills shrink decode/closure work to one, then
598    /// end the ticket with its phase's platform cap, never Rejected.
599    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    /// Write out buffered idempotent rows.
671    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    /// Delete a finished or abandoned job's rows, then `vs` unless the pack
691    /// became a member (kind-2 self-cleaning, R-148). Local rows only, so
692    /// several pages fit one fire.
693    #[allow(clippy::too_many_lines)] // Keep peer lifetime guards and the Gone transition together.
694    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 &current.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    /// One bounded read of the pack, bound to the job's etag.
853    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    /// A packlist is one small window. An advance can consume at most seven
901    /// tickets, so a list naming more than a lookup and those can never pass
902    /// the MKPL rule: it is capped at once, without storing the list.
903    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)] // One loop owns the slice's stop rules.
924    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        // Timer slices also use the current configuration, independently of
932        // the consuming advance's guard.
933        if job.pack_len > self.h.cfg.max_pack_bytes {
934            return Err(Stop::Reject(super::PACK_CAP_MESSAGE));
935        }
936        // A fresh job or source restart must discard provisional rows first.
937        // Delete one page with the guarded checkpoint, then try again.
938        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            // Frames are still needed for the advance's closure check.
962            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                    // Nested base decoding must not retain an idle pack window
1034                    // alongside the acquired source frame and decoder scratch.
1035                    // Save the same post-entry boundary; this slice commits only
1036                    // after `entry` succeeds, as on every other checkpoint.
1037                    #[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                    // One window of progress per slice: the resumed window
1051                    // and the next. Only an entry boundary can be saved.
1052                    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    /// Verify and record one decoded entry.
1114    #[allow(clippy::too_many_lines)] // Depth, budget, rows and checks share one entry's state.
1115    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        // A replayed entry meets its own row; a real duplicate meets an
1187        // earlier offset. Only the first occurrence counts (native staging).
1188        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            // The summary is written after its pages. Decode's durable phase
1232            // transition is the full-pack completion marker; partial pages
1233            // are provisional and replay identically after a lost reply.
1234            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    /// Put `base`'s canonical bytes in the cache: an earlier frame of this
1268    /// pack is re-read from storage, anything else is a repository member.
1269    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    /// Resolve an external base through this repository's members only, and
1333    /// charge it (once per distinct base, however the slices fall).
1334    async fn resolve_external(
1335        &self,
1336        st: &mut SliceState,
1337        job: &mut VerifyJobV1,
1338        base: Hash,
1339    ) -> Result<(), Stop> {
1340        // A base this slice already resolved is in the member cache: no call.
1341        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        // The memory bound of the retained members. The decode budget itself
1361        // is charged per distinct base by `charge_bases`, from persisted rows,
1362        // so it does not depend on what earlier slices retained.
1363        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        // Dependencies and byte charges are staged for the enclosing flush. The LRU owns
1392        // the needed base; release the duplicate member chain before an outer
1393        // in-pack frame can start its decoder.
1394        st.memo = MemberCache::default();
1395        Ok(())
1396    }
1397
1398    /// Charge each newly retained member object once: the row records the
1399    /// entry that first needed it, so a replayed entry charges again and a
1400    /// later entry does not.
1401    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            // Object-keyed zero-size rows carry depth; location-keyed rows
1454            // charge bytes. Equal bytes at different locations count twice.
1455            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    /// Look owed closure children up in this repository's members. A child
1471    /// found in this pack's frames or in a member is settled and its row goes
1472    /// (the member is remembered in `satisfying`); what stays is exactly what
1473    /// no pack this job can see holds, the rows an advance and the final
1474    /// recheck read.
1475    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                    // Reserve the lesser of the lookup's worst case and a full
1512                    // slice. A larger single-id lookup is a terminal cap.
1513                    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    /// Relay this pack's index rows, one page of frames per slice, after the
1563    /// decode reached `Done` (SPEC-PACKFILE §11).
1564    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    /// R-130: `Verified` only once the relay delivered every index row.
1656    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    /// The guarded `Pending` to `Verified` transition, monotone.
1669    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            // Verify owns scan after decoding; checkpoint its inventory position.
1687            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    /// One final look after the lag bound, so an advance can tell a child
1760    /// that is still catching up from one that is truly open.
1761    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}