Skip to main content

mkit_server/indexed/
checkpoint.rs

1//! Job and checkpoint rows of scheduled verification (`vc`, WP-4.8).
2//!
3//! One job per (repository, pack) lives in the consuming ref shard. Every row
4//! of a job sits under one prefix, so cleanup is a single range. Frame, child
5//! and base rows are written idempotently ahead of the guarded job row that
6//! commits a slice. Settled children are deleted with that checkpoint, so a
7//! crash cannot discard the satisfying-pack dependency.
8
9use super::state::VerificationV1;
10use crate::repo::RepoName;
11use crate::store::{
12    codec::CODEC_V1,
13    index::{IndexEntry, IndexValue},
14    keys,
15};
16use crate::{Batch, NamespaceStore, Partition, StoreError, Value};
17use mkit_core::hash::Hash;
18use serde::{Deserialize, Serialize};
19
20/// New progress per slice: one window, plus the current window on resume.
21pub const WINDOW_BYTES: u64 = 16 << 20;
22/// Entries one slice decodes before its next checkpoint, at most.
23pub const DEFAULT_ENTRY_CAP: u32 = 4096;
24
25/// Cursor bytes are hex strings; hashes use strict fixed-size JSON arrays.
26mod hex {
27    pub(super) mod bytes {
28        use serde::{Deserialize, Deserializer, Serialize, Serializer};
29        pub(in super::super) fn serialize<S: Serializer>(
30            bytes: &[u8],
31            s: S,
32        ) -> Result<S::Ok, S::Error> {
33            mkit_core::hash::to_hex_bytes(bytes).serialize(s)
34        }
35        pub(in super::super) fn deserialize<'de, D: Deserializer<'de>>(
36            d: D,
37        ) -> Result<Vec<u8>, D::Error> {
38            let text = String::deserialize(d)?;
39            if text.len() % 2 != 0 || !text.is_ascii() {
40                return Err(serde::de::Error::custom("malformed hex"));
41            }
42            (0..text.len())
43                .step_by(2)
44                .map(|i| u8::from_str_radix(&text[i..i + 2], 16).map_err(serde::de::Error::custom))
45                .collect()
46        }
47    }
48}
49
50/// Where a job is. `Recheck` and `Watch` follow `Verified`: the pack is
51/// usable by an advance from `Recheck` on.
52#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
53#[serde(rename_all = "snake_case")]
54pub enum Phase {
55    /// Window-by-window decode, hashing, signatures and frame rows.
56    #[default]
57    Decode,
58    /// Owed closure children looked up in the repository's members.
59    ClosureResolve,
60    /// Index rows relayed to their shards, only after the pack decoded to `Done`.
61    EmitIndex,
62    /// Wait until every emitted relay row has been delivered (R-130).
63    AwaitDelivery,
64    /// Extraction slot: WP-4.10b fills it; this WP fails closed (R-163).
65    Extract,
66    /// Guarded `Pending` to `Verified` transition of `vs`.
67    Verify,
68    /// Final closure recheck once the membership lag window has passed.
69    Recheck,
70    /// Finished: waits for the ticket to close, then cleans its rows.
71    Watch,
72}
73
74/// What kind of upload a job verifies.
75#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
76#[serde(rename_all = "snake_case")]
77pub enum Kind {
78    /// The first window is not read yet.
79    #[default]
80    Unknown,
81    /// An MKIT pack.
82    Pack,
83    /// An MKPL packlist node.
84    Packlist,
85}
86
87/// A terminal result that is never persisted as `Rejected`: its cause is the
88/// repository's membership or the platform, not the pack (R-148).
89#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
90#[serde(rename_all = "snake_case")]
91pub enum Outcome {
92    /// An external base was not a member after the lag window.
93    BaseMissing,
94    /// A fresh global denial: permission failure, distinct from unavailable storage.
95    Blocked,
96    /// An external base lookup hit an index cap.
97    BaseCapped,
98    /// A closure lookup hit an index cap.
99    ClosureCapped,
100    /// Whole-group closure was still open after the membership lag window.
101    ClosureMissing,
102    /// A named pack was still not a member after the membership lag window.
103    PacklistMissing,
104    /// The claimed head has a type that cannot be a history tip.
105    OpenClosure,
106    /// In-pack plus external chain depth passed the cap.
107    ExternalTooDeep,
108    /// The decode budget, or one object past the Worker's resident cap.
109    DecodeBudget,
110    /// The pack needs extraction, which the Worker cannot do yet (WP-4.10b).
111    ExtractionUnavailable,
112    /// A selected object became blocked before a holder could be queued.
113    ObjectBlocked,
114}
115
116/// One immutable member of the Advance that claimed an extraction group.
117#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
118#[serde(deny_unknown_fields)]
119pub struct ExtractionGroupMember {
120    /// Immutable source pack identity.
121    pub pack: Hash,
122    /// Ticket identity selected by the Advance.
123    pub ticket: Hash,
124    /// Declared pack length.
125    pub bytes: u64,
126    /// Source age for repository-local resolution failures.
127    pub created_at_ms: u64,
128    /// Native still stages this member's selection facts, but does not extract
129    /// objects whose first owner was already verified.
130    pub already_verified: bool,
131}
132
133/// The persisted state of one job, guarded by `vc` sub-class 0.
134#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
135#[serde(deny_unknown_fields)]
136#[allow(clippy::struct_excessive_bools)] // Independent pack facts, lifetime and hydration markers.
137pub struct VerifyJobV1 {
138    /// Monotone across every mutation and retained Gone header.
139    pub generation: u64,
140    /// Cleanup retained this header after deleting its facts and bodies.
141    pub gone: bool,
142    /// Immutable satisfying-member and packlist body in this pack's vc4 range.
143    pub member_body_id: Option<Hash>,
144    /// Only hydrated jobs may replace the member lists.
145    #[serde(skip)]
146    pub members_loaded: bool,
147    /// First consuming Advance's immutable head, checked before effects.
148    pub extraction_head: Option<Hash>,
149    /// Bounded extraction cursors; bulk parts, receipts and offsets stay in vc4.
150    #[serde(default)]
151    pub extraction: Option<ExtractionV1>,
152    /// Ordered, atomically claimed first-Advance extraction context. Standalone
153    /// jobs await a consuming Advance before an extraction group is claimed.
154    #[serde(default)]
155    pub extraction_group: Vec<ExtractionGroupMember>,
156    /// The ticket this job serves; the job lives while that ticket does.
157    pub ticket_id: Hash,
158    /// That ticket's creation time: the start of the membership lag window.
159    pub created_at_ms: u64,
160    /// Declared pack length.
161    pub pack_len: u64,
162    /// Current phase.
163    pub phase: Phase,
164    /// Upload type, known after the first window.
165    pub kind: Kind,
166    /// Pack format version, known after the first window.
167    pub version: u32,
168    /// `WindowCursor` bytes at the last checkpoint; empty before the first.
169    #[serde(with = "hex::bytes")]
170    pub cursor: Vec<u8>,
171    /// Object-store etag every window read is bound to (SPEC-PACKFILE ยง11).
172    pub etag: Option<String>,
173    /// Entries decoded up to `cursor`.
174    pub entries: u64,
175    /// Sum of first-occurrence decoded sizes.
176    pub in_pack_bytes: u64,
177    /// Sum of distinct external base and chain-intermediate sizes.
178    pub external_bytes: u64,
179    /// Windows read so far, for `Retry-After`.
180    pub windows_done: u32,
181    /// Slices started on the current cursor.
182    pub attempts: u32,
183    /// Entries one slice may decode; halved after repeated failures.
184    pub entry_cap: u32,
185    /// Closure ids per slice; repeated interrupted passes shrink this to one.
186    pub closure_cap: u32,
187    /// Times a source change restarted the job.
188    pub restarts: u8,
189    /// A bad signature was seen; reported only if the pack decodes to `Done`.
190    pub bad_signature: bool,
191    /// An entry needs extraction (see [`Outcome::ExtractionUnavailable`]).
192    pub extract_needed: bool,
193    /// Last examined closure id (or emit page); advances with guarded deletions.
194    #[serde(with = "hex::bytes")]
195    pub scan: Vec<u8>,
196    /// Owed closure children seen in the current pass.
197    pub owed: u64,
198    /// Whether the current closure pass is the final recheck.
199    pub final_pass: bool,
200    /// Distinct member packs that satisfied a child.
201    #[serde(skip)]
202    pub satisfying: Vec<Hash>,
203    /// The last relay sequence this job enqueued.
204    pub last_relay_seq: Option<u64>,
205    /// When the final closure recheck finished.
206    pub closure_final_at_ms: Option<u64>,
207    /// The packs a packlist names.
208    #[serde(skip)]
209    pub packlist: Vec<Hash>,
210    /// Verified MKPL predecessor, checkpointed with its decoded header.
211    pub packlist_prev: Option<Hash>,
212    /// A terminal non-persisted result.
213    pub outcome: Option<Outcome>,
214}
215
216/// Resumable driver progress. Every bulk row is keyed by the frozen group digest.
217#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
218#[serde(deny_unknown_fields)]
219pub struct ExtractionV1 {
220    /// Immutable decoded member descriptors captured before effects.
221    pub sources: Vec<ExtractionSource>,
222    /// Domain separated group identity.
223    pub group: Hash,
224    /// Scan, selection, source verification, upload, offsets, enqueue or delivery.
225    pub stage: u8,
226    /// Group scan's member cursor.
227    pub member: usize,
228    /// Frame scan cursor.
229    #[serde(with = "hex::bytes")]
230    pub scan: Vec<u8>,
231    /// Distinct canonical objects in the frozen group.
232    pub staged_objects: u64,
233    /// Union canonical bytes.
234    pub staged_bytes: u64,
235    /// Union selected content bytes.
236    pub selected_bytes: u64,
237    /// Current selected object.
238    pub object: Option<Hash>,
239    /// Declared current object length.
240    pub length: u64,
241    /// Next chunk to resolve.
242    pub chunk: u32,
243    /// Bytes already consumed from the current canonical chunk.
244    pub chunk_offset: u64,
245    /// Verified content bytes stored in bounded local fragments.
246    pub written: u64,
247    /// Parts whose CV was computed.
248    pub cvs: u32,
249    /// Content root computed before any publication.
250    pub root: Option<Hash>,
251    /// Root pinned opaque backend session.
252    #[serde(with = "hex::bytes")]
253    pub session: Vec<u8>,
254    /// Parts committed to the backend.
255    pub uploaded: u32,
256    /// Atomic holder outbox sequence.
257    pub relay: Option<u64>,
258    /// Resumable member reconstruction; ancestry and bytes remain in vc4.
259    pub reconstruction: Option<MemberCursor>,
260}
261
262/// One member resolution; only its immediate canonical parent remains live.
263#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
264#[serde(deny_unknown_fields)]
265pub struct MemberCursor {
266    /// Immutable source being resolved.
267    pub target: Hash,
268    /// Current descent frontier.
269    pub next: Hash,
270    /// Previous pack and offset for same-pack preference.
271    pub preferred: Option<(Hash, u64)>,
272    /// Stack row cursor.
273    pub level: u32,
274    /// Whether the current descent still follows a staged pack.
275    pub local: bool,
276    /// Decode the stored chain toward its requested source.
277    pub ascending: bool,
278    /// Immediate canonical parent (id, length, total depth, source pack).
279    pub canonical: Option<(Hash, u64, u32, Hash)>,
280    /// Native canonical byte cost, including discarded ancestors.
281    pub bytes: u64,
282}
283
284/// Identity of one verified source's selection facts.
285#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
286#[serde(deny_unknown_fields)]
287pub struct ExtractionSource {
288    /// Ordered immutable group member.
289    pub member: ExtractionGroupMember,
290    /// Bounded content addressed `ETag` token.
291    pub etag: Option<String>,
292    /// Canonical source pack version.
293    pub version: u32,
294    /// Completed frame count.
295    pub entries: u64,
296    /// Canonical bytes this source decoded.
297    pub decoded: u64,
298    /// Immutable member lists validated by the frozen closure barrier.
299    pub member_body_id: Option<Hash>,
300}
301
302impl VerifyJobV1 {
303    pub(super) fn closure_retry(&self) -> bool {
304        matches!(
305            self.outcome,
306            Some(Outcome::ClosureMissing | Outcome::PacklistMissing | Outcome::BaseCapped)
307        ) && self
308            .extraction
309            .as_ref()
310            .is_some_and(|x| x.object.is_none() && (x.stage <= 2 || (10..=13).contains(&x.stage)))
311    }
312
313    /// A job for `ticket`, at the start of `Decode`.
314    #[must_use]
315    pub fn new(ticket_id: Hash, created_at_ms: u64, pack_len: u64, entry_cap: u32) -> Self {
316        Self {
317            ticket_id,
318            created_at_ms,
319            pack_len,
320            entry_cap,
321            closure_cap: 4,
322            members_loaded: true,
323            ..Self::default()
324        }
325    }
326
327    /// The job's progress reset for a fresh run: a source changed.
328    pub fn restart(&mut self) {
329        let fresh = Self::new(
330            self.ticket_id,
331            self.created_at_ms,
332            self.pack_len,
333            self.entry_cap,
334        );
335        *self = Self {
336            restarts: self.restarts.saturating_add(1),
337            extraction_group: self.extraction_group.clone(),
338            extraction_head: self.extraction_head,
339            ..fresh
340        };
341    }
342
343    /// Whether an advance may rely on the pack: verified, extracted, indexed.
344    #[must_use]
345    pub fn usable(&self) -> bool {
346        !self.gone && self.outcome.is_none() && matches!(self.phase, Phase::Recheck | Phase::Watch)
347    }
348}
349
350/// Encode a job with the metadata codec version byte.
351///
352/// # Panics
353/// Serializing this fixed DTO into a `Vec` cannot fail.
354#[must_use]
355pub fn encode_job(job: &VerifyJobV1) -> Value {
356    let mut bytes = vec![CODEC_V1];
357    serde_json::to_writer(&mut bytes, job).expect("job DTO serializes");
358    Value::new(bytes)
359}
360
361/// Decode a job, failing closed on corrupt or future values.
362pub fn decode_job(value: &Value) -> Result<VerifyJobV1, StoreError> {
363    let Some((&CODEC_V1, body)) = value.as_bytes().split_first() else {
364        return Err(StoreError::Corrupt("bad verification job version".into()));
365    };
366    let job: VerifyJobV1 = serde_json::from_slice(body)
367        .map_err(|_| StoreError::Corrupt("bad verification job".into()))?;
368    validate_header(&job, value)?;
369    Ok(job)
370}
371
372/// Every guarded header is small even with seven source snapshots. Window
373/// cursors are at most 4KiB; extraction starts only after that cursor clears.
374pub const MAX_JOB_HEADER_BYTES: usize = 16 << 10;
375
376fn validate_header(job: &VerifyJobV1, raw: &Value) -> Result<(), StoreError> {
377    let bounded = raw.as_bytes().len() <= MAX_JOB_HEADER_BYTES
378        && job.cursor.len() <= 4096
379        && job.scan.len() <= 324
380        && job.etag.as_ref().is_none_or(|e| e.len() <= 64)
381        && job.extraction_group.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
382        && job.extraction.as_ref().is_none_or(|x| {
383            job.cursor.is_empty()
384                && x.scan.len() <= 324
385                && x.session.len() <= 1024
386                && x.cvs <= 10_000
387                && x.uploaded <= 10_000
388                && x.reconstruction
389                    .as_ref()
390                    .is_none_or(|r| r.canonical.is_none_or(|(_, n, _, _)| n <= 8 << 20))
391                && x.stage <= 13
392                && x.member <= x.sources.len()
393                && x.sources.len() <= crate::store::outbox::MAX_TICKETS_PER_ADVANCE
394                && x.sources
395                    .iter()
396                    .all(|s| s.etag.as_ref().is_none_or(|e| e.len() <= 64))
397        });
398    if bounded {
399        Ok(())
400    } else {
401        Err(StoreError::Corrupt("oversized job header".into()))
402    }
403}
404
405#[derive(Serialize, Deserialize)]
406#[serde(deny_unknown_fields)]
407struct MemberLists {
408    satisfying: Vec<Hash>,
409    packlist: Vec<Hash>,
410}
411
412fn member_id(bytes: &[u8]) -> Hash {
413    let mut h = mkit_core::hash::Hasher::new();
414    h.update(b"mkit-job-members:v1");
415    h.update(bytes);
416    h.finalize()
417}
418fn member_key(repo: &RepoName, pack: &Hash, id: &Hash) -> crate::Key {
419    keys::verify_row(repo, pack, keys::VC_CANDIDATE, Some(id))
420}
421
422/// Append a generation-bumped header and, only when changed, its immutable
423/// lists. The caller must guard the exact prior header and its deadline.
424pub fn write_job(
425    mut batch: Batch,
426    job: &mut VerifyJobV1,
427    prior: Option<&Value>,
428    repo: &RepoName,
429    pack: &Hash,
430) -> Result<Batch, StoreError> {
431    let old = prior.map(decode_job).transpose()?;
432    job.generation = old
433        .as_ref()
434        .map_or(0, |j| j.generation)
435        .checked_add(1)
436        .ok_or_else(|| StoreError::Corrupt("job generation overflow".into()))?;
437    let old_body = old.as_ref().and_then(|j| j.member_body_id);
438    if job.members_loaded {
439        if job.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
440            || job.packlist.len()
441                > crate::store::index::MAX_LOOKUP_IDS
442                    + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
443        {
444            return Err(StoreError::Corrupt("oversized job member lists".into()));
445        }
446        job.member_body_id = if job.satisfying.is_empty() && job.packlist.is_empty() {
447            None
448        } else {
449            let mut bytes = vec![CODEC_V1];
450            serde_json::to_writer(
451                &mut bytes,
452                &MemberLists {
453                    satisfying: job.satisfying.clone(),
454                    packlist: job.packlist.clone(),
455                },
456            )
457            .map_err(StoreError::unavailable)?;
458            let id = member_id(&bytes);
459            if Some(id) != old_body {
460                batch = batch.put(member_key(repo, pack, &id), Value::new(bytes));
461            }
462            Some(id)
463        };
464    } else if job.member_body_id != old_body {
465        return Err(StoreError::Corrupt(
466            "unloaded job member lists changed".into(),
467        ));
468    }
469    let header = encode_job(job);
470    validate_header(job, &header)?;
471    Ok(batch.put(keys::verify_job(repo, pack), header))
472}
473
474/// Hydrate a header without expanding any peer's guarded value. Bodies are
475/// content addressed and remain immutable until the job's guarded cleanup.
476pub async fn hydrate_job<S: NamespaceStore>(
477    store: &S,
478    source: &Partition,
479    repo: &RepoName,
480    pack: &Hash,
481    job: &mut VerifyJobV1,
482) -> Result<(), StoreError> {
483    if !job.gone
484        && let Some(id) = job.member_body_id
485    {
486        let raw = store
487            .get(source, &member_key(repo, pack, &id))
488            .await?
489            .ok_or_else(|| StoreError::Unavailable("job member body disappeared".into()))?;
490        let Some((&CODEC_V1, bytes)) = raw.as_bytes().split_first() else {
491            return Err(StoreError::Corrupt("bad job member body version".into()));
492        };
493        if member_id(raw.as_bytes()) != id {
494            return Err(StoreError::Corrupt("bad job member body digest".into()));
495        }
496        let lists: MemberLists = serde_json::from_slice(bytes)
497            .map_err(|_| StoreError::Corrupt("bad job member body".into()))?;
498        if lists.satisfying.len() > crate::store::index::MAX_LOOKUP_IDS
499            || lists.packlist.len()
500                > crate::store::index::MAX_LOOKUP_IDS
501                    + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
502        {
503            return Err(StoreError::Corrupt("oversized job member body".into()));
504        }
505        job.satisfying = lists.satisfying;
506        job.packlist = lists.packlist;
507    }
508    job.members_loaded = true;
509    Ok(())
510}
511
512/// A pack entry as the job recorded it. First occurrence wins.
513#[derive(Debug, Clone, Copy, PartialEq, Eq)]
514pub struct FrameRow {
515    /// Location, sizes and in-pack depth, as the index row will carry them.
516    pub value: IndexValue,
517    /// Object type tag, for the head check.
518    pub object_type: u8,
519    /// The external base this entry's delta chain ends at, if any.
520    pub external: Option<Hash>,
521}
522
523/// Encode a frame row: type, optional external base, then the index value.
524pub fn encode_frame(id: &Hash, row: &FrameRow) -> Result<Value, StoreError> {
525    let index = crate::store::codec::encode_object_index(id, &row.value)?;
526    let mut bytes = vec![row.object_type, u8::from(row.external.is_some())];
527    if let Some(base) = &row.external {
528        bytes.extend_from_slice(base);
529    }
530    bytes.extend_from_slice(index.as_bytes());
531    Ok(Value::new(bytes))
532}
533
534/// Decode a frame row for object `id`.
535pub fn decode_frame(id: &Hash, value: &Value) -> Result<FrameRow, StoreError> {
536    let corrupt = || StoreError::Corrupt("bad verification frame row".into());
537    let bytes = value.as_bytes();
538    let (&object_type, rest) = bytes.split_first().ok_or_else(corrupt)?;
539    let (external, rest) = match rest.split_first() {
540        Some((0, rest)) => (None, rest),
541        Some((1, rest)) => {
542            let (base, rest) = rest.split_first_chunk::<32>().ok_or_else(corrupt)?;
543            (Some(*base), rest)
544        }
545        _ => return Err(corrupt()),
546    };
547    Ok(FrameRow {
548        value: crate::store::codec::decode_object_index(id, &Value::new(rest.to_vec()))?,
549        object_type,
550        external,
551    })
552}
553
554/// An external base row: location-keyed rows charge size once using the first
555/// entry marker; object-keyed zero-size rows retain the base's chain depth.
556#[derive(Debug, Clone, Copy, PartialEq, Eq)]
557pub struct BaseRow {
558    /// Decoded size of the base object.
559    pub size: u64,
560    /// Total delta depth of the base in its member pack.
561    pub depth: u32,
562    /// Index of the first entry that needed it.
563    pub entry: u64,
564}
565
566/// Encode a base row.
567#[must_use]
568pub fn encode_base(row: &BaseRow) -> Value {
569    let mut bytes = row.size.to_be_bytes().to_vec();
570    bytes.extend_from_slice(&row.depth.to_be_bytes());
571    bytes.extend_from_slice(&row.entry.to_be_bytes());
572    Value::new(bytes)
573}
574
575/// Decode a base row.
576pub fn decode_base(value: &Value) -> Result<BaseRow, StoreError> {
577    let bytes: &[u8; 20] = value
578        .as_bytes()
579        .try_into()
580        .map_err(|_| StoreError::Corrupt("bad verification base row".into()))?;
581    Ok(BaseRow {
582        size: u64::from_be_bytes(bytes[..8].try_into().unwrap_or_default()),
583        depth: u32::from_be_bytes(bytes[8..12].try_into().unwrap_or_default()),
584        entry: u64::from_be_bytes(bytes[12..].try_into().unwrap_or_default()),
585    })
586}
587
588/// The index entry a frame row stands for.
589#[must_use]
590pub fn index_entry(id: Hash, row: &FrameRow) -> IndexEntry {
591    IndexEntry {
592        object: id,
593        value: row.value,
594    }
595}
596
597/// The reference of a job's timer: `repo 00 pack`.
598#[must_use]
599pub fn timer_reference(repo: &RepoName, pack: &Hash) -> Vec<u8> {
600    let mut reference = repo.as_str().as_bytes().to_vec();
601    reference.push(0);
602    reference.extend_from_slice(pack);
603    reference
604}
605
606/// Split a timer reference back into its repository and pack.
607#[must_use]
608pub fn parse_reference(reference: &[u8]) -> Option<(RepoName, Hash)> {
609    let sep = reference.iter().position(|&b| b == 0)?;
610    let pack: Hash = reference[sep + 1..].try_into().ok()?;
611    let repo = RepoName::new(String::from_utf8(reference[..sep].to_vec()).ok()?).ok()?;
612    Some((repo, pack))
613}
614
615type Stored<T> = Option<(T, Value)>;
616
617/// The job row and `vs` in one read.
618pub async fn read_job<S: NamespaceStore>(
619    store: &S,
620    source: &Partition,
621    repo: &RepoName,
622    pack: &Hash,
623) -> Result<(Stored<VerifyJobV1>, Stored<VerificationV1>), StoreError> {
624    let rows = store
625        .get_many(
626            source,
627            &[keys::verify_job(repo, pack), keys::verification(repo, pack)],
628        )
629        .await?;
630    let [job, state] = <[_; 2]>::try_from(rows)
631        .map_err(|_| StoreError::Corrupt("short verification read".into()))?;
632    let job = if let Some(raw) = job {
633        let mut job = decode_job(&raw)?;
634        hydrate_job(store, source, repo, pack, &mut job).await?;
635        Some((job, raw))
636    } else {
637        None
638    };
639    Ok((
640        job,
641        state
642            .map(|raw| super::state::decode(&raw).map(|state| (state, raw)))
643            .transpose()?,
644    ))
645}
646
647#[cfg(test)]
648mod tests {
649    use super::*;
650
651    #[test]
652    fn job_frame_and_base_codecs_round_trip() {
653        let mut job = VerifyJobV1::new([1; 32], 5, 99, 4096);
654        job.cursor = vec![0xab, 0x01];
655        job.satisfying = vec![[2; 32]];
656        job.outcome = Some(Outcome::BaseMissing);
657        let header = decode_job(&encode_job(&job)).unwrap();
658        assert!(header.satisfying.is_empty());
659        assert!(!header.members_loaded);
660        assert_eq!(header.outcome, job.outcome);
661        assert!(decode_job(&Value::new(b"\x02{}".to_vec())).is_err());
662        let frame = FrameRow {
663            value: IndexValue {
664                frame_offset: 12,
665                frame_length: 40,
666                wire_type: 0x02,
667                decoded_size: 7,
668                chain_depth: 2,
669                delta_base: Some([3; 32]),
670            },
671            object_type: 3,
672            external: Some([4; 32]),
673        };
674        let id = [9; 32];
675        assert_eq!(
676            decode_frame(&id, &encode_frame(&id, &frame).unwrap()).unwrap(),
677            frame
678        );
679        let base = BaseRow {
680            size: 1,
681            depth: 2,
682            entry: 3,
683        };
684        assert_eq!(decode_base(&encode_base(&base)).unwrap(), base);
685        let name = RepoName::new("a").unwrap();
686        assert_eq!(
687            parse_reference(&timer_reference(&name, &[7; 32])),
688            Some((name, [7; 32]))
689        );
690        job.restart();
691        assert_eq!(
692            (job.restarts, job.phase, job.cursor.len()),
693            (1, Phase::Decode, 0)
694        );
695    }
696}