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