Skip to main content

icydb_core/db/integrity/
progress_store.rs

1//! Module: db::integrity::progress_store
2//! Responsibility: independently persist one bounded record per Deep job.
3//! Does not own: inspected database state, commit markers, journals, or advancement semantics.
4//! Boundary: current-form job codec -> physically separate stable BTreeMap allocation.
5
6use crate::{
7    db::{
8        codec::{finalize_hash_sha256, new_hash_sha256_prefixed},
9        database_format::crc32c,
10        integrity::{
11            IntegrityJob, IntegrityJobError, IntegrityJobId, IntegrityJobOwner, IntegrityJobState,
12            progress_codec::{
13                MAX_INTEGRITY_JOB_PAYLOAD_BYTES, decode_integrity_job_payload,
14                encode_integrity_job_payload,
15            },
16        },
17        mutation_job::{
18            MAX_MUTATION_JOB_RECORD_BYTES, MutationJobError, MutationJobId, MutationJobRecord,
19            MutationJobStatus, decode_mutation_job_payload, encode_mutation_job_payload,
20        },
21        resumable_job::{
22            ResumableJobError, ResumableJobId, ResumableJobRecord, ResumableJobStatus,
23            decode_resumable_job_payload, encode_resumable_job_payload,
24        },
25    },
26    error::InternalError,
27    traits::CanisterKind,
28};
29use candid::CandidType;
30use ic_memory::RuntimeMemory;
31use ic_memory::ic_stable_structures::{
32    BTreeMap as StableBTreeMap, DefaultMemoryImpl, Storable, storable::Bound,
33};
34#[cfg(not(test))]
35use ic_memory::open_default_memory_manager_memory_by_key;
36use serde::Deserialize;
37use sha2::Digest;
38use std::borrow::Cow;
39#[cfg(test)]
40use std::cell::RefCell;
41use std::ops::Bound::{Excluded, Unbounded};
42
43const PROGRESS_HEADER_KEY: ProgressRecordKey = ProgressRecordKey([0; 32]);
44const PROGRESS_HEADER_MAGIC: &[u8; 8] = b"ICYIPROG";
45const PROGRESS_HEADER_VERSION: u8 = 1;
46const PROGRESS_HEADER_BYTES: usize = 8 + 1 + 4;
47const JOB_RECORD_MAGIC: &[u8; 8] = b"ICYIJPTH";
48const JOB_RECORD_VERSION: u8 = 1;
49const JOB_RECORD_HEADER_BYTES: usize = 8 + 1 + 4 + 4;
50const RESUMABLE_JOB_KEY_DOMAIN: &[u8] = b"icydb.resumable-job.progress-key.v1";
51const RESUMABLE_JOB_RECORD_MAGIC: &[u8; 8] = b"ICYRJOB1";
52const RESUMABLE_JOB_RECORD_VERSION: u8 = 1;
53const RESUMABLE_JOB_RECORD_HEADER_BYTES: usize = 8 + 1 + 4 + 4;
54const MUTATION_JOB_KEY_DOMAIN: &[u8] = b"icydb.mutation-job.progress-key.v1";
55const MUTATION_JOB_RECORD_MAGIC: &[u8; 8] = b"ICYMJOB1";
56const MUTATION_JOB_RECORD_VERSION: u8 = 1;
57const MUTATION_JOB_RECORD_HEADER_BYTES: usize = 8 + 1 + 4 + 4;
58const MUTATION_PROGRESS_BEFORE_DIGEST_DOMAIN: &[u8] = b"icydb.mutation-job.progress-before.v1";
59const MAX_PROGRESS_RECORD_BYTES: u32 = 512 * 1024;
60const MAX_PROGRESS_JOBS_GLOBAL: u64 = 64;
61const MAX_PROGRESS_JOBS_NON_INTEGRITY: u64 = 56;
62const PROGRESS_JOBS_INTEGRITY_RESERVATION: u64 =
63    MAX_PROGRESS_JOBS_GLOBAL - MAX_PROGRESS_JOBS_NON_INTEGRITY;
64const MAX_PROGRESS_JOBS_PER_OWNER: u64 = 8;
65
66/// Stable family of one retained record in the shared progress allocation.
67#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
68pub enum ProgressJobFamily {
69    /// Engine-owned Deep integrity inspection.
70    Integrity,
71    /// Application-owned generic resumable job.
72    Resumable,
73    /// Engine-owned fixed SQL mutation job.
74    Mutation,
75}
76
77/// Bounded normalized lifecycle of one retained progress record.
78#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
79pub enum ProgressJobLifecycle {
80    /// More work may advance.
81    Active,
82    /// Integrity advancement is frozen before its terminal receipt is published.
83    TerminalPending,
84    /// Successful exhaustion is retained.
85    Completed,
86    /// Protected source authority changed.
87    Invalidated,
88    /// Mutation authority or policy requires a fresh job.
89    RestartRequired,
90    /// Another integrity terminal outcome is retained.
91    Terminal,
92}
93
94/// One privacy-bounded retained-record inventory entry.
95#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
96pub struct ProgressJobInventoryRecord {
97    /// Durable progress family.
98    pub family: ProgressJobFamily,
99    /// Family-owned job identity bytes.
100    pub job_id: [u8; 32],
101    /// Normalized retained lifecycle.
102    pub lifecycle: ProgressJobLifecycle,
103    /// Family sequence when its current record carries one.
104    pub sequence: Option<u64>,
105}
106
107/// Complete bounded inventory of the shared excluded progress allocation.
108#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
109pub struct ProgressJobInventory {
110    /// Total retained records across all families.
111    pub retained_count: u64,
112    /// Hard retained-record capacity.
113    pub hard_limit: u64,
114    /// Slots protected for integrity work from non-integrity starts.
115    pub reserved_integrity_headroom: u64,
116    /// Retained Deep integrity records.
117    pub integrity_count: u64,
118    /// Retained generic resumable records.
119    pub resumable_count: u64,
120    /// Retained fixed SQL mutation records.
121    pub mutation_count: u64,
122    /// Complete encoded bytes across every retained record.
123    pub retained_record_bytes: u64,
124    /// Every validated retained record in stable key order.
125    pub records: Vec<ProgressJobInventoryRecord>,
126}
127
128#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
129struct ProgressRecordKey([u8; 32]);
130
131impl ProgressRecordKey {
132    const fn from_job_id(job_id: IntegrityJobId) -> Self {
133        Self(job_id.to_bytes())
134    }
135
136    fn from_resumable_job_id(job_id: ResumableJobId) -> Result<Self, ResumableJobError> {
137        let mut hasher = new_hash_sha256_prefixed(RESUMABLE_JOB_KEY_DOMAIN);
138        hasher.update(job_id.to_bytes());
139        let key = finalize_hash_sha256(hasher);
140        if key == PROGRESS_HEADER_KEY.0 {
141            return Err(ResumableJobError::InvalidJobId);
142        }
143        Ok(Self(key))
144    }
145
146    fn from_mutation_job_id(job_id: MutationJobId) -> Result<Self, MutationJobError> {
147        job_id.validate()?;
148        let mut hasher = new_hash_sha256_prefixed(MUTATION_JOB_KEY_DOMAIN);
149        hasher.update(job_id.to_bytes());
150        let key = finalize_hash_sha256(hasher);
151        if key == PROGRESS_HEADER_KEY.0 {
152            return Err(MutationJobError::InvalidJobId);
153        }
154        Ok(Self(key))
155    }
156
157    const fn to_bytes(self) -> [u8; 32] {
158        self.0
159    }
160}
161
162impl Storable for ProgressRecordKey {
163    fn to_bytes(&self) -> Cow<'_, [u8]> {
164        Cow::Borrowed(&self.0)
165    }
166
167    fn from_bytes(bytes: Cow<'_, [u8]>) -> Self {
168        let mut key = [0; 32];
169        if bytes.len() == key.len() {
170            key.copy_from_slice(bytes.as_ref());
171        }
172        Self(key)
173    }
174
175    fn into_bytes(self) -> Vec<u8> {
176        self.0.to_vec()
177    }
178
179    const BOUND: Bound = Bound::Bounded {
180        max_size: 32,
181        is_fixed_size: true,
182    };
183}
184
185#[derive(Clone, Debug, Eq, PartialEq)]
186struct ProgressRecordBytes(Vec<u8>);
187
188impl Storable for ProgressRecordBytes {
189    fn to_bytes(&self) -> Cow<'_, [u8]> {
190        Cow::Borrowed(self.0.as_slice())
191    }
192
193    fn from_bytes(bytes: Cow<'_, [u8]>) -> Self {
194        Self(bytes.into_owned())
195    }
196
197    fn into_bytes(self) -> Vec<u8> {
198        self.0
199    }
200
201    const BOUND: Bound = Bound::Bounded {
202        max_size: MAX_PROGRESS_RECORD_BYTES,
203        is_fixed_size: false,
204    };
205}
206
207pub(super) enum InsertJobResult {
208    Inserted,
209    Occupied(Box<IntegrityJob>),
210}
211
212pub(in crate::db) enum InsertMutationJobResult {
213    Inserted,
214    Occupied(Box<MutationJobRecord>),
215}
216
217/// Exact mutation-progress compare-and-replace effect carried by one commit marker.
218///
219/// Every constructor validates the complete before/after proof. Store
220/// preflight and apply therefore compare only the already-canonical bytes and
221/// do not repeat bounded record decoding on the commit hot path.
222#[derive(Clone, Debug)]
223pub(in crate::db) struct MutationProgressRecordOp {
224    key: ProgressRecordKey,
225    job_id: MutationJobId,
226    expected_sequence: u64,
227    expected_before_digest: [u8; 32],
228    before: Vec<u8>,
229    after: Vec<u8>,
230}
231
232impl MutationProgressRecordOp {
233    /// Build one exact successor replacement from two validated current records.
234    pub(in crate::db) fn replace(
235        before: &MutationJobRecord,
236        after: &MutationJobRecord,
237    ) -> Result<Self, MutationJobError> {
238        let job_id = before.state().job_id;
239        let expected_sequence = before.state().sequence;
240        let key = ProgressRecordKey::from_mutation_job_id(job_id)?;
241        let before = encode_mutation_job_record(before)?;
242        Self::from_encoded(
243            key.to_bytes(),
244            job_id,
245            expected_sequence,
246            mutation_progress_before_digest(&before),
247            before,
248            encode_mutation_job_record(after)?,
249        )
250    }
251
252    /// Reconstruct and validate one current marker-owned replacement.
253    pub(in crate::db) fn from_encoded(
254        key: [u8; 32],
255        job_id: MutationJobId,
256        expected_sequence: u64,
257        expected_before_digest: [u8; 32],
258        before: Vec<u8>,
259        after: Vec<u8>,
260    ) -> Result<Self, MutationJobError> {
261        let operation = Self {
262            key: ProgressRecordKey(key),
263            job_id,
264            expected_sequence,
265            expected_before_digest,
266            before,
267            after,
268        };
269        operation.validate()?;
270        Ok(operation)
271    }
272
273    pub(in crate::db) const fn key(&self) -> [u8; 32] {
274        self.key.to_bytes()
275    }
276
277    pub(in crate::db) const fn job_id(&self) -> MutationJobId {
278        self.job_id
279    }
280
281    pub(in crate::db) const fn expected_sequence(&self) -> u64 {
282        self.expected_sequence
283    }
284
285    pub(in crate::db) const fn expected_before_digest(&self) -> [u8; 32] {
286        self.expected_before_digest
287    }
288
289    pub(in crate::db) const fn before_bytes(&self) -> &[u8] {
290        self.before.as_slice()
291    }
292
293    pub(in crate::db) const fn after_bytes(&self) -> &[u8] {
294        self.after.as_slice()
295    }
296
297    pub(in crate::db) fn validate(&self) -> Result<(), MutationJobError> {
298        self.job_id.validate()?;
299        if ProgressRecordKey::from_mutation_job_id(self.job_id)? != self.key
300            || mutation_progress_before_digest(&self.before) != self.expected_before_digest
301        {
302            return Err(MutationJobError::CorruptProgressStore);
303        }
304        let before = decode_mutation_job_record(&self.before, self.job_id)?;
305        let after = decode_mutation_job_record(&self.after, self.job_id)?;
306        let Some(expected_after_sequence) = self.expected_sequence.checked_add(1) else {
307            return Err(MutationJobError::CounterOverflow);
308        };
309        if before.state().sequence != self.expected_sequence
310            || after.state().sequence != expected_after_sequence
311            || before.canonical_intent() != after.canonical_intent()
312        {
313            return Err(MutationJobError::CorruptProgressStore);
314        }
315        Ok(())
316    }
317}
318
319pub(super) struct ProgressScanPage {
320    pub(super) job_ids: Vec<IntegrityJobId>,
321    pub(super) exhausted: bool,
322}
323
324pub(in crate::db) struct InspectionProgressStore {
325    map: StableBTreeMap<ProgressRecordKey, ProgressRecordBytes, RuntimeMemory<DefaultMemoryImpl>>,
326}
327
328impl InspectionProgressStore {
329    fn open(memory: RuntimeMemory<DefaultMemoryImpl>) -> Result<Self, IntegrityJobError> {
330        let mut store = Self {
331            map: StableBTreeMap::init(memory),
332        };
333        if store.map.is_empty() {
334            store.map.insert(
335                PROGRESS_HEADER_KEY,
336                ProgressRecordBytes(encode_progress_header()),
337            );
338        } else {
339            let header = store
340                .map
341                .get(&PROGRESS_HEADER_KEY)
342                .ok_or(IntegrityJobError::CorruptProgressHeader)?;
343            decode_progress_header(&header.0)?;
344            if store.job_count()? > MAX_PROGRESS_JOBS_GLOBAL {
345                return Err(IntegrityJobError::CorruptProgressHeader);
346            }
347        }
348        Ok(store)
349    }
350
351    pub(super) fn load(&self, job_id: IntegrityJobId) -> Result<IntegrityJob, IntegrityJobError> {
352        let raw = self
353            .map
354            .get(&ProgressRecordKey::from_job_id(job_id))
355            .ok_or(IntegrityJobError::JobNotFound)?;
356        decode_job_record(&raw.0, job_id)
357    }
358
359    pub(super) fn insert_new(
360        &mut self,
361        job: &IntegrityJob,
362    ) -> Result<InsertJobResult, IntegrityJobError> {
363        job.validate()?;
364        let key = ProgressRecordKey::from_job_id(job.id);
365        if key == PROGRESS_HEADER_KEY {
366            return Err(IntegrityJobError::CorruptProgressRecord);
367        }
368        if let Some(raw) = self.map.get(&key) {
369            return decode_job_record(&raw.0, job.id)
370                .map(Box::new)
371                .map(InsertJobResult::Occupied);
372        }
373        if self.job_count()? >= MAX_PROGRESS_JOBS_GLOBAL
374            || self.owner_job_count(&job.owner)? >= MAX_PROGRESS_JOBS_PER_OWNER
375        {
376            return Err(IntegrityJobError::CapacityExceeded);
377        }
378        self.map
379            .insert(key, ProgressRecordBytes(encode_job_record(job)?));
380        Ok(InsertJobResult::Inserted)
381    }
382
383    pub(super) fn replace(&mut self, job: &IntegrityJob) -> Result<(), IntegrityJobError> {
384        job.validate()?;
385        let key = ProgressRecordKey::from_job_id(job.id);
386        if !self.map.contains_key(&key) {
387            return Err(IntegrityJobError::JobNotFound);
388        }
389        self.map
390            .insert(key, ProgressRecordBytes(encode_job_record(job)?));
391        Ok(())
392    }
393
394    pub(in crate::db) fn load_resumable(
395        &self,
396        job_id: ResumableJobId,
397    ) -> Result<ResumableJobRecord, ResumableJobError> {
398        let key = ProgressRecordKey::from_resumable_job_id(job_id)?;
399        let raw = self.map.get(&key).ok_or(ResumableJobError::NotFound)?;
400        decode_resumable_job_record(&raw.0, job_id)
401    }
402
403    pub(in crate::db) fn insert_resumable(
404        &mut self,
405        record: &ResumableJobRecord,
406    ) -> Result<(), ResumableJobError> {
407        record.validate()?;
408        let key = ProgressRecordKey::from_resumable_job_id(record.state().job_id)?;
409        if self.map.contains_key(&key) {
410            return Err(ResumableJobError::AlreadyExists);
411        }
412        if self.job_count().map_err(map_integrity_store_error)? >= MAX_PROGRESS_JOBS_NON_INTEGRITY {
413            return Err(ResumableJobError::CapacityExceeded);
414        }
415        self.map.insert(
416            key,
417            ProgressRecordBytes(encode_resumable_job_record(record)?),
418        );
419        Ok(())
420    }
421
422    pub(in crate::db) fn replace_resumable(
423        &mut self,
424        record: &ResumableJobRecord,
425    ) -> Result<(), ResumableJobError> {
426        record.validate()?;
427        let key = ProgressRecordKey::from_resumable_job_id(record.state().job_id)?;
428        if !self.map.contains_key(&key) {
429            return Err(ResumableJobError::NotFound);
430        }
431        self.map.insert(
432            key,
433            ProgressRecordBytes(encode_resumable_job_record(record)?),
434        );
435        Ok(())
436    }
437
438    pub(in crate::db) fn remove_resumable(
439        &mut self,
440        job_id: ResumableJobId,
441    ) -> Result<(), ResumableJobError> {
442        let key = ProgressRecordKey::from_resumable_job_id(job_id)?;
443        let Some(raw) = self.map.get(&key) else {
444            return Ok(());
445        };
446        decode_resumable_job_record(&raw.0, job_id)?;
447        let _ = self.map.remove(&key);
448        Ok(())
449    }
450
451    pub(in crate::db) fn load_mutation(
452        &self,
453        job_id: MutationJobId,
454    ) -> Result<MutationJobRecord, MutationJobError> {
455        let key = ProgressRecordKey::from_mutation_job_id(job_id)?;
456        let raw = self.map.get(&key).ok_or(MutationJobError::NotFound)?;
457        decode_mutation_job_record(&raw.0, job_id)
458    }
459
460    pub(in crate::db) fn insert_mutation(
461        &mut self,
462        record: &MutationJobRecord,
463    ) -> Result<InsertMutationJobResult, MutationJobError> {
464        record.validate()?;
465        let key = ProgressRecordKey::from_mutation_job_id(record.state().job_id)?;
466        if let Some(raw) = self.map.get(&key) {
467            return decode_mutation_job_record(&raw.0, record.state().job_id)
468                .map(Box::new)
469                .map(InsertMutationJobResult::Occupied);
470        }
471        if self.job_count().map_err(map_mutation_store_error)? >= MAX_PROGRESS_JOBS_NON_INTEGRITY {
472            return Err(MutationJobError::CapacityExceeded);
473        }
474        self.map.insert(
475            key,
476            ProgressRecordBytes(encode_mutation_job_record(record)?),
477        );
478        Ok(InsertMutationJobResult::Inserted)
479    }
480
481    #[cfg_attr(
482        not(test),
483        expect(
484            dead_code,
485            reason = "mutation-job advancement owns sequence-checked replacement"
486        )
487    )]
488    pub(in crate::db) fn replace_mutation(
489        &mut self,
490        record: &MutationJobRecord,
491    ) -> Result<(), MutationJobError> {
492        record.validate()?;
493        let key = ProgressRecordKey::from_mutation_job_id(record.state().job_id)?;
494        if !self.map.contains_key(&key) {
495            return Err(MutationJobError::NotFound);
496        }
497        self.map.insert(
498            key,
499            ProgressRecordBytes(encode_mutation_job_record(record)?),
500        );
501        Ok(())
502    }
503
504    fn preflight_mutation_progress(
505        &self,
506        operation: &MutationProgressRecordOp,
507    ) -> Result<(), MutationJobError> {
508        let current = self
509            .map
510            .get(&operation.key)
511            .ok_or(MutationJobError::CorruptProgressStore)?;
512        if current.0 != operation.before {
513            return Err(MutationJobError::CorruptProgressStore);
514        }
515        Ok(())
516    }
517
518    fn apply_mutation_progress(
519        &mut self,
520        operation: &MutationProgressRecordOp,
521    ) -> Result<(), MutationJobError> {
522        let current = self
523            .map
524            .get(&operation.key)
525            .ok_or(MutationJobError::CorruptProgressStore)?;
526        if current.0 == operation.after {
527            return Ok(());
528        }
529        if current.0 != operation.before {
530            return Err(MutationJobError::CorruptProgressStore);
531        }
532        self.map
533            .insert(operation.key, ProgressRecordBytes(operation.after.clone()));
534        Ok(())
535    }
536
537    fn apply_preflighted_mutation_progress(&mut self, operation: &MutationProgressRecordOp) {
538        self.map
539            .insert(operation.key, ProgressRecordBytes(operation.after.clone()));
540    }
541
542    /// Replace one mutation record without target-row work using the same exact
543    /// before/after proof as marker-owned mutation progress.
544    pub(in crate::db) fn replace_mutation_progress(
545        &mut self,
546        operation: &MutationProgressRecordOp,
547    ) -> Result<(), MutationJobError> {
548        self.apply_mutation_progress(operation)
549    }
550
551    fn verify_mutation_progress(
552        &self,
553        operation: &MutationProgressRecordOp,
554    ) -> Result<(), MutationJobError> {
555        operation.validate()?;
556        let current = self
557            .map
558            .get(&operation.key)
559            .ok_or(MutationJobError::CorruptProgressStore)?;
560        if current.0 != operation.after {
561            return Err(MutationJobError::CorruptProgressStore);
562        }
563        Ok(())
564    }
565
566    pub(in crate::db) fn acknowledge_mutation(
567        &mut self,
568        job_id: MutationJobId,
569        expected_sequence: u64,
570    ) -> Result<(), MutationJobError> {
571        let key = ProgressRecordKey::from_mutation_job_id(job_id)?;
572        let Some(raw) = self.map.get(&key) else {
573            return Ok(());
574        };
575        let record = decode_mutation_job_record(&raw.0, job_id)?;
576        if record.state().sequence != expected_sequence {
577            return Err(MutationJobError::StaleSequence {
578                expected: expected_sequence,
579                actual: record.state().sequence,
580            });
581        }
582        if record.state().status == crate::db::MutationJobStatus::Active {
583            return Err(MutationJobError::Active);
584        }
585        let _ = self.map.remove(&key);
586        Ok(())
587    }
588
589    /// Remove exactly one current initial mutation record.
590    ///
591    /// The caller supplies the engine-format validator so the shared store
592    /// remains independent of SQL continuation semantics while still making
593    /// validation and removal one indivisible store borrow.
594    #[cfg(feature = "sql")]
595    pub(in crate::db) fn cancel_unadvanced_mutation(
596        &mut self,
597        job_id: MutationJobId,
598        expected_sequence: u64,
599        validate_initial_continuation: impl FnOnce(&[u8]) -> Result<(), MutationJobError>,
600    ) -> Result<(), MutationJobError> {
601        let key = ProgressRecordKey::from_mutation_job_id(job_id)?;
602        let Some(raw) = self.map.get(&key) else {
603            return Ok(());
604        };
605        let record = decode_mutation_job_record(&raw.0, job_id)?;
606        let continuation = record.ensure_cancelable_at_sequence(expected_sequence)?;
607        validate_initial_continuation(continuation)?;
608        let _ = self.map.remove(&key);
609        Ok(())
610    }
611
612    /// Decode every retained slot before returning one complete bounded inventory.
613    pub(in crate::db) fn inventory(&self) -> Result<ProgressJobInventory, MutationJobError> {
614        let retained_count = self.job_count().map_err(map_mutation_store_error)?;
615        let record_capacity =
616            usize::try_from(retained_count).map_err(|_| MutationJobError::CorruptProgressStore)?;
617        let mut records = Vec::with_capacity(record_capacity);
618        let mut integrity_count = 0_u64;
619        let mut resumable_count = 0_u64;
620        let mut mutation_count = 0_u64;
621        let mut retained_record_bytes = 0_u64;
622
623        for entry in self.map.iter() {
624            let key = *entry.key();
625            if key == PROGRESS_HEADER_KEY {
626                continue;
627            }
628            let bytes = &entry.value().0;
629            retained_record_bytes = retained_record_bytes
630                .checked_add(
631                    u64::try_from(bytes.len())
632                        .map_err(|_| MutationJobError::CorruptProgressStore)?,
633                )
634                .ok_or(MutationJobError::CorruptProgressStore)?;
635            let record = if bytes.starts_with(JOB_RECORD_MAGIC) {
636                let job_id = IntegrityJobId::try_from_bytes(key.to_bytes())
637                    .map_err(|_| MutationJobError::CorruptProgressStore)?;
638                let job = decode_job_record(bytes, job_id)
639                    .map_err(|_| MutationJobError::CorruptProgressStore)?;
640                integrity_count = integrity_count
641                    .checked_add(1)
642                    .ok_or(MutationJobError::CorruptProgressStore)?;
643                ProgressJobInventoryRecord {
644                    family: ProgressJobFamily::Integrity,
645                    job_id: job.id.to_bytes(),
646                    lifecycle: match job.state {
647                        IntegrityJobState::InProgress => ProgressJobLifecycle::Active,
648                        IntegrityJobState::TerminalPending(_) => {
649                            ProgressJobLifecycle::TerminalPending
650                        }
651                        IntegrityJobState::Terminal { .. } => ProgressJobLifecycle::Terminal,
652                    },
653                    sequence: Some(job.pages_completed),
654                }
655            } else if bytes.starts_with(RESUMABLE_JOB_RECORD_MAGIC) {
656                let job = decode_resumable_job_record_for_inventory(bytes, key)?;
657                resumable_count = resumable_count
658                    .checked_add(1)
659                    .ok_or(MutationJobError::CorruptProgressStore)?;
660                ProgressJobInventoryRecord {
661                    family: ProgressJobFamily::Resumable,
662                    job_id: job.state().job_id.to_bytes(),
663                    lifecycle: match job.state().status {
664                        ResumableJobStatus::Active => ProgressJobLifecycle::Active,
665                        ResumableJobStatus::Completed => ProgressJobLifecycle::Completed,
666                        ResumableJobStatus::Invalidated => ProgressJobLifecycle::Invalidated,
667                    },
668                    sequence: Some(job.state().sequence),
669                }
670            } else if bytes.starts_with(MUTATION_JOB_RECORD_MAGIC) {
671                let job = decode_mutation_job_record_for_inventory(bytes, key)?;
672                mutation_count = mutation_count
673                    .checked_add(1)
674                    .ok_or(MutationJobError::CorruptProgressStore)?;
675                ProgressJobInventoryRecord {
676                    family: ProgressJobFamily::Mutation,
677                    job_id: job.state().job_id.to_bytes(),
678                    lifecycle: match job.state().status {
679                        MutationJobStatus::Active => ProgressJobLifecycle::Active,
680                        MutationJobStatus::Completed => ProgressJobLifecycle::Completed,
681                        MutationJobStatus::RestartRequired(_) => {
682                            ProgressJobLifecycle::RestartRequired
683                        }
684                    },
685                    sequence: Some(job.state().sequence),
686                }
687            } else {
688                return Err(MutationJobError::CorruptProgressStore);
689            };
690            records.push(record);
691        }
692
693        let decoded_count = integrity_count
694            .checked_add(resumable_count)
695            .and_then(|count| count.checked_add(mutation_count))
696            .ok_or(MutationJobError::CorruptProgressStore)?;
697        if decoded_count != retained_count || u64::try_from(records.len()) != Ok(retained_count) {
698            return Err(MutationJobError::CorruptProgressStore);
699        }
700
701        Ok(ProgressJobInventory {
702            retained_count,
703            hard_limit: MAX_PROGRESS_JOBS_GLOBAL,
704            reserved_integrity_headroom: PROGRESS_JOBS_INTEGRITY_RESERVATION,
705            integrity_count,
706            resumable_count,
707            mutation_count,
708            retained_record_bytes,
709            records,
710        })
711    }
712
713    pub(super) fn remove(&mut self, job_id: IntegrityJobId) -> Result<(), IntegrityJobError> {
714        if self
715            .map
716            .remove(&ProgressRecordKey::from_job_id(job_id))
717            .is_none()
718        {
719            return Err(IntegrityJobError::JobNotFound);
720        }
721        Ok(())
722    }
723
724    pub(super) fn scan_after(
725        &self,
726        checkpoint: Option<IntegrityJobId>,
727        limit: usize,
728    ) -> Result<ProgressScanPage, IntegrityJobError> {
729        if limit == 0 {
730            return Err(IntegrityJobError::CapacityExceeded);
731        }
732        let lower = checkpoint.map_or(PROGRESS_HEADER_KEY, ProgressRecordKey::from_job_id);
733        let mut job_ids = Vec::with_capacity(limit);
734        let mut has_more = false;
735        for entry in self.map.range((Excluded(lower), Unbounded)) {
736            let Ok(job_id) = integrity_job_id_from_record(&entry.value().0) else {
737                continue;
738            };
739            if job_ids.len() == limit {
740                has_more = true;
741                break;
742            }
743            job_ids.push(job_id);
744        }
745        Ok(ProgressScanPage {
746            job_ids,
747            exhausted: !has_more,
748        })
749    }
750
751    fn job_count(&self) -> Result<u64, IntegrityJobError> {
752        self.map
753            .len()
754            .checked_sub(1)
755            .ok_or(IntegrityJobError::CorruptProgressHeader)
756    }
757
758    fn owner_job_count(&self, owner: &IntegrityJobOwner) -> Result<u64, IntegrityJobError> {
759        let mut count = 0_u64;
760        for entry in self.map.iter() {
761            if *entry.key() == PROGRESS_HEADER_KEY {
762                continue;
763            }
764            // Corrupt records already consume one slot from the global hard
765            // capacity, but their owner cannot be trusted. Skipping them here
766            // isolates the failed job without allowing unbounded progress
767            // growth or blocking every other owner from starting work.
768            let Ok(job_id) = IntegrityJobId::try_from_bytes(entry.key().0) else {
769                continue;
770            };
771            let Ok(job) = decode_job_record(&entry.value().0, job_id) else {
772                continue;
773            };
774            if job.owner == *owner {
775                count = count
776                    .checked_add(1)
777                    .ok_or(IntegrityJobError::CapacityExceeded)?;
778            }
779        }
780        Ok(count)
781    }
782}
783
784fn encode_progress_header() -> Vec<u8> {
785    let mut bytes = Vec::with_capacity(PROGRESS_HEADER_BYTES);
786    bytes.extend_from_slice(PROGRESS_HEADER_MAGIC);
787    bytes.push(PROGRESS_HEADER_VERSION);
788    let checksum = crc32c(bytes.as_slice());
789    bytes.extend_from_slice(&checksum.to_be_bytes());
790    bytes
791}
792
793fn decode_progress_header(bytes: &[u8]) -> Result<(), IntegrityJobError> {
794    if bytes.len() != PROGRESS_HEADER_BYTES
795        || !bytes.starts_with(PROGRESS_HEADER_MAGIC)
796        || bytes[PROGRESS_HEADER_MAGIC.len()] != PROGRESS_HEADER_VERSION
797    {
798        return Err(IntegrityJobError::IncompatibleProgressFormat);
799    }
800    let checksum_offset = PROGRESS_HEADER_MAGIC.len() + 1;
801    let mut checksum = [0; 4];
802    checksum.copy_from_slice(&bytes[checksum_offset..]);
803    if u32::from_be_bytes(checksum) != crc32c(&bytes[..checksum_offset]) {
804        return Err(IntegrityJobError::CorruptProgressHeader);
805    }
806    Ok(())
807}
808
809fn encode_job_record(job: &IntegrityJob) -> Result<Vec<u8>, IntegrityJobError> {
810    let payload =
811        encode_integrity_job_payload(job).map_err(|_| IntegrityJobError::CapacityExceeded)?;
812    let total_len = JOB_RECORD_HEADER_BYTES
813        .checked_add(payload.len())
814        .ok_or(IntegrityJobError::CapacityExceeded)?;
815    if total_len > MAX_PROGRESS_RECORD_BYTES as usize {
816        return Err(IntegrityJobError::CapacityExceeded);
817    }
818    let payload_len =
819        u32::try_from(payload.len()).map_err(|_| IntegrityJobError::CapacityExceeded)?;
820    let mut bytes = Vec::with_capacity(total_len);
821    bytes.extend_from_slice(JOB_RECORD_MAGIC);
822    bytes.push(JOB_RECORD_VERSION);
823    bytes.extend_from_slice(&payload_len.to_be_bytes());
824    bytes.extend_from_slice(&crc32c(payload.as_slice()).to_be_bytes());
825    bytes.extend_from_slice(&payload);
826    Ok(bytes)
827}
828
829fn decode_job_record(
830    bytes: &[u8],
831    expected_id: IntegrityJobId,
832) -> Result<IntegrityJob, IntegrityJobError> {
833    if bytes.len() < JOB_RECORD_HEADER_BYTES
834        || !bytes.starts_with(JOB_RECORD_MAGIC)
835        || bytes[JOB_RECORD_MAGIC.len()] != JOB_RECORD_VERSION
836    {
837        return Err(IntegrityJobError::IncompatibleProgressFormat);
838    }
839    if bytes.len() > MAX_PROGRESS_RECORD_BYTES as usize {
840        return Err(IntegrityJobError::CorruptProgressRecord);
841    }
842    let payload_len_offset = JOB_RECORD_MAGIC.len() + 1;
843    let checksum_offset = payload_len_offset + 4;
844    let payload_offset = checksum_offset + 4;
845    let mut payload_len = [0; 4];
846    payload_len.copy_from_slice(&bytes[payload_len_offset..checksum_offset]);
847    if u32::from_be_bytes(payload_len) as usize != bytes.len() - payload_offset {
848        return Err(IntegrityJobError::CorruptProgressRecord);
849    }
850    let payload = &bytes[payload_offset..];
851    let mut checksum = [0; 4];
852    checksum.copy_from_slice(&bytes[checksum_offset..payload_offset]);
853    if u32::from_be_bytes(checksum) != crc32c(payload) {
854        return Err(IntegrityJobError::CorruptProgressRecord);
855    }
856    if payload.len() > MAX_INTEGRITY_JOB_PAYLOAD_BYTES {
857        return Err(IntegrityJobError::CorruptProgressRecord);
858    }
859    let job = decode_integrity_job_payload(payload)
860        .map_err(|_| IntegrityJobError::CorruptProgressRecord)?;
861    if job.id != expected_id {
862        return Err(IntegrityJobError::CorruptProgressRecord);
863    }
864    Ok(job)
865}
866
867fn integrity_job_id_from_record(bytes: &[u8]) -> Result<IntegrityJobId, IntegrityJobError> {
868    if bytes.len() < JOB_RECORD_HEADER_BYTES || !bytes.starts_with(JOB_RECORD_MAGIC) {
869        return Err(IntegrityJobError::CorruptProgressRecord);
870    }
871    let payload_len_offset = JOB_RECORD_MAGIC.len() + 1;
872    let checksum_offset = payload_len_offset + 4;
873    let payload_offset = checksum_offset + 4;
874    let payload = bytes
875        .get(payload_offset..)
876        .ok_or(IntegrityJobError::CorruptProgressRecord)?;
877    let job = decode_integrity_job_payload(payload)
878        .map_err(|_| IntegrityJobError::CorruptProgressRecord)?;
879    decode_job_record(bytes, job.id).map(|job| job.id)
880}
881
882fn encode_resumable_job_record(record: &ResumableJobRecord) -> Result<Vec<u8>, ResumableJobError> {
883    let payload = encode_resumable_job_payload(record)?;
884    let total_len = RESUMABLE_JOB_RECORD_HEADER_BYTES
885        .checked_add(payload.len())
886        .ok_or(ResumableJobError::PayloadTooLarge)?;
887    if total_len > MAX_PROGRESS_RECORD_BYTES as usize {
888        return Err(ResumableJobError::PayloadTooLarge);
889    }
890    let payload_len =
891        u32::try_from(payload.len()).map_err(|_| ResumableJobError::PayloadTooLarge)?;
892    let mut bytes = Vec::with_capacity(total_len);
893    bytes.extend_from_slice(RESUMABLE_JOB_RECORD_MAGIC);
894    bytes.push(RESUMABLE_JOB_RECORD_VERSION);
895    bytes.extend_from_slice(&payload_len.to_be_bytes());
896    bytes.extend_from_slice(&crc32c(payload.as_slice()).to_be_bytes());
897    bytes.extend_from_slice(&payload);
898    Ok(bytes)
899}
900
901fn decode_resumable_job_record(
902    bytes: &[u8],
903    expected_id: ResumableJobId,
904) -> Result<ResumableJobRecord, ResumableJobError> {
905    let record = decode_resumable_job_record_unbound(bytes)?;
906    if record.state().job_id != expected_id {
907        return Err(ResumableJobError::CorruptProgressStore);
908    }
909    Ok(record)
910}
911
912fn decode_resumable_job_record_unbound(
913    bytes: &[u8],
914) -> Result<ResumableJobRecord, ResumableJobError> {
915    if bytes.len() < RESUMABLE_JOB_RECORD_HEADER_BYTES
916        || !bytes.starts_with(RESUMABLE_JOB_RECORD_MAGIC)
917        || bytes[RESUMABLE_JOB_RECORD_MAGIC.len()] != RESUMABLE_JOB_RECORD_VERSION
918    {
919        return Err(ResumableJobError::IncompatibleProgressFormat);
920    }
921    if bytes.len() > MAX_PROGRESS_RECORD_BYTES as usize {
922        return Err(ResumableJobError::CorruptProgressStore);
923    }
924    let payload_len_offset = RESUMABLE_JOB_RECORD_MAGIC.len() + 1;
925    let checksum_offset = payload_len_offset + 4;
926    let payload_offset = checksum_offset + 4;
927    let mut payload_len = [0; 4];
928    payload_len.copy_from_slice(&bytes[payload_len_offset..checksum_offset]);
929    if u32::from_be_bytes(payload_len) as usize != bytes.len() - payload_offset {
930        return Err(ResumableJobError::CorruptProgressStore);
931    }
932    let payload = &bytes[payload_offset..];
933    let mut checksum = [0; 4];
934    checksum.copy_from_slice(&bytes[checksum_offset..payload_offset]);
935    if u32::from_be_bytes(checksum) != crc32c(payload) {
936        return Err(ResumableJobError::CorruptProgressStore);
937    }
938    decode_resumable_job_payload(payload)
939}
940
941fn decode_resumable_job_record_for_inventory(
942    bytes: &[u8],
943    key: ProgressRecordKey,
944) -> Result<ResumableJobRecord, MutationJobError> {
945    let record = decode_resumable_job_record_unbound(bytes)
946        .map_err(|_| MutationJobError::CorruptProgressStore)?;
947    let expected_key = ProgressRecordKey::from_resumable_job_id(record.state().job_id)
948        .map_err(|_| MutationJobError::CorruptProgressStore)?;
949    if key != expected_key {
950        return Err(MutationJobError::CorruptProgressStore);
951    }
952    Ok(record)
953}
954
955fn encode_mutation_job_record(record: &MutationJobRecord) -> Result<Vec<u8>, MutationJobError> {
956    let payload = encode_mutation_job_payload(record)?;
957    let total_len = MUTATION_JOB_RECORD_HEADER_BYTES
958        .checked_add(payload.len())
959        .ok_or(MutationJobError::CapacityExceeded)?;
960    if total_len > MAX_MUTATION_JOB_RECORD_BYTES {
961        return Err(mutation_record_size_error(total_len));
962    }
963    let payload_len =
964        u32::try_from(payload.len()).map_err(|_| MutationJobError::CapacityExceeded)?;
965    let mut bytes = Vec::with_capacity(total_len);
966    bytes.extend_from_slice(MUTATION_JOB_RECORD_MAGIC);
967    bytes.push(MUTATION_JOB_RECORD_VERSION);
968    bytes.extend_from_slice(&payload_len.to_be_bytes());
969    bytes.extend_from_slice(&crc32c(payload.as_slice()).to_be_bytes());
970    bytes.extend_from_slice(&payload);
971    Ok(bytes)
972}
973
974fn decode_mutation_job_record(
975    bytes: &[u8],
976    expected_id: MutationJobId,
977) -> Result<MutationJobRecord, MutationJobError> {
978    let record = decode_mutation_job_record_unbound(bytes)?;
979    if record.state().job_id != expected_id {
980        return Err(MutationJobError::CorruptProgressStore);
981    }
982    Ok(record)
983}
984
985fn decode_mutation_job_record_unbound(bytes: &[u8]) -> Result<MutationJobRecord, MutationJobError> {
986    if bytes.len() < MUTATION_JOB_RECORD_HEADER_BYTES
987        || !bytes.starts_with(MUTATION_JOB_RECORD_MAGIC)
988        || bytes[MUTATION_JOB_RECORD_MAGIC.len()] != MUTATION_JOB_RECORD_VERSION
989    {
990        return Err(MutationJobError::IncompatibleProgressFormat);
991    }
992    if bytes.len() > MAX_MUTATION_JOB_RECORD_BYTES {
993        return Err(MutationJobError::CorruptProgressStore);
994    }
995    let payload_len_offset = MUTATION_JOB_RECORD_MAGIC.len() + 1;
996    let checksum_offset = payload_len_offset + 4;
997    let payload_offset = checksum_offset + 4;
998    let mut payload_len = [0; 4];
999    payload_len.copy_from_slice(&bytes[payload_len_offset..checksum_offset]);
1000    if u32::from_be_bytes(payload_len) as usize != bytes.len() - payload_offset {
1001        return Err(MutationJobError::CorruptProgressStore);
1002    }
1003    let payload = &bytes[payload_offset..];
1004    let mut checksum = [0; 4];
1005    checksum.copy_from_slice(&bytes[checksum_offset..payload_offset]);
1006    if u32::from_be_bytes(checksum) != crc32c(payload) {
1007        return Err(MutationJobError::CorruptProgressStore);
1008    }
1009    decode_mutation_job_payload(payload)
1010}
1011
1012fn decode_mutation_job_record_for_inventory(
1013    bytes: &[u8],
1014    key: ProgressRecordKey,
1015) -> Result<MutationJobRecord, MutationJobError> {
1016    let record = decode_mutation_job_record_unbound(bytes)
1017        .map_err(|_| MutationJobError::CorruptProgressStore)?;
1018    let expected_key = ProgressRecordKey::from_mutation_job_id(record.state().job_id)
1019        .map_err(|_| MutationJobError::CorruptProgressStore)?;
1020    if key != expected_key {
1021        return Err(MutationJobError::CorruptProgressStore);
1022    }
1023    Ok(record)
1024}
1025
1026fn mutation_progress_before_digest(bytes: &[u8]) -> [u8; 32] {
1027    let mut hasher = new_hash_sha256_prefixed(MUTATION_PROGRESS_BEFORE_DIGEST_DOMAIN);
1028    hasher.update(bytes);
1029    finalize_hash_sha256(hasher)
1030}
1031
1032fn mutation_record_size_error(observed: usize) -> MutationJobError {
1033    MutationJobError::PayloadTooLarge {
1034        kind: crate::db::MutationJobPayloadKind::Record,
1035        limit: u64::try_from(MAX_MUTATION_JOB_RECORD_BYTES).unwrap_or(u64::MAX),
1036        observed: u64::try_from(observed).unwrap_or(u64::MAX),
1037    }
1038}
1039
1040const fn map_integrity_store_error(error: IntegrityJobError) -> ResumableJobError {
1041    match error {
1042        IntegrityJobError::IncompatibleProgressFormat => {
1043            ResumableJobError::IncompatibleProgressFormat
1044        }
1045        IntegrityJobError::CapacityExceeded => ResumableJobError::CapacityExceeded,
1046        _ => ResumableJobError::CorruptProgressStore,
1047    }
1048}
1049
1050const fn map_mutation_store_error(error: IntegrityJobError) -> MutationJobError {
1051    match error {
1052        IntegrityJobError::IncompatibleProgressFormat => {
1053            MutationJobError::IncompatibleProgressFormat
1054        }
1055        IntegrityJobError::CapacityExceeded => MutationJobError::CapacityExceeded,
1056        _ => MutationJobError::CorruptProgressStore,
1057    }
1058}
1059
1060pub(in crate::db) fn with_progress_store<C: CanisterKind, R>(
1061    f: impl FnOnce(&mut InspectionProgressStore) -> Result<R, IntegrityJobError>,
1062) -> Result<R, IntegrityJobError> {
1063    let memory = progress_memory::<C>()?;
1064    let mut store = InspectionProgressStore::open(memory)?;
1065    f(&mut store)
1066}
1067
1068pub(in crate::db) fn with_resumable_progress_store<C: CanisterKind, R>(
1069    f: impl FnOnce(&mut InspectionProgressStore) -> Result<R, ResumableJobError>,
1070) -> Result<R, ResumableJobError> {
1071    let memory = progress_memory::<C>().map_err(map_integrity_store_error)?;
1072    let mut store = InspectionProgressStore::open(memory).map_err(map_integrity_store_error)?;
1073    f(&mut store)
1074}
1075
1076pub(in crate::db) fn with_mutation_progress_store<C: CanisterKind, R>(
1077    f: impl FnOnce(&mut InspectionProgressStore) -> Result<R, MutationJobError>,
1078) -> Result<R, MutationJobError> {
1079    let memory = progress_memory::<C>().map_err(map_mutation_store_error)?;
1080    let mut store = InspectionProgressStore::open(memory).map_err(map_mutation_store_error)?;
1081    f(&mut store)
1082}
1083
1084pub(in crate::db) fn preflight_mutation_progress_record_op<C: CanisterKind>(
1085    operation: &MutationProgressRecordOp,
1086) -> Result<(), InternalError> {
1087    with_mutation_progress_store::<C, _>(|store| store.preflight_mutation_progress(operation))
1088        .map_err(|_| InternalError::commit_corruption())
1089}
1090
1091pub(in crate::db) fn apply_mutation_progress_record_op<C: CanisterKind>(
1092    operation: &MutationProgressRecordOp,
1093) -> Result<(), InternalError> {
1094    with_mutation_progress_store::<C, _>(|store| store.apply_mutation_progress(operation))
1095        .map_err(|_| InternalError::commit_corruption())
1096}
1097
1098/// Mechanically publish one mutation-progress replacement already proved by preflight.
1099pub(in crate::db) fn apply_preflighted_mutation_progress_record_op<C: CanisterKind>(
1100    operation: &MutationProgressRecordOp,
1101) -> Result<(), InternalError> {
1102    with_mutation_progress_store::<C, _>(|store| {
1103        store.apply_preflighted_mutation_progress(operation);
1104        Ok(())
1105    })
1106    .map_err(|_| InternalError::commit_corruption())
1107}
1108
1109pub(in crate::db) fn replace_mutation_progress_record_op<C: CanisterKind>(
1110    operation: &MutationProgressRecordOp,
1111) -> Result<(), MutationJobError> {
1112    with_mutation_progress_store::<C, _>(|store| store.replace_mutation_progress(operation))
1113}
1114
1115pub(in crate::db) fn verify_mutation_progress_record_op<C: CanisterKind>(
1116    operation: &MutationProgressRecordOp,
1117) -> Result<(), InternalError> {
1118    with_mutation_progress_store::<C, _>(|store| store.verify_mutation_progress(operation))
1119        .map_err(|_| InternalError::recovery_effect_verification_failed())
1120}
1121
1122#[cfg(test)]
1123fn progress_memory<C: CanisterKind>() -> Result<RuntimeMemory<DefaultMemoryImpl>, IntegrityJobError>
1124{
1125    let resolved_id = C::integrity_progress_memory_id().map_err(|_| IntegrityJobError::Internal)?;
1126    thread_local! {
1127        static MEMORIES: RefCell<
1128            Vec<(u8, &'static str, RuntimeMemory<DefaultMemoryImpl>)>
1129        > = const { RefCell::new(Vec::new()) };
1130    }
1131
1132    MEMORIES.with(|memories| {
1133        let mut memories = memories.borrow_mut();
1134        if let Some((_, _, memory)) = memories
1135            .iter()
1136            .find(|(id, key, _)| *id == resolved_id && *key == C::INTEGRITY_PROGRESS_STABLE_KEY)
1137        {
1138            return Ok(memory.clone());
1139        }
1140        let memory = crate::testing::test_memory(resolved_id);
1141        memories.push((
1142            resolved_id,
1143            C::INTEGRITY_PROGRESS_STABLE_KEY,
1144            memory.clone(),
1145        ));
1146        Ok(memory)
1147    })
1148}
1149
1150#[cfg(not(test))]
1151fn progress_memory<C: CanisterKind>() -> Result<RuntimeMemory<DefaultMemoryImpl>, IntegrityJobError>
1152{
1153    open_default_memory_manager_memory_by_key(C::INTEGRITY_PROGRESS_STABLE_KEY)
1154        .map_err(|_| IntegrityJobError::Internal)
1155}
1156
1157#[cfg(test)]
1158mod tests {
1159    use super::*;
1160    use crate::{
1161        db::{
1162            MutationJobAdvanceRequest, MutationJobIdempotencyKey, MutationJobPhase,
1163            MutationJobRestartReason, MutationJobStatus, ReadSetRevisionProof,
1164            ReadSetStoreIdentity, ReadSetStoreRevision,
1165            integrity::progress_codec::current_job_codec_fixture,
1166            mutation_job::MutationJobTransition,
1167        },
1168        testing::test_memory,
1169    };
1170    use ic_memory::ic_stable_structures::Memory;
1171
1172    fn current_resumable_record() -> ResumableJobRecord {
1173        let proof = ReadSetRevisionProof::from_parts(
1174            [1; 16],
1175            7,
1176            1,
1177            [2; 32],
1178            vec![ReadSetStoreRevision::new(
1179                ReadSetStoreIdentity::from_bytes([3; 32]),
1180                11,
1181                13,
1182            )],
1183        )
1184        .expect("bounded canonical proof should admit");
1185        ResumableJobRecord::new(
1186            ResumableJobId::try_from_bytes([4; 32])
1187                .expect("nonzero resumable job identity should admit"),
1188            proof,
1189            vec![5, 6],
1190        )
1191        .expect("current resumable record should admit")
1192    }
1193
1194    fn mutation_job_id(byte: u8) -> MutationJobId {
1195        MutationJobId::try_from_bytes([byte; 32]).expect("nonzero mutation job id should admit")
1196    }
1197
1198    fn current_mutation_record(byte: u8) -> MutationJobRecord {
1199        MutationJobRecord::new(mutation_job_id(byte), vec![1, 2, 3], vec![4, 5])
1200            .expect("current mutation record should admit")
1201    }
1202
1203    fn current_integrity_record(byte: u8) -> IntegrityJob {
1204        let mut job = current_job_codec_fixture();
1205        let job_id = IntegrityJobId::try_from_bytes([byte; 32])
1206            .expect("nonzero integrity job id should admit");
1207        job.id = job_id;
1208        match &mut job.last_receipt.receipt {
1209            crate::db::IntegrityJobReceipt::Page(page) => page.job_id = job_id,
1210            crate::db::IntegrityJobReceipt::Abort(receipt) => receipt.job_id = job_id,
1211        }
1212        job.validate()
1213            .expect("rewritten integrity job should admit");
1214        job
1215    }
1216
1217    fn insert_mutation_record(store: &mut InspectionProgressStore, record: &MutationJobRecord) {
1218        assert!(matches!(
1219            store
1220                .insert_mutation(record)
1221                .expect("mutation record should insert"),
1222            InsertMutationJobResult::Inserted,
1223        ));
1224    }
1225
1226    fn mutation_request(byte: u8, sequence: u64, key: &str) -> MutationJobAdvanceRequest {
1227        MutationJobAdvanceRequest::new(
1228            mutation_job_id(byte),
1229            sequence,
1230            MutationJobIdempotencyKey::new(key).expect("bounded replay key should admit"),
1231        )
1232    }
1233
1234    #[test]
1235    fn progress_header_rejects_future_version_and_checksum_corruption() {
1236        let mut future = encode_progress_header();
1237        future[PROGRESS_HEADER_MAGIC.len()] = PROGRESS_HEADER_VERSION + 1;
1238        assert_eq!(
1239            decode_progress_header(&future),
1240            Err(IntegrityJobError::IncompatibleProgressFormat),
1241        );
1242
1243        let mut corrupt = encode_progress_header();
1244        let last = corrupt
1245            .last_mut()
1246            .expect("current progress header has a checksum");
1247        *last ^= 0xff;
1248        assert_eq!(
1249            decode_progress_header(&corrupt),
1250            Err(IntegrityJobError::CorruptProgressHeader),
1251        );
1252    }
1253
1254    #[test]
1255    fn current_job_record_uses_the_direct_bounded_payload() {
1256        let job = current_job_codec_fixture();
1257        let encoded = encode_job_record(&job).expect("current job should encode");
1258
1259        assert_eq!(encoded[JOB_RECORD_MAGIC.len()], 1);
1260        assert!(!encoded[JOB_RECORD_HEADER_BYTES..].starts_with(b"DIDL"));
1261        assert_eq!(
1262            decode_job_record(&encoded, job.id).expect("current job should decode"),
1263            job,
1264        );
1265
1266        let mut corrupt = encoded;
1267        let last = corrupt
1268            .last_mut()
1269            .expect("current job record has a payload");
1270        *last ^= 0xff;
1271        assert_eq!(
1272            decode_job_record(&corrupt, job.id),
1273            Err(IntegrityJobError::CorruptProgressRecord),
1274        );
1275    }
1276
1277    #[test]
1278    fn current_resumable_record_is_direct_bounded_and_checksum_protected() {
1279        let record = current_resumable_record();
1280        let encoded =
1281            encode_resumable_job_record(&record).expect("current resumable record should encode");
1282
1283        assert_eq!(encoded.len(), 175);
1284        assert_eq!(encoded[RESUMABLE_JOB_RECORD_MAGIC.len()], 1);
1285        assert!(!encoded[RESUMABLE_JOB_RECORD_HEADER_BYTES..].starts_with(b"DIDL"));
1286        assert_eq!(
1287            decode_resumable_job_record(&encoded, record.state().job_id)
1288                .expect("current resumable record should decode"),
1289            record,
1290        );
1291
1292        let mut future = encoded.clone();
1293        future[RESUMABLE_JOB_RECORD_MAGIC.len()] = RESUMABLE_JOB_RECORD_VERSION + 1;
1294        assert_eq!(
1295            decode_resumable_job_record(&future, record.state().job_id),
1296            Err(ResumableJobError::IncompatibleProgressFormat),
1297        );
1298
1299        let mut corrupt = encoded;
1300        let last = corrupt
1301            .last_mut()
1302            .expect("current resumable record has a payload");
1303        *last ^= 0xff;
1304        assert_eq!(
1305            decode_resumable_job_record(&corrupt, record.state().job_id),
1306            Err(ResumableJobError::CorruptProgressStore),
1307        );
1308    }
1309
1310    #[test]
1311    fn current_mutation_record_is_distinct_bounded_and_checksum_protected() {
1312        let record = current_mutation_record(7);
1313        let encoded =
1314            encode_mutation_job_record(&record).expect("current mutation record should encode");
1315
1316        assert_eq!(encoded.len(), 97);
1317        assert_eq!(encoded[MUTATION_JOB_RECORD_MAGIC.len()], 1);
1318        assert!(!encoded[MUTATION_JOB_RECORD_HEADER_BYTES..].starts_with(b"DIDL"));
1319        assert_eq!(
1320            decode_mutation_job_record(&encoded, record.state().job_id)
1321                .expect("current mutation record should decode"),
1322            record,
1323        );
1324
1325        let mut future = encoded.clone();
1326        future[MUTATION_JOB_RECORD_MAGIC.len()] = MUTATION_JOB_RECORD_VERSION + 1;
1327        assert_eq!(
1328            decode_mutation_job_record(&future, record.state().job_id),
1329            Err(MutationJobError::IncompatibleProgressFormat),
1330        );
1331
1332        let mut corrupt = encoded;
1333        let last = corrupt
1334            .last_mut()
1335            .expect("current mutation record has a payload");
1336        *last ^= 0xff;
1337        assert_eq!(
1338            decode_mutation_job_record(&corrupt, record.state().job_id),
1339            Err(MutationJobError::CorruptProgressStore),
1340        );
1341
1342        let mut oversized = vec![0; MAX_MUTATION_JOB_RECORD_BYTES + 1];
1343        oversized[..MUTATION_JOB_RECORD_MAGIC.len()].copy_from_slice(MUTATION_JOB_RECORD_MAGIC);
1344        oversized[MUTATION_JOB_RECORD_MAGIC.len()] = MUTATION_JOB_RECORD_VERSION;
1345        assert_eq!(
1346            decode_mutation_job_record(&oversized, record.state().job_id),
1347            Err(MutationJobError::CorruptProgressStore),
1348        );
1349    }
1350
1351    #[test]
1352    fn mutation_progress_replacement_is_exact_idempotent_and_fail_closed() {
1353        let before = current_mutation_record(21);
1354        let (after, _) = before
1355            .apply_transition(
1356                &mutation_request(21, 0, "atomic-forward"),
1357                MutationJobTransition::new(
1358                    MutationJobStatus::Active,
1359                    MutationJobPhase::Forward,
1360                    vec![9],
1361                    8,
1362                    3,
1363                    0,
1364                ),
1365            )
1366            .expect("bounded atomic successor should admit");
1367        let operation = MutationProgressRecordOp::replace(&before, &after)
1368            .expect("exact mutation progress replacement should admit");
1369        let mut store = InspectionProgressStore::open(test_memory(252))
1370            .expect("isolated progress store should open");
1371        assert!(matches!(
1372            store
1373                .insert_mutation(&before)
1374                .expect("before record should insert"),
1375            InsertMutationJobResult::Inserted,
1376        ));
1377
1378        store
1379            .preflight_mutation_progress(&operation)
1380            .expect("exact before bytes should preflight");
1381        store
1382            .apply_mutation_progress(&operation)
1383            .expect("exact before bytes should advance");
1384        store
1385            .apply_mutation_progress(&operation)
1386            .expect("exact after bytes should replay idempotently");
1387        store
1388            .verify_mutation_progress(&operation)
1389            .expect("exact after bytes should verify");
1390        assert_eq!(
1391            store
1392                .load_mutation(before.state().job_id)
1393                .expect("advanced record should load"),
1394            after,
1395        );
1396        assert_eq!(
1397            store.preflight_mutation_progress(&operation),
1398            Err(MutationJobError::CorruptProgressStore),
1399            "opening a new marker against after-state must not reset progress",
1400        );
1401
1402        let (unexpected, _) = after
1403            .apply_transition(
1404                &mutation_request(21, 1, "unexpected"),
1405                MutationJobTransition::new(
1406                    MutationJobStatus::Active,
1407                    MutationJobPhase::Forward,
1408                    vec![10],
1409                    1,
1410                    0,
1411                    0,
1412                ),
1413            )
1414            .expect("third valid state should admit");
1415        store
1416            .replace_mutation(&unexpected)
1417            .expect("test should install neither-side state");
1418        assert_eq!(
1419            store.apply_mutation_progress(&operation),
1420            Err(MutationJobError::CorruptProgressStore),
1421        );
1422        assert_eq!(
1423            store.verify_mutation_progress(&operation),
1424            Err(MutationJobError::CorruptProgressStore),
1425        );
1426    }
1427
1428    #[test]
1429    fn mutation_record_sizes_are_fixed_for_current_and_maximal_states() {
1430        let initial = current_mutation_record(8);
1431        let (active, _) = initial
1432            .apply_transition(
1433                &mutation_request(8, 0, "forward-0"),
1434                MutationJobTransition::new(
1435                    MutationJobStatus::Active,
1436                    MutationJobPhase::Verify,
1437                    vec![6],
1438                    13,
1439                    4,
1440                    0,
1441                ),
1442            )
1443            .expect("bounded active transition should admit");
1444        let (completed, _) = active
1445            .apply_transition(
1446                &mutation_request(8, 1, "verify-0"),
1447                MutationJobTransition::new(
1448                    MutationJobStatus::Completed,
1449                    MutationJobPhase::Verify,
1450                    Vec::new(),
1451                    9,
1452                    0,
1453                    0,
1454                ),
1455            )
1456            .expect("bounded completion should admit");
1457        let (restart, _) = initial
1458            .apply_transition(
1459                &mutation_request(8, 0, "restart"),
1460                MutationJobTransition::new(
1461                    MutationJobStatus::RestartRequired(
1462                        MutationJobRestartReason::AcceptedSchemaChanged,
1463                    ),
1464                    MutationJobPhase::Forward,
1465                    Vec::new(),
1466                    0,
1467                    0,
1468                    0,
1469                ),
1470            )
1471            .expect("bounded restart should admit");
1472        let maximal_initial = MutationJobRecord::new(
1473            mutation_job_id(9),
1474            vec![1; crate::db::MAX_MUTATION_JOB_INTENT_BYTES],
1475            vec![2; crate::db::MAX_MUTATION_JOB_CONTINUATION_BYTES],
1476        )
1477        .expect("maximum initial record should admit");
1478        let (maximal_active, _) = maximal_initial
1479            .apply_transition(
1480                &MutationJobAdvanceRequest::new(
1481                    mutation_job_id(9),
1482                    0,
1483                    MutationJobIdempotencyKey::new(
1484                        "k".repeat(crate::db::MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES),
1485                    )
1486                    .expect("maximum replay key should admit"),
1487                ),
1488                MutationJobTransition::new(
1489                    MutationJobStatus::Active,
1490                    MutationJobPhase::Forward,
1491                    vec![2; crate::db::MAX_MUTATION_JOB_CONTINUATION_BYTES],
1492                    crate::db::MAX_MUTATION_JOB_STEP_KEYS_SCANNED,
1493                    crate::db::MAX_MUTATION_JOB_STEP_ROWS_UPDATED,
1494                    0,
1495                ),
1496            )
1497            .expect("maximum active record should admit");
1498
1499        assert_eq!(
1500            encode_mutation_job_record(&initial).map(|bytes| bytes.len()),
1501            Ok(97)
1502        );
1503        assert_eq!(
1504            encode_mutation_job_record(&active).map(|bytes| bytes.len()),
1505            Ok(167)
1506        );
1507        assert_eq!(
1508            encode_mutation_job_record(&completed).map(|bytes| bytes.len()),
1509            Ok(165),
1510        );
1511        assert_eq!(
1512            encode_mutation_job_record(&restart).map(|bytes| bytes.len()),
1513            Ok(166)
1514        );
1515        assert_eq!(
1516            encode_mutation_job_record(&maximal_initial).map(|bytes| bytes.len()),
1517            Ok(18_524),
1518        );
1519        assert_eq!(
1520            encode_mutation_job_record(&maximal_active).map(|bytes| bytes.len()),
1521            Ok(18_842),
1522        );
1523    }
1524
1525    #[test]
1526    fn mutation_key_domain_and_shared_capacity_reservation_are_enforced() {
1527        let shared_bytes = [11; 32];
1528        let mutation_key = ProgressRecordKey::from_mutation_job_id(
1529            MutationJobId::try_from_bytes(shared_bytes).expect("mutation id should admit"),
1530        )
1531        .expect("mutation progress key should derive");
1532        let resumable_key = ProgressRecordKey::from_resumable_job_id(
1533            ResumableJobId::try_from_bytes(shared_bytes).expect("resumable id should admit"),
1534        )
1535        .expect("resumable progress key should derive");
1536        let integrity_key = ProgressRecordKey::from_job_id(
1537            IntegrityJobId::try_from_bytes(shared_bytes).expect("integrity id should admit"),
1538        );
1539        assert_ne!(mutation_key, resumable_key);
1540        assert_ne!(mutation_key, integrity_key);
1541
1542        let mut store = InspectionProgressStore::open(test_memory(251))
1543            .expect("isolated progress store should open");
1544        store
1545            .insert_resumable(&current_resumable_record())
1546            .expect("generic job should consume one shared slot");
1547        for byte in 1..=54 {
1548            assert!(matches!(
1549                store
1550                    .insert_mutation(&current_mutation_record(byte))
1551                    .expect("record inside shared capacity should insert"),
1552                InsertMutationJobResult::Inserted,
1553            ));
1554        }
1555        assert_eq!(
1556            store
1557                .inventory()
1558                .expect("55 current records should inventory")
1559                .retained_count,
1560            55,
1561        );
1562        assert!(matches!(
1563            store
1564                .insert_mutation(&current_mutation_record(55))
1565                .expect("the 56th non-integrity record should insert"),
1566            InsertMutationJobResult::Inserted,
1567        ));
1568        assert_eq!(
1569            store
1570                .inventory()
1571                .expect("56 current records should inventory")
1572                .retained_count,
1573            56,
1574        );
1575        assert!(matches!(
1576            store.insert_mutation(&current_mutation_record(56)),
1577            Err(MutationJobError::CapacityExceeded),
1578        ));
1579
1580        for byte in 200..=206 {
1581            assert!(matches!(
1582                store
1583                    .insert_new(&current_integrity_record(byte))
1584                    .expect("reserved integrity record should insert"),
1585                InsertJobResult::Inserted,
1586            ));
1587        }
1588        assert_eq!(
1589            store
1590                .inventory()
1591                .expect("63 current records should inventory")
1592                .retained_count,
1593            63,
1594        );
1595        assert!(matches!(
1596            store
1597                .insert_new(&current_integrity_record(207))
1598                .expect("the 64th integrity record should insert"),
1599            InsertJobResult::Inserted,
1600        ));
1601        let full = store.inventory().expect("full store should inventory");
1602        assert_eq!(full.retained_count, 64);
1603        assert_eq!(full.hard_limit, 64);
1604        assert_eq!(full.reserved_integrity_headroom, 8);
1605        assert_eq!(full.integrity_count, 8);
1606        assert_eq!(full.resumable_count, 1);
1607        assert_eq!(full.mutation_count, 55);
1608        assert!(matches!(
1609            store.insert_new(&current_integrity_record(208)),
1610            Err(IntegrityJobError::CapacityExceeded),
1611        ));
1612    }
1613
1614    #[test]
1615    fn progress_stable_growth_is_measured_at_reservation_boundaries() {
1616        const STABLE_PAGE_BYTES: u64 = 65_536;
1617
1618        let memory = test_memory(249);
1619        let mut store = InspectionProgressStore::open(memory.clone())
1620            .expect("isolated progress store should open");
1621        let mut bytes_at_fifty_five = 0;
1622        let mut bytes_at_fifty_six = 0;
1623        for byte in 1..=56 {
1624            assert!(matches!(
1625                store
1626                    .insert_mutation(&current_mutation_record(byte))
1627                    .expect("record inside shared capacity should insert"),
1628                InsertMutationJobResult::Inserted,
1629            ));
1630            if byte == 55 {
1631                bytes_at_fifty_five = memory.size() * STABLE_PAGE_BYTES;
1632            } else if byte == 56 {
1633                bytes_at_fifty_six = memory.size() * STABLE_PAGE_BYTES;
1634            }
1635        }
1636        let mut bytes_at_sixty_three = 0;
1637        for byte in 200..=207 {
1638            assert!(matches!(
1639                store
1640                    .insert_new(&current_integrity_record(byte))
1641                    .expect("record inside integrity reservation should insert"),
1642                InsertJobResult::Inserted,
1643            ));
1644            if byte == 206 {
1645                bytes_at_sixty_three = memory.size() * STABLE_PAGE_BYTES;
1646            }
1647        }
1648        let bytes_at_sixty_four = memory.size() * STABLE_PAGE_BYTES;
1649
1650        assert_eq!(
1651            (
1652                bytes_at_fifty_five,
1653                bytes_at_fifty_six,
1654                bytes_at_sixty_three,
1655                bytes_at_sixty_four,
1656            ),
1657            (38_993_920, 38_993_920, 43_319_296, 43_319_296),
1658        );
1659    }
1660
1661    #[test]
1662    fn mutation_store_load_replay_replace_and_acknowledge_are_exact() {
1663        let mut store = InspectionProgressStore::open(test_memory(250))
1664            .expect("isolated progress store should open");
1665        let initial = current_mutation_record(10);
1666        assert!(matches!(
1667            store
1668                .insert_mutation(&initial)
1669                .expect("initial mutation record should insert"),
1670            InsertMutationJobResult::Inserted,
1671        ));
1672        assert!(matches!(
1673            store
1674                .insert_mutation(&initial)
1675                .expect("duplicate identity should load retained record"),
1676            InsertMutationJobResult::Occupied(record) if *record == initial,
1677        ));
1678        assert_eq!(
1679            store.load_mutation(mutation_job_id(10)),
1680            Ok(initial.clone())
1681        );
1682        assert_eq!(
1683            store.acknowledge_mutation(mutation_job_id(10), 0),
1684            Err(MutationJobError::Active),
1685        );
1686
1687        let request = mutation_request(10, 0, "restart");
1688        let (terminal, receipt) = initial
1689            .apply_transition(
1690                &request,
1691                MutationJobTransition::new(
1692                    MutationJobStatus::RestartRequired(
1693                        MutationJobRestartReason::BatchPolicyChanged,
1694                    ),
1695                    MutationJobPhase::Forward,
1696                    Vec::new(),
1697                    0,
1698                    0,
1699                    0,
1700                ),
1701            )
1702            .expect("terminal transition should admit");
1703        store
1704            .replace_mutation(&terminal)
1705            .expect("terminal replacement should persist");
1706        assert_eq!(
1707            store.load_mutation(mutation_job_id(10)).and_then(|record| {
1708                let replay = record.exact_replay(&request)?;
1709                Ok(replay.cloned())
1710            }),
1711            Ok(Some(receipt)),
1712        );
1713        assert_eq!(
1714            store.acknowledge_mutation(mutation_job_id(10), 0),
1715            Err(MutationJobError::StaleSequence {
1716                expected: 0,
1717                actual: 1,
1718            }),
1719        );
1720        assert_eq!(store.acknowledge_mutation(mutation_job_id(10), 1), Ok(()));
1721        assert_eq!(store.acknowledge_mutation(mutation_job_id(10), 1), Ok(()));
1722        assert_eq!(
1723            store.load_mutation(mutation_job_id(10)),
1724            Err(MutationJobError::NotFound),
1725        );
1726    }
1727
1728    #[cfg(feature = "sql")]
1729    #[test]
1730    fn mutation_cancellation_is_exact_zero_state_and_absent_idempotent() {
1731        let mut store = InspectionProgressStore::open(test_memory(248))
1732            .expect("isolated progress store should open");
1733        let initial = current_mutation_record(31);
1734        insert_mutation_record(&mut store, &initial);
1735        assert_eq!(
1736            store.cancel_unadvanced_mutation(mutation_job_id(31), 1, |_| Ok(())),
1737            Err(MutationJobError::StaleSequence {
1738                expected: 1,
1739                actual: 0,
1740            }),
1741        );
1742        assert_eq!(
1743            store.cancel_unadvanced_mutation(mutation_job_id(31), 0, |_| Ok(())),
1744            Ok(()),
1745        );
1746        assert_eq!(
1747            store.cancel_unadvanced_mutation(mutation_job_id(31), 0, |_| {
1748                Err(MutationJobError::CorruptProgressStore)
1749            }),
1750            Ok(()),
1751            "an absent retry must not invoke continuation validation",
1752        );
1753
1754        let advanced_initial = current_mutation_record(32);
1755        let (advanced, _) = advanced_initial
1756            .apply_transition(
1757                &mutation_request(32, 0, "advanced"),
1758                MutationJobTransition::new(
1759                    MutationJobStatus::Active,
1760                    MutationJobPhase::Forward,
1761                    vec![6],
1762                    1,
1763                    0,
1764                    0,
1765                ),
1766            )
1767            .expect("advanced state should admit");
1768        insert_mutation_record(&mut store, &advanced);
1769        for expected_sequence in [0, 1] {
1770            assert_eq!(
1771                store.cancel_unadvanced_mutation(
1772                    mutation_job_id(32),
1773                    expected_sequence,
1774                    |_| Ok(()),
1775                ),
1776                Err(MutationJobError::StaleSequence {
1777                    expected: 0,
1778                    actual: 1,
1779                }),
1780            );
1781        }
1782
1783        let terminal_initial = current_mutation_record(33);
1784        let (terminal, _) = terminal_initial
1785            .apply_transition(
1786                &mutation_request(33, 0, "terminal"),
1787                MutationJobTransition::new(
1788                    MutationJobStatus::RestartRequired(
1789                        MutationJobRestartReason::BatchPolicyChanged,
1790                    ),
1791                    MutationJobPhase::Forward,
1792                    Vec::new(),
1793                    0,
1794                    0,
1795                    0,
1796                ),
1797            )
1798            .expect("terminal state should admit");
1799        insert_mutation_record(&mut store, &terminal);
1800        assert_eq!(
1801            store.cancel_unadvanced_mutation(mutation_job_id(33), 1, |_| Ok(())),
1802            Err(MutationJobError::StaleSequence {
1803                expected: 0,
1804                actual: 1,
1805            }),
1806        );
1807
1808        let malformed_continuation = current_mutation_record(34);
1809        insert_mutation_record(&mut store, &malformed_continuation);
1810        assert_eq!(
1811            store.cancel_unadvanced_mutation(mutation_job_id(34), 0, |_| {
1812                Err(MutationJobError::CorruptProgressStore)
1813            }),
1814            Err(MutationJobError::CorruptProgressStore),
1815        );
1816        assert_eq!(
1817            store.load_mutation(mutation_job_id(34)),
1818            Ok(malformed_continuation),
1819            "failed validation must retain the record",
1820        );
1821    }
1822
1823    #[test]
1824    fn progress_inventory_is_complete_family_bounded_and_fail_closed() {
1825        let mut store = InspectionProgressStore::open(test_memory(247))
1826            .expect("isolated progress store should open");
1827        let integrity = current_integrity_record(201);
1828        let resumable = current_resumable_record();
1829        let mutation = current_mutation_record(35);
1830        assert!(matches!(
1831            store
1832                .insert_new(&integrity)
1833                .expect("integrity record should insert"),
1834            InsertJobResult::Inserted,
1835        ));
1836        store
1837            .insert_resumable(&resumable)
1838            .expect("resumable record should insert");
1839        assert!(matches!(
1840            store
1841                .insert_mutation(&mutation)
1842                .expect("mutation record should insert"),
1843            InsertMutationJobResult::Inserted,
1844        ));
1845
1846        let inventory = store.inventory().expect("valid records should inventory");
1847        assert_eq!(inventory.retained_count, 3);
1848        assert_eq!(inventory.hard_limit, 64);
1849        assert_eq!(inventory.reserved_integrity_headroom, 8);
1850        assert_eq!(inventory.integrity_count, 1);
1851        assert_eq!(inventory.resumable_count, 1);
1852        assert_eq!(inventory.mutation_count, 1);
1853        let expected_record_bytes = encode_job_record(&integrity)
1854            .expect("integrity record should encode")
1855            .len()
1856            .saturating_add(
1857                encode_resumable_job_record(&resumable)
1858                    .expect("resumable record should encode")
1859                    .len(),
1860            )
1861            .saturating_add(
1862                encode_mutation_job_record(&mutation)
1863                    .expect("mutation record should encode")
1864                    .len(),
1865            );
1866        assert_eq!(
1867            inventory.retained_record_bytes,
1868            u64::try_from(expected_record_bytes).expect("bounded records fit u64")
1869        );
1870        assert_eq!(inventory.records.len(), 3);
1871        assert!(inventory.records.iter().all(|record| {
1872            record.lifecycle == ProgressJobLifecycle::Active && record.sequence == Some(0)
1873        }));
1874        assert!(inventory.records.iter().any(|record| {
1875            record.family == ProgressJobFamily::Integrity
1876                && record.job_id == integrity.id.to_bytes()
1877        }));
1878        assert!(inventory.records.iter().any(|record| {
1879            record.family == ProgressJobFamily::Resumable
1880                && record.job_id == resumable.state().job_id.to_bytes()
1881        }));
1882        assert!(inventory.records.iter().any(|record| {
1883            record.family == ProgressJobFamily::Mutation
1884                && record.job_id == mutation.state().job_id.to_bytes()
1885        }));
1886
1887        store.map.insert(
1888            ProgressRecordKey([202; 32]),
1889            ProgressRecordBytes(b"undecodable-retained-slot".to_vec()),
1890        );
1891        assert_eq!(
1892            store.inventory(),
1893            Err(MutationJobError::CorruptProgressStore),
1894            "one undecodable slot must fail the whole inventory",
1895        );
1896    }
1897
1898    #[test]
1899    fn integrity_scan_skips_other_progress_record_families() {
1900        let mut store = InspectionProgressStore::open(test_memory(252))
1901            .expect("isolated progress store should open");
1902        let integrity = current_job_codec_fixture();
1903        assert!(matches!(
1904            store
1905                .insert_new(&integrity)
1906                .expect("integrity job should insert"),
1907            InsertJobResult::Inserted,
1908        ));
1909        store
1910            .insert_resumable(&current_resumable_record())
1911            .expect("generic resumable job should insert");
1912        assert!(matches!(
1913            store
1914                .insert_mutation(&current_mutation_record(12))
1915                .expect("mutation job should insert"),
1916            InsertMutationJobResult::Inserted,
1917        ));
1918
1919        let page = store
1920            .scan_after(None, 8)
1921            .expect("integrity scan should ignore other record families");
1922        assert_eq!(page.job_ids, vec![integrity.id]);
1923        assert!(page.exhausted);
1924    }
1925}