1use 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#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
68pub enum ProgressJobFamily {
69 Integrity,
71 Resumable,
73 Mutation,
75}
76
77#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
79pub enum ProgressJobLifecycle {
80 Active,
82 TerminalPending,
84 Completed,
86 Invalidated,
88 RestartRequired,
90 Terminal,
92}
93
94#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
96pub struct ProgressJobInventoryRecord {
97 pub family: ProgressJobFamily,
99 pub job_id: [u8; 32],
101 pub lifecycle: ProgressJobLifecycle,
103 pub sequence: Option<u64>,
105}
106
107#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
109pub struct ProgressJobInventory {
110 pub retained_count: u64,
112 pub hard_limit: u64,
114 pub reserved_integrity_headroom: u64,
116 pub integrity_count: u64,
118 pub resumable_count: u64,
120 pub mutation_count: u64,
122 pub retained_record_bytes: u64,
124 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#[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 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 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 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 #[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 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 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
1098pub(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(¤t_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(¤t_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(¤t_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(¤t_mutation_record(56)),
1577 Err(MutationJobError::CapacityExceeded),
1578 ));
1579
1580 for byte in 200..=206 {
1581 assert!(matches!(
1582 store
1583 .insert_new(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_resumable_record())
1911 .expect("generic resumable job should insert");
1912 assert!(matches!(
1913 store
1914 .insert_mutation(¤t_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}