1#[cfg(feature = "sql")]
7mod intent;
8
9use crate::db::codec::{ByteDecodeError, ByteReader};
10use candid::CandidType;
11use serde::Deserialize;
12use std::{error::Error as StdError, fmt};
13
14#[cfg(feature = "sql")]
15pub(in crate::db) use intent::CanonicalMutationIntent;
16
17pub const MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES: usize = 256;
19pub const MAX_MUTATION_JOB_CONTINUATION_BYTES: usize = 2 * 1024;
21pub const MAX_MUTATION_JOB_INTENT_BYTES: usize = 16 * 1024;
23pub const MAX_MUTATION_JOB_RECEIPT_BYTES: usize = 8 * 1024;
25pub const MAX_MUTATION_JOB_RECORD_BYTES: usize = 64 * 1024;
27
28pub const MAX_MUTATION_JOB_STEP_KEYS_SCANNED: u64 = 4_096;
30pub const MAX_MUTATION_JOB_STEP_ROWS_UPDATED: u64 =
32 crate::db::executor::MAX_MUTATION_PROGRESS_BATCH_ROWS_AT_MAX_INDEX_FANOUT as u64;
33
34#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd)]
40pub struct MutationJobId([u8; 32]);
41
42impl MutationJobId {
43 pub fn try_from_bytes(bytes: [u8; 32]) -> Result<Self, MutationJobError> {
45 if bytes == [0; 32] {
46 return Err(MutationJobError::InvalidJobId);
47 }
48 Ok(Self(bytes))
49 }
50
51 #[must_use]
53 pub const fn to_bytes(self) -> [u8; 32] {
54 self.0
55 }
56
57 pub(in crate::db) fn validate(self) -> Result<(), MutationJobError> {
58 if self.0 == [0; 32] {
59 return Err(MutationJobError::InvalidJobId);
60 }
61 Ok(())
62 }
63}
64
65#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
67pub struct MutationJobIdempotencyKey(String);
68
69impl MutationJobIdempotencyKey {
70 pub fn new(value: impl Into<String>) -> Result<Self, MutationJobError> {
72 let value = value.into();
73 if value.is_empty() || value.len() > MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES {
74 return Err(MutationJobError::InvalidIdempotencyKey);
75 }
76 Ok(Self(value))
77 }
78
79 #[must_use]
81 pub const fn as_str(&self) -> &str {
82 self.0.as_str()
83 }
84
85 const fn validate(&self) -> Result<(), MutationJobError> {
86 if self.0.is_empty() || self.0.len() > MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES {
87 return Err(MutationJobError::InvalidIdempotencyKey);
88 }
89 Ok(())
90 }
91}
92
93#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
95pub enum MutationJobRestartReason {
96 AcceptedSchemaChanged,
98 TargetAllocationChanged,
100 IntentIneligible,
102 BatchPolicyChanged,
104 UnsupportedContinuation,
106 ManagedTimestampRegression,
108 CandidateExceedsBatchPolicy,
110 ExecutionBudgetPolicyExceeded,
112}
113
114#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
116pub enum MutationJobTargetFailureReason {
117 StagingByteBudgetExceeded,
119 Other,
121}
122
123#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
125pub enum MutationJobStatus {
126 Active,
128 Completed,
130 RestartRequired(MutationJobRestartReason),
132}
133
134#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
136pub enum MutationJobPhase {
137 Forward,
139 Verify,
141}
142
143#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
145pub struct MutationJobState {
146 pub job_id: MutationJobId,
148 pub sequence: u64,
150 pub status: MutationJobStatus,
152 pub phase: MutationJobPhase,
154 pub keys_scanned_total: u64,
156 pub rows_updated_total: u64,
158 pub verify_restarts_total: u64,
160}
161
162#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
164pub struct MutationJobAdvanceRequest {
165 pub job_id: MutationJobId,
167 pub expected_sequence: u64,
169 pub idempotency_key: MutationJobIdempotencyKey,
171}
172
173impl MutationJobAdvanceRequest {
174 #[must_use]
176 pub const fn new(
177 job_id: MutationJobId,
178 expected_sequence: u64,
179 idempotency_key: MutationJobIdempotencyKey,
180 ) -> Self {
181 Self {
182 job_id,
183 expected_sequence,
184 idempotency_key,
185 }
186 }
187
188 fn validate(&self) -> Result<(), MutationJobError> {
189 self.job_id.validate()?;
190 self.idempotency_key.validate()
191 }
192}
193
194#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
196pub struct MutationJobAdvanceReceipt {
197 pub request_sequence: u64,
199 pub committed_sequence: u64,
201 pub status: MutationJobStatus,
203 pub phase: MutationJobPhase,
205 pub keys_scanned: u64,
207 pub rows_updated: u64,
209 pub keys_scanned_total: u64,
211 pub rows_updated_total: u64,
213 pub verify_restarts_total: u64,
215}
216
217#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
219pub enum MutationJobPayloadKind {
220 Intent,
222 Continuation,
224 Receipt,
226 Record,
228}
229
230#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
232pub enum MutationJobError {
233 InvalidJobId,
235 InvalidIdempotencyKey,
237 IdentityConflict,
239 NotFound,
241 StaleSequence { expected: u64, actual: u64 },
243 Active,
245 Completed,
247 RestartRequired(MutationJobRestartReason),
249 PayloadTooLarge {
251 kind: MutationJobPayloadKind,
252 limit: u64,
253 observed: u64,
254 },
255 AuthorityMismatch,
257 IneligibleIntent,
259 CapacityExceeded,
261 CounterOverflow,
263 CorruptProgressStore,
265 IncompatibleProgressFormat,
267 CommitCorruption,
269 TargetMutationFailed(MutationJobTargetFailureReason),
271 TargetQueryFailed,
273 Internal,
275 ExecutionBudgetExceeded {
277 resource: u64,
278 limit: u64,
279 observed: u64,
280 scope: u64,
281 lane: u64,
282 normalized_shape_fingerprint_prefix: u64,
283 },
284}
285
286#[cfg(feature = "sql")]
287impl MutationJobError {
288 pub(in crate::db) fn from_internal_error(error: &crate::error::InternalError) -> Self {
291 use icydb_diagnostic_code::{
292 DiagnosticDetail, DiagnosticFactTag as Fact, RuntimeBoundaryCode,
293 };
294
295 if !matches!(
296 error.diagnostic().detail(),
297 Some(DiagnosticDetail::RuntimeBoundary {
298 boundary: RuntimeBoundaryCode::ExecutionBudgetExceeded,
299 })
300 ) {
301 return Self::Internal;
302 }
303 let facts = error.diagnostic_facts();
304 let [
305 Some(resource),
306 Some(limit),
307 Some(observed),
308 Some(scope),
309 Some(lane),
310 Some(prefix),
311 ] = [
312 Fact::BudgetResource,
313 Fact::Limit,
314 Fact::Actual,
315 Fact::ExecutionBudgetScope,
316 Fact::ExecutionLane,
317 Fact::QueryShapeFingerprintPrefix,
318 ]
319 .map(|tag| {
320 facts
321 .iter()
322 .find(|(fact, _)| *fact == tag)
323 .map(|(_, value)| *value)
324 })
325 else {
326 return Self::Internal;
327 };
328 Self::ExecutionBudgetExceeded {
329 resource,
330 limit,
331 observed,
332 scope,
333 lane,
334 normalized_shape_fingerprint_prefix: prefix,
335 }
336 }
337
338 pub(in crate::db) fn from_query_error(error: crate::db::query::intent::QueryError) -> Self {
339 match error {
340 crate::db::query::intent::QueryError::Execute(error) => {
341 if matches!(
342 error.diagnostic().detail(),
343 Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary { .. })
344 ) {
345 Self::IneligibleIntent
346 } else {
347 Self::from_internal_error(error.as_internal())
348 }
349 }
350 _ => Self::IneligibleIntent,
351 }
352 }
353}
354
355impl fmt::Display for MutationJobError {
356 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
357 formatter.write_str("mutation job operation failed")
358 }
359}
360
361impl From<ByteDecodeError> for MutationJobError {
362 fn from(_: ByteDecodeError) -> Self {
363 Self::CorruptProgressStore
364 }
365}
366
367impl StdError for MutationJobError {}
368
369#[derive(Clone, Debug, Eq, PartialEq)]
370struct RetainedMutationJobReceipt {
371 receipt: MutationJobAdvanceReceipt,
372 idempotency_key: MutationJobIdempotencyKey,
373}
374
375#[derive(Clone, Debug, Eq, PartialEq)]
376pub(in crate::db) struct MutationJobRecord {
377 state: MutationJobState,
378 canonical_intent: Vec<u8>,
379 engine_continuation: Vec<u8>,
380 last_receipt: Option<RetainedMutationJobReceipt>,
381}
382
383impl MutationJobRecord {
384 pub(in crate::db) fn new(
385 job_id: MutationJobId,
386 canonical_intent: Vec<u8>,
387 engine_continuation: Vec<u8>,
388 ) -> Result<Self, MutationJobError> {
389 let record = Self {
390 state: MutationJobState {
391 job_id,
392 sequence: 0,
393 status: MutationJobStatus::Active,
394 phase: MutationJobPhase::Forward,
395 keys_scanned_total: 0,
396 rows_updated_total: 0,
397 verify_restarts_total: 0,
398 },
399 canonical_intent,
400 engine_continuation,
401 last_receipt: None,
402 };
403 record.validate()?;
404 Ok(record)
405 }
406
407 pub(in crate::db) const fn state(&self) -> &MutationJobState {
408 &self.state
409 }
410
411 pub(in crate::db) const fn canonical_intent(&self) -> &[u8] {
412 self.canonical_intent.as_slice()
413 }
414
415 pub(in crate::db) const fn engine_continuation(&self) -> &[u8] {
416 self.engine_continuation.as_slice()
417 }
418
419 pub(in crate::db) fn ensure_cancelable_at_sequence(
421 &self,
422 expected_sequence: u64,
423 ) -> Result<&[u8], MutationJobError> {
424 self.validate()?;
425 if self.state.sequence != expected_sequence {
426 return Err(MutationJobError::StaleSequence {
427 expected: expected_sequence,
428 actual: self.state.sequence,
429 });
430 }
431 if self.state.sequence != 0 {
432 return Err(MutationJobError::StaleSequence {
433 expected: 0,
434 actual: self.state.sequence,
435 });
436 }
437 if self.state.status != MutationJobStatus::Active
438 || self.state.phase != MutationJobPhase::Forward
439 || self.state.keys_scanned_total != 0
440 || self.state.rows_updated_total != 0
441 || self.state.verify_restarts_total != 0
442 || self.last_receipt.is_some()
443 {
444 return Err(MutationJobError::CorruptProgressStore);
445 }
446 Ok(self.engine_continuation())
447 }
448
449 pub(in crate::db) fn exact_replay(
450 &self,
451 request: &MutationJobAdvanceRequest,
452 ) -> Result<Option<&MutationJobAdvanceReceipt>, MutationJobError> {
453 self.validate()?;
454 request.validate()?;
455 if request.job_id != self.state.job_id {
456 return Err(MutationJobError::NotFound);
457 }
458 Ok(self.last_receipt.as_ref().and_then(|retained| {
459 (retained.receipt.request_sequence == request.expected_sequence
460 && retained.idempotency_key == request.idempotency_key)
461 .then_some(&retained.receipt)
462 }))
463 }
464
465 pub(in crate::db) fn ensure_can_advance(
466 &self,
467 request: &MutationJobAdvanceRequest,
468 ) -> Result<(), MutationJobError> {
469 self.validate()?;
470 request.validate()?;
471 if request.job_id != self.state.job_id {
472 return Err(MutationJobError::NotFound);
473 }
474 if self.state.sequence != request.expected_sequence {
475 return Err(MutationJobError::StaleSequence {
476 expected: request.expected_sequence,
477 actual: self.state.sequence,
478 });
479 }
480 match self.state.status {
481 MutationJobStatus::Active => Ok(()),
482 MutationJobStatus::Completed => Err(MutationJobError::Completed),
483 MutationJobStatus::RestartRequired(reason) => {
484 Err(MutationJobError::RestartRequired(reason))
485 }
486 }
487 }
488
489 pub(in crate::db) fn apply_transition(
490 &self,
491 request: &MutationJobAdvanceRequest,
492 transition: MutationJobTransition,
493 ) -> Result<(Self, MutationJobAdvanceReceipt), MutationJobError> {
494 self.ensure_can_advance(request)?;
495 transition.validate(self.state.phase)?;
496 let committed_sequence = self
497 .state
498 .sequence
499 .checked_add(1)
500 .ok_or(MutationJobError::CounterOverflow)?;
501 let keys_scanned_total = self
502 .state
503 .keys_scanned_total
504 .checked_add(transition.keys_scanned)
505 .ok_or(MutationJobError::CounterOverflow)?;
506 let rows_updated_total = self
507 .state
508 .rows_updated_total
509 .checked_add(transition.rows_updated)
510 .ok_or(MutationJobError::CounterOverflow)?;
511 let verify_restarts_total = self
512 .state
513 .verify_restarts_total
514 .checked_add(transition.verify_restarts)
515 .ok_or(MutationJobError::CounterOverflow)?;
516 let receipt = MutationJobAdvanceReceipt {
517 request_sequence: request.expected_sequence,
518 committed_sequence,
519 status: transition.status,
520 phase: transition.phase,
521 keys_scanned: transition.keys_scanned,
522 rows_updated: transition.rows_updated,
523 keys_scanned_total,
524 rows_updated_total,
525 verify_restarts_total,
526 };
527 let record = Self {
528 state: MutationJobState {
529 job_id: self.state.job_id,
530 sequence: committed_sequence,
531 status: transition.status,
532 phase: transition.phase,
533 keys_scanned_total,
534 rows_updated_total,
535 verify_restarts_total,
536 },
537 canonical_intent: self.canonical_intent.clone(),
538 engine_continuation: transition.engine_continuation,
539 last_receipt: Some(RetainedMutationJobReceipt {
540 receipt: receipt.clone(),
541 idempotency_key: request.idempotency_key.clone(),
542 }),
543 };
544 record.validate()?;
545 Ok((record, receipt))
546 }
547
548 pub(in crate::db) fn validate(&self) -> Result<(), MutationJobError> {
549 self.state.job_id.validate()?;
550 validate_nonempty_bytes(
551 &self.canonical_intent,
552 MAX_MUTATION_JOB_INTENT_BYTES,
553 MutationJobPayloadKind::Intent,
554 )?;
555 validate_bytes(
556 &self.engine_continuation,
557 MAX_MUTATION_JOB_CONTINUATION_BYTES,
558 MutationJobPayloadKind::Continuation,
559 )?;
560 if self.state.rows_updated_total > self.state.keys_scanned_total
561 || matches!(self.state.status, MutationJobStatus::Active)
562 && self.engine_continuation.is_empty()
563 || matches!(self.state.status, MutationJobStatus::Completed)
564 && (self.state.phase != MutationJobPhase::Verify
565 || !self.engine_continuation.is_empty())
566 || matches!(self.state.status, MutationJobStatus::RestartRequired(_))
567 && !self.engine_continuation.is_empty()
568 {
569 return Err(MutationJobError::CorruptProgressStore);
570 }
571 match &self.last_receipt {
572 None => {
573 if self.state.sequence != 0
574 || self.state.status != MutationJobStatus::Active
575 || self.state.phase != MutationJobPhase::Forward
576 || self.state.keys_scanned_total != 0
577 || self.state.rows_updated_total != 0
578 || self.state.verify_restarts_total != 0
579 {
580 return Err(MutationJobError::CorruptProgressStore);
581 }
582 }
583 Some(retained) => {
584 retained.idempotency_key.validate()?;
585 validate_receipt(&retained.receipt)?;
586 if retained.receipt.committed_sequence != self.state.sequence
587 || retained.receipt.request_sequence.checked_add(1)
588 != Some(retained.receipt.committed_sequence)
589 || retained.receipt.status != self.state.status
590 || retained.receipt.phase != self.state.phase
591 || retained.receipt.keys_scanned_total != self.state.keys_scanned_total
592 || retained.receipt.rows_updated_total != self.state.rows_updated_total
593 || retained.receipt.verify_restarts_total != self.state.verify_restarts_total
594 {
595 return Err(MutationJobError::CorruptProgressStore);
596 }
597 let receipt_bytes = retained_receipt_encoded_len(retained)?;
598 if receipt_bytes > MAX_MUTATION_JOB_RECEIPT_BYTES {
599 return Err(payload_too_large(
600 MutationJobPayloadKind::Receipt,
601 MAX_MUTATION_JOB_RECEIPT_BYTES,
602 receipt_bytes,
603 ));
604 }
605 }
606 }
607 Ok(())
608 }
609}
610
611#[derive(Clone, Debug, Eq, PartialEq)]
612pub(in crate::db) struct MutationJobTransition {
613 status: MutationJobStatus,
614 phase: MutationJobPhase,
615 engine_continuation: Vec<u8>,
616 keys_scanned: u64,
617 rows_updated: u64,
618 verify_restarts: u64,
619}
620
621impl MutationJobTransition {
622 pub(in crate::db) const fn new(
623 status: MutationJobStatus,
624 phase: MutationJobPhase,
625 engine_continuation: Vec<u8>,
626 keys_scanned: u64,
627 rows_updated: u64,
628 verify_restarts: u64,
629 ) -> Self {
630 Self {
631 status,
632 phase,
633 engine_continuation,
634 keys_scanned,
635 rows_updated,
636 verify_restarts,
637 }
638 }
639
640 fn validate(&self, previous_phase: MutationJobPhase) -> Result<(), MutationJobError> {
641 validate_bytes(
642 &self.engine_continuation,
643 MAX_MUTATION_JOB_CONTINUATION_BYTES,
644 MutationJobPayloadKind::Continuation,
645 )?;
646 let expected_verify_restarts = u64::from(
647 previous_phase == MutationJobPhase::Verify
648 && self.phase == MutationJobPhase::Forward
649 && self.status == MutationJobStatus::Active,
650 );
651 if self.keys_scanned > MAX_MUTATION_JOB_STEP_KEYS_SCANNED
652 || self.rows_updated > MAX_MUTATION_JOB_STEP_ROWS_UPDATED
653 || self.rows_updated > self.keys_scanned
654 || matches!(self.status, MutationJobStatus::Active)
655 && self.engine_continuation.is_empty()
656 || previous_phase == MutationJobPhase::Verify && self.rows_updated != 0
657 || self.verify_restarts != expected_verify_restarts
658 || matches!(self.status, MutationJobStatus::Completed)
659 && (previous_phase != MutationJobPhase::Verify
660 || self.phase != MutationJobPhase::Verify
661 || self.rows_updated != 0
662 || !self.engine_continuation.is_empty())
663 || matches!(self.status, MutationJobStatus::RestartRequired(_))
664 && (self.phase != previous_phase
665 || !self.engine_continuation.is_empty()
666 || self.keys_scanned != 0
667 || self.rows_updated != 0)
668 {
669 return Err(MutationJobError::CorruptProgressStore);
670 }
671 Ok(())
672 }
673}
674
675pub(in crate::db) fn encode_mutation_job_payload(
676 record: &MutationJobRecord,
677) -> Result<Vec<u8>, MutationJobError> {
678 record.validate()?;
679 let mut bytes = Vec::new();
680 bytes.extend_from_slice(&record.state.job_id.to_bytes());
681 bytes.extend_from_slice(&record.state.sequence.to_be_bytes());
682 write_status(&mut bytes, record.state.status);
683 write_phase(&mut bytes, record.state.phase);
684 bytes.extend_from_slice(&record.state.keys_scanned_total.to_be_bytes());
685 bytes.extend_from_slice(&record.state.rows_updated_total.to_be_bytes());
686 bytes.extend_from_slice(&record.state.verify_restarts_total.to_be_bytes());
687 write_bytes(&mut bytes, &record.canonical_intent)?;
688 write_bytes(&mut bytes, &record.engine_continuation)?;
689 match &record.last_receipt {
690 None => bytes.push(0),
691 Some(retained) => {
692 bytes.push(1);
693 let receipt = &retained.receipt;
694 bytes.extend_from_slice(&receipt.request_sequence.to_be_bytes());
695 bytes.extend_from_slice(&receipt.committed_sequence.to_be_bytes());
696 write_status(&mut bytes, receipt.status);
697 write_phase(&mut bytes, receipt.phase);
698 bytes.extend_from_slice(&receipt.keys_scanned.to_be_bytes());
699 bytes.extend_from_slice(&receipt.rows_updated.to_be_bytes());
700 bytes.extend_from_slice(&receipt.keys_scanned_total.to_be_bytes());
701 bytes.extend_from_slice(&receipt.rows_updated_total.to_be_bytes());
702 bytes.extend_from_slice(&receipt.verify_restarts_total.to_be_bytes());
703 write_bytes(&mut bytes, retained.idempotency_key.as_str().as_bytes())?;
704 }
705 }
706 Ok(bytes)
707}
708
709pub(in crate::db) fn decode_mutation_job_payload(
710 bytes: &[u8],
711) -> Result<MutationJobRecord, MutationJobError> {
712 if bytes.len() > MAX_MUTATION_JOB_RECORD_BYTES {
713 return Err(MutationJobError::CorruptProgressStore);
714 }
715 let mut reader = ByteReader::new(bytes);
716 let job_id = MutationJobId::try_from_bytes(reader.read_array()?)
717 .map_err(|_| MutationJobError::CorruptProgressStore)?;
718 let state = MutationJobState {
719 job_id,
720 sequence: reader.read_u64()?,
721 status: read_status(&mut reader)?,
722 phase: read_phase(&mut reader)?,
723 keys_scanned_total: reader.read_u64()?,
724 rows_updated_total: reader.read_u64()?,
725 verify_restarts_total: reader.read_u64()?,
726 };
727 let canonical_intent = reader
728 .read_bounded_len_prefixed_bytes(MAX_MUTATION_JOB_INTENT_BYTES)?
729 .to_vec();
730 let engine_continuation = reader
731 .read_bounded_len_prefixed_bytes(MAX_MUTATION_JOB_CONTINUATION_BYTES)?
732 .to_vec();
733 let last_receipt = match reader.read_u8()? {
734 0 => None,
735 1 => {
736 let receipt = MutationJobAdvanceReceipt {
737 request_sequence: reader.read_u64()?,
738 committed_sequence: reader.read_u64()?,
739 status: read_status(&mut reader)?,
740 phase: read_phase(&mut reader)?,
741 keys_scanned: reader.read_u64()?,
742 rows_updated: reader.read_u64()?,
743 keys_scanned_total: reader.read_u64()?,
744 rows_updated_total: reader.read_u64()?,
745 verify_restarts_total: reader.read_u64()?,
746 };
747 let idempotency_key = MutationJobIdempotencyKey::new(
748 reader.read_bounded_string(MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES)?,
749 )
750 .map_err(|_| MutationJobError::CorruptProgressStore)?;
751 Some(RetainedMutationJobReceipt {
752 receipt,
753 idempotency_key,
754 })
755 }
756 _ => return Err(MutationJobError::CorruptProgressStore),
757 };
758 reader.finish()?;
759 let record = MutationJobRecord {
760 state,
761 canonical_intent,
762 engine_continuation,
763 last_receipt,
764 };
765 record
766 .validate()
767 .map_err(|_| MutationJobError::CorruptProgressStore)?;
768 Ok(record)
769}
770
771fn validate_receipt(receipt: &MutationJobAdvanceReceipt) -> Result<(), MutationJobError> {
772 if receipt.keys_scanned > MAX_MUTATION_JOB_STEP_KEYS_SCANNED
773 || receipt.rows_updated > MAX_MUTATION_JOB_STEP_ROWS_UPDATED
774 || receipt.rows_updated > receipt.keys_scanned
775 || receipt.rows_updated_total > receipt.keys_scanned_total
776 || receipt.keys_scanned > receipt.keys_scanned_total
777 || receipt.rows_updated > receipt.rows_updated_total
778 || matches!(receipt.status, MutationJobStatus::Completed)
779 && (receipt.phase != MutationJobPhase::Verify || receipt.rows_updated != 0)
780 || matches!(receipt.status, MutationJobStatus::RestartRequired(_))
781 && (receipt.keys_scanned != 0 || receipt.rows_updated != 0)
782 {
783 return Err(MutationJobError::CorruptProgressStore);
784 }
785 Ok(())
786}
787
788fn retained_receipt_encoded_len(
789 retained: &RetainedMutationJobReceipt,
790) -> Result<usize, MutationJobError> {
791 let status_bytes = match retained.receipt.status {
792 MutationJobStatus::RestartRequired(_) => 2,
793 MutationJobStatus::Active | MutationJobStatus::Completed => 1,
794 };
795 8_usize
796 .checked_add(8)
797 .and_then(|value| value.checked_add(status_bytes))
798 .and_then(|value| value.checked_add(1))
799 .and_then(|value| value.checked_add(5 * 8))
800 .and_then(|value| value.checked_add(4))
801 .and_then(|value| value.checked_add(retained.idempotency_key.as_str().len()))
802 .ok_or(MutationJobError::CounterOverflow)
803}
804
805fn validate_nonempty_bytes(
806 value: &[u8],
807 limit: usize,
808 kind: MutationJobPayloadKind,
809) -> Result<(), MutationJobError> {
810 if value.is_empty() {
811 return Err(MutationJobError::CorruptProgressStore);
812 }
813 validate_bytes(value, limit, kind)
814}
815
816fn validate_bytes(
817 value: &[u8],
818 limit: usize,
819 kind: MutationJobPayloadKind,
820) -> Result<(), MutationJobError> {
821 if value.len() > limit {
822 return Err(payload_too_large(kind, limit, value.len()));
823 }
824 Ok(())
825}
826
827fn payload_too_large(
828 kind: MutationJobPayloadKind,
829 limit: usize,
830 observed: usize,
831) -> MutationJobError {
832 MutationJobError::PayloadTooLarge {
833 kind,
834 limit: u64::try_from(limit).unwrap_or(u64::MAX),
835 observed: u64::try_from(observed).unwrap_or(u64::MAX),
836 }
837}
838
839fn write_status(bytes: &mut Vec<u8>, status: MutationJobStatus) {
840 match status {
841 MutationJobStatus::Active => bytes.push(0),
842 MutationJobStatus::Completed => bytes.push(1),
843 MutationJobStatus::RestartRequired(reason) => {
844 bytes.push(2);
845 bytes.push(match reason {
846 MutationJobRestartReason::AcceptedSchemaChanged => 0,
847 MutationJobRestartReason::TargetAllocationChanged => 1,
848 MutationJobRestartReason::IntentIneligible => 2,
849 MutationJobRestartReason::BatchPolicyChanged => 3,
850 MutationJobRestartReason::UnsupportedContinuation => 4,
851 MutationJobRestartReason::ManagedTimestampRegression => 5,
852 MutationJobRestartReason::CandidateExceedsBatchPolicy => 6,
853 MutationJobRestartReason::ExecutionBudgetPolicyExceeded => 7,
854 });
855 }
856 }
857}
858
859fn read_status(reader: &mut ByteReader<'_>) -> Result<MutationJobStatus, MutationJobError> {
860 match reader.read_u8()? {
861 0 => Ok(MutationJobStatus::Active),
862 1 => Ok(MutationJobStatus::Completed),
863 2 => Ok(MutationJobStatus::RestartRequired(
864 match reader.read_u8()? {
865 0 => MutationJobRestartReason::AcceptedSchemaChanged,
866 1 => MutationJobRestartReason::TargetAllocationChanged,
867 2 => MutationJobRestartReason::IntentIneligible,
868 3 => MutationJobRestartReason::BatchPolicyChanged,
869 4 => MutationJobRestartReason::UnsupportedContinuation,
870 5 => MutationJobRestartReason::ManagedTimestampRegression,
871 6 => MutationJobRestartReason::CandidateExceedsBatchPolicy,
872 7 => MutationJobRestartReason::ExecutionBudgetPolicyExceeded,
873 _ => return Err(MutationJobError::CorruptProgressStore),
874 },
875 )),
876 _ => Err(MutationJobError::CorruptProgressStore),
877 }
878}
879
880fn write_phase(bytes: &mut Vec<u8>, phase: MutationJobPhase) {
881 bytes.push(match phase {
882 MutationJobPhase::Forward => 0,
883 MutationJobPhase::Verify => 1,
884 });
885}
886
887fn read_phase(reader: &mut ByteReader<'_>) -> Result<MutationJobPhase, MutationJobError> {
888 match reader.read_u8()? {
889 0 => Ok(MutationJobPhase::Forward),
890 1 => Ok(MutationJobPhase::Verify),
891 _ => Err(MutationJobError::CorruptProgressStore),
892 }
893}
894
895fn write_bytes(bytes: &mut Vec<u8>, value: &[u8]) -> Result<(), MutationJobError> {
896 let len = u32::try_from(value.len()).map_err(|_| MutationJobError::Internal)?;
897 bytes.extend_from_slice(&len.to_be_bytes());
898 bytes.extend_from_slice(value);
899 Ok(())
900}
901
902#[cfg(test)]
903mod tests {
904 use super::*;
905
906 fn job_id() -> MutationJobId {
907 MutationJobId::try_from_bytes([7; 32]).expect("nonzero mutation job id should admit")
908 }
909
910 fn request(sequence: u64, key: &str) -> MutationJobAdvanceRequest {
911 MutationJobAdvanceRequest::new(
912 job_id(),
913 sequence,
914 MutationJobIdempotencyKey::new(key).expect("bounded replay key should admit"),
915 )
916 }
917
918 fn initial_record() -> MutationJobRecord {
919 MutationJobRecord::new(job_id(), vec![1, 2, 3], vec![4, 5])
920 .expect("bounded mutation record should admit")
921 }
922
923 #[test]
924 fn identities_and_variable_components_enforce_current_bounds() {
925 assert_eq!(
926 MutationJobId::try_from_bytes([0; 32]),
927 Err(MutationJobError::InvalidJobId),
928 );
929 assert_eq!(
930 MutationJobIdempotencyKey::new(""),
931 Err(MutationJobError::InvalidIdempotencyKey),
932 );
933 assert!(MutationJobIdempotencyKey::new("k".repeat(256)).is_ok());
934 assert_eq!(
935 MutationJobIdempotencyKey::new("k".repeat(257)),
936 Err(MutationJobError::InvalidIdempotencyKey),
937 );
938 assert_eq!(
939 MutationJobRecord::new(job_id(), vec![1], Vec::new()),
940 Err(MutationJobError::CorruptProgressStore),
941 );
942
943 assert!(MutationJobRecord::new(job_id(), vec![1; 16 * 1024], vec![2; 2 * 1024]).is_ok());
944 assert!(matches!(
945 MutationJobRecord::new(job_id(), vec![1; 16 * 1024 + 1], Vec::new()),
946 Err(MutationJobError::PayloadTooLarge {
947 kind: MutationJobPayloadKind::Intent,
948 ..
949 }),
950 ));
951 assert!(matches!(
952 MutationJobRecord::new(job_id(), vec![1], vec![2; 2 * 1024 + 1]),
953 Err(MutationJobError::PayloadTooLarge {
954 kind: MutationJobPayloadKind::Continuation,
955 ..
956 }),
957 ));
958
959 let maximum_key =
960 MutationJobIdempotencyKey::new("k".repeat(MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES))
961 .expect("maximum replay key should admit");
962 let request = MutationJobAdvanceRequest::new(job_id(), 0, maximum_key);
963 let (record, _) = initial_record()
964 .apply_transition(
965 &request,
966 MutationJobTransition::new(
967 MutationJobStatus::Active,
968 MutationJobPhase::Forward,
969 vec![7],
970 1,
971 0,
972 0,
973 ),
974 )
975 .expect("maximum replay identity should retain");
976 assert_eq!(
977 record
978 .last_receipt
979 .as_ref()
980 .map(retained_receipt_encoded_len),
981 Some(Ok(318)),
982 );
983
984 let maximum_key =
985 MutationJobIdempotencyKey::new("k".repeat(MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES))
986 .expect("maximum replay key should admit");
987 let (restart, _) = initial_record()
988 .apply_transition(
989 &MutationJobAdvanceRequest::new(job_id(), 0, maximum_key),
990 MutationJobTransition::new(
991 MutationJobStatus::RestartRequired(
992 MutationJobRestartReason::BatchPolicyChanged,
993 ),
994 MutationJobPhase::Forward,
995 Vec::new(),
996 0,
997 0,
998 0,
999 ),
1000 )
1001 .expect("maximum restart receipt should retain");
1002 assert_eq!(
1003 restart
1004 .last_receipt
1005 .as_ref()
1006 .map(retained_receipt_encoded_len),
1007 Some(Ok(319)),
1008 );
1009 }
1010
1011 #[test]
1012 fn current_payload_round_trips_every_lifecycle() {
1013 let initial = initial_record();
1014 let (active, _) = initial
1015 .apply_transition(
1016 &request(0, "forward-0"),
1017 MutationJobTransition::new(
1018 MutationJobStatus::Active,
1019 MutationJobPhase::Verify,
1020 vec![6],
1021 13,
1022 4,
1023 0,
1024 ),
1025 )
1026 .expect("bounded active transition should admit");
1027 let (completed, _) = active
1028 .apply_transition(
1029 &request(1, "verify-0"),
1030 MutationJobTransition::new(
1031 MutationJobStatus::Completed,
1032 MutationJobPhase::Verify,
1033 Vec::new(),
1034 9,
1035 0,
1036 0,
1037 ),
1038 )
1039 .expect("clean terminal transition should admit");
1040 let (restart, _) = initial
1041 .apply_transition(
1042 &request(0, "restart"),
1043 MutationJobTransition::new(
1044 MutationJobStatus::RestartRequired(
1045 MutationJobRestartReason::ManagedTimestampRegression,
1046 ),
1047 MutationJobPhase::Forward,
1048 Vec::new(),
1049 0,
1050 0,
1051 0,
1052 ),
1053 )
1054 .expect("typed restart transition should admit");
1055 let (oversized_candidate, _) = initial
1056 .apply_transition(
1057 &request(0, "candidate-exceeds-policy"),
1058 MutationJobTransition::new(
1059 MutationJobStatus::RestartRequired(
1060 MutationJobRestartReason::CandidateExceedsBatchPolicy,
1061 ),
1062 MutationJobPhase::Forward,
1063 Vec::new(),
1064 0,
1065 0,
1066 0,
1067 ),
1068 )
1069 .expect("candidate policy restart should admit");
1070 let (execution_budget, _) = initial
1071 .apply_transition(
1072 &request(0, "execution-budget-policy"),
1073 MutationJobTransition::new(
1074 MutationJobStatus::RestartRequired(
1075 MutationJobRestartReason::ExecutionBudgetPolicyExceeded,
1076 ),
1077 MutationJobPhase::Forward,
1078 Vec::new(),
1079 0,
1080 0,
1081 0,
1082 ),
1083 )
1084 .expect("execution budget policy restart should admit");
1085
1086 for record in [
1087 initial,
1088 active,
1089 completed,
1090 restart,
1091 oversized_candidate,
1092 execution_budget,
1093 ] {
1094 let bytes = encode_mutation_job_payload(&record)
1095 .expect("current mutation payload should encode");
1096 assert!(!bytes.starts_with(b"DIDL"));
1097 assert_eq!(
1098 decode_mutation_job_payload(&bytes)
1099 .expect("current mutation payload should decode"),
1100 record,
1101 );
1102 }
1103 }
1104
1105 #[test]
1106 fn public_state_request_receipt_and_error_are_candid_compatible() {
1107 let state = initial_record().state().clone();
1108 let request = request(0, "candid-request");
1109 let receipt = MutationJobAdvanceReceipt {
1110 request_sequence: 0,
1111 committed_sequence: 1,
1112 status: MutationJobStatus::Active,
1113 phase: MutationJobPhase::Forward,
1114 keys_scanned: 8,
1115 rows_updated: 3,
1116 keys_scanned_total: 8,
1117 rows_updated_total: 3,
1118 verify_restarts_total: 0,
1119 };
1120 let error = MutationJobError::TargetMutationFailed(
1121 MutationJobTargetFailureReason::StagingByteBudgetExceeded,
1122 );
1123
1124 let state_bytes = candid::encode_one(&state).expect("mutation state should encode");
1125 let request_bytes = candid::encode_one(&request).expect("mutation request should encode");
1126 let receipt_bytes = candid::encode_one(&receipt).expect("mutation receipt should encode");
1127 let error_bytes = candid::encode_one(&error).expect("mutation error should encode");
1128 assert_eq!(
1129 candid::decode_one::<MutationJobState>(&state_bytes)
1130 .expect("mutation state should decode"),
1131 state,
1132 );
1133 assert_eq!(
1134 candid::decode_one::<MutationJobAdvanceRequest>(&request_bytes)
1135 .expect("mutation request should decode"),
1136 request,
1137 );
1138 assert_eq!(
1139 candid::decode_one::<MutationJobAdvanceReceipt>(&receipt_bytes)
1140 .expect("mutation receipt should decode"),
1141 receipt,
1142 );
1143 assert_eq!(
1144 candid::decode_one::<MutationJobError>(&error_bytes)
1145 .expect("mutation error should decode"),
1146 error,
1147 );
1148 }
1149
1150 #[test]
1151 fn exact_replay_precedes_stale_and_terminal_rejection() {
1152 let initial = initial_record();
1153 let (verifying, _) = initial
1154 .apply_transition(
1155 &request(0, "forward-0"),
1156 MutationJobTransition::new(
1157 MutationJobStatus::Active,
1158 MutationJobPhase::Verify,
1159 vec![7],
1160 8,
1161 3,
1162 0,
1163 ),
1164 )
1165 .expect("Forward exhaustion should enter Verify");
1166 let terminal_request = request(1, "verify-0");
1167 let (completed, receipt) = verifying
1168 .apply_transition(
1169 &terminal_request,
1170 MutationJobTransition::new(
1171 MutationJobStatus::Completed,
1172 MutationJobPhase::Verify,
1173 Vec::new(),
1174 8,
1175 0,
1176 0,
1177 ),
1178 )
1179 .expect("terminal transition should admit");
1180
1181 assert_eq!(
1182 completed
1183 .exact_replay(&terminal_request)
1184 .expect("exact replay lookup should succeed"),
1185 Some(&receipt),
1186 );
1187 assert_eq!(
1188 completed.ensure_can_advance(&request(1, "different")),
1189 Err(MutationJobError::StaleSequence {
1190 expected: 1,
1191 actual: 2,
1192 }),
1193 );
1194 assert_eq!(
1195 completed.ensure_can_advance(&request(2, "next")),
1196 Err(MutationJobError::Completed),
1197 );
1198 }
1199
1200 #[test]
1201 fn payload_decode_is_bounded_fallible_and_rejects_trailing_bytes() {
1202 let bytes = encode_mutation_job_payload(&initial_record())
1203 .expect("current mutation payload should encode");
1204 for end in 0..bytes.len() {
1205 assert_eq!(
1206 decode_mutation_job_payload(&bytes[..end]),
1207 Err(MutationJobError::CorruptProgressStore),
1208 "truncation at byte {end} must retain corruption classification",
1209 );
1210 }
1211 let mut trailing = bytes;
1212 trailing.push(0);
1213 assert_eq!(
1214 decode_mutation_job_payload(&trailing),
1215 Err(MutationJobError::CorruptProgressStore),
1216 );
1217 let mut unknown_status = encode_mutation_job_payload(&initial_record())
1218 .expect("current mutation payload should encode");
1219 unknown_status[32 + 8] = u8::MAX;
1220 assert_eq!(
1221 decode_mutation_job_payload(&unknown_status),
1222 Err(MutationJobError::CorruptProgressStore),
1223 );
1224 let mut zero_job_id = encode_mutation_job_payload(&initial_record())
1225 .expect("current mutation payload should encode");
1226 zero_job_id[..32].fill(0);
1227 assert_eq!(
1228 decode_mutation_job_payload(&zero_job_id),
1229 Err(MutationJobError::CorruptProgressStore),
1230 );
1231
1232 let initial = initial_record();
1233 let mut bytes =
1234 encode_mutation_job_payload(&initial).expect("current mutation payload should encode");
1235 let intent_len_offset = 32 + 8 + 1 + 1 + 3 * 8;
1236 let continuation_len_offset = intent_len_offset + 4 + initial.canonical_intent.len();
1237 let continuation_offset = continuation_len_offset + 4;
1238 let continuation_end = continuation_offset + initial.engine_continuation.len();
1239 bytes[continuation_len_offset..continuation_offset].fill(0);
1240 bytes.drain(continuation_offset..continuation_end);
1241 assert_eq!(
1242 decode_mutation_job_payload(&bytes),
1243 Err(MutationJobError::CorruptProgressStore),
1244 );
1245 }
1246
1247 #[test]
1248 fn transition_totals_fail_closed_on_overflow() {
1249 assert_eq!(
1250 initial_record().apply_transition(
1251 &request(0, "empty-active-continuation"),
1252 MutationJobTransition::new(
1253 MutationJobStatus::Active,
1254 MutationJobPhase::Forward,
1255 Vec::new(),
1256 1,
1257 0,
1258 0,
1259 ),
1260 ),
1261 Err(MutationJobError::CorruptProgressStore),
1262 );
1263
1264 let mut record = initial_record();
1265 record.state.keys_scanned_total = u64::MAX;
1266 record.state.rows_updated_total = u64::MAX;
1267 record.last_receipt = Some(RetainedMutationJobReceipt {
1268 receipt: MutationJobAdvanceReceipt {
1269 request_sequence: 0,
1270 committed_sequence: 1,
1271 status: MutationJobStatus::Active,
1272 phase: MutationJobPhase::Forward,
1273 keys_scanned: 1,
1274 rows_updated: 1,
1275 keys_scanned_total: u64::MAX,
1276 rows_updated_total: u64::MAX,
1277 verify_restarts_total: 0,
1278 },
1279 idempotency_key: MutationJobIdempotencyKey::new("prior")
1280 .expect("bounded replay key should admit"),
1281 });
1282 record.state.sequence = 1;
1283 assert_eq!(
1284 record.apply_transition(
1285 &request(1, "overflow"),
1286 MutationJobTransition::new(
1287 MutationJobStatus::Active,
1288 MutationJobPhase::Forward,
1289 vec![7],
1290 1,
1291 1,
1292 0,
1293 ),
1294 ),
1295 Err(MutationJobError::CounterOverflow),
1296 );
1297
1298 let mut sequence_record = initial_record();
1299 sequence_record.state.sequence = u64::MAX;
1300 sequence_record.last_receipt = Some(RetainedMutationJobReceipt {
1301 receipt: MutationJobAdvanceReceipt {
1302 request_sequence: u64::MAX - 1,
1303 committed_sequence: u64::MAX,
1304 status: MutationJobStatus::Active,
1305 phase: MutationJobPhase::Forward,
1306 keys_scanned: 0,
1307 rows_updated: 0,
1308 keys_scanned_total: 0,
1309 rows_updated_total: 0,
1310 verify_restarts_total: 0,
1311 },
1312 idempotency_key: MutationJobIdempotencyKey::new("prior-sequence")
1313 .expect("bounded replay key should admit"),
1314 });
1315 assert_eq!(
1316 sequence_record.apply_transition(
1317 &request(u64::MAX, "sequence-overflow"),
1318 MutationJobTransition::new(
1319 MutationJobStatus::Active,
1320 MutationJobPhase::Forward,
1321 vec![7],
1322 0,
1323 0,
1324 0,
1325 ),
1326 ),
1327 Err(MutationJobError::CounterOverflow),
1328 );
1329 }
1330}