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 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#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
64pub enum ProgressJobFamily {
65 Integrity,
67 Resumable,
69 Mutation,
71}
72
73#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
75pub enum ProgressJobLifecycle {
76 Active,
78 TerminalPending,
80 Completed,
82 Invalidated,
84 RestartRequired,
86 Terminal,
88}
89
90#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
92pub struct ProgressJobInventoryRecord {
93 pub family: ProgressJobFamily,
95 pub job_id: [u8; 32],
97 pub lifecycle: ProgressJobLifecycle,
99 pub sequence: Option<u64>,
101}
102
103#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
105pub struct ProgressJobInventory {
106 pub retained_count: u64,
108 pub hard_limit: u64,
110 pub reserved_integrity_headroom: u64,
112 pub integrity_count: u64,
114 pub resumable_count: u64,
116 pub mutation_count: u64,
118 pub retained_record_bytes: u64,
120 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#[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 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 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 #[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 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 #[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 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 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
820fn 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
1058pub(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 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(¤t_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(¤t_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(¤t_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(¤t_mutation_record(56)),
1601 Err(MutationJobError::CapacityExceeded),
1602 ));
1603
1604 for byte in 200..=206 {
1605 assert!(matches!(
1606 store
1607 .insert_new(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_resumable_record())
1935 .expect("generic resumable job should insert");
1936 assert!(matches!(
1937 store
1938 .insert_mutation(¤t_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}