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