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