Skip to main content

icydb_core/db/
mutation_job.rs

1//! Module: db::mutation_job
2//! Responsibility: bounded durable mutation-job lifecycle, replay, and current payload codec.
3//! Does not own: SQL lowering, target traversal, row mutation, or commit-marker recovery.
4//! Boundary: trusted mutation coordinator -> excluded progress-record envelope.
5
6#[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
16/// Maximum UTF-8 bytes in one mutation-job idempotency key.
17pub const MAX_MUTATION_JOB_IDEMPOTENCY_KEY_BYTES: usize = 256;
18/// Maximum current engine-continuation bytes retained by one mutation job.
19pub const MAX_MUTATION_JOB_CONTINUATION_BYTES: usize = 2 * 1024;
20/// Maximum canonical accepted-intent bytes retained by one mutation job.
21pub const MAX_MUTATION_JOB_INTENT_BYTES: usize = 16 * 1024;
22/// Maximum encoded replay-receipt bytes retained by one mutation job.
23pub const MAX_MUTATION_JOB_RECEIPT_BYTES: usize = 8 * 1024;
24/// Maximum complete encoded mutation-job record, including its storage envelope.
25pub const MAX_MUTATION_JOB_RECORD_BYTES: usize = 64 * 1024;
26
27/// Maximum authoritative keys examined by one mutation-job advance.
28pub const MAX_MUTATION_JOB_STEP_KEYS_SCANNED: u64 = 4_096;
29/// Maximum target rows changed by one mutation-job advance.
30pub 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/// Nonzero application-owned identity for one durable mutation job incarnation.
34///
35/// An application must allocate a fresh identity for every logical job. An
36/// identity is never reusable after cancellation, acknowledgement, failure,
37/// or an absent-record response.
38#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd)]
39pub struct MutationJobId([u8; 32]);
40
41impl MutationJobId {
42    /// Admit one nonzero application-owned identity.
43    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    /// Return the application-owned identity bytes.
51    #[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/// Bounded request identity reused exactly after a lost response.
65#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
66pub struct MutationJobIdempotencyKey(String);
67
68impl MutationJobIdempotencyKey {
69    /// Admit one nonempty bounded UTF-8 request identity.
70    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    /// Borrow the request identity.
79    #[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/// Terminal reason why a valid retained mutation job cannot continue safely.
93#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
94pub enum MutationJobRestartReason {
95    /// The accepted schema identity changed after the job started.
96    AcceptedSchemaChanged,
97    /// The target allocation or database incarnation changed.
98    TargetAllocationChanged,
99    /// The frozen canonical intent is no longer eligible.
100    IntentIneligible,
101    /// The engine-owned batch-policy identity changed.
102    BatchPolicyChanged,
103    /// The retained current record names an unsupported internal continuation.
104    UnsupportedContinuation,
105    /// The current managed-write time would move a target row backward.
106    ManagedTimestampRegression,
107    /// One valid mutation candidate cannot fit the current fixed page policy.
108    CandidateExceedsBatchPolicy,
109    /// Admitted page work exceeded the execution policy calibrated for it.
110    ExecutionBudgetPolicyExceeded,
111}
112
113/// Bounded reason why target mutation execution failed before commit.
114#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
115pub enum MutationJobTargetFailureReason {
116    /// Writer and page-packer staged-byte authorities diverged.
117    StagingByteBudgetExceeded,
118    /// Another target mutation failure occurred without safe semantic detail.
119    Other,
120}
121
122/// Durable lifecycle of one mutation job.
123#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
124pub enum MutationJobStatus {
125    /// One bounded Forward or Verify step may still run.
126    Active,
127    /// Stable clean Verify exhaustion committed.
128    Completed,
129    /// Authority or policy drift requires a new job.
130    RestartRequired(MutationJobRestartReason),
131}
132
133/// Current convergence phase of one mutation job.
134#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
135pub enum MutationJobPhase {
136    /// Scan authoritative keys and converge stale rows.
137    Forward,
138    /// Prove a clean scan at one unchanged durable target revision.
139    Verify,
140}
141
142/// Public bounded state for one retained mutation job.
143#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
144pub struct MutationJobState {
145    /// Application-owned job identity.
146    pub job_id: MutationJobId,
147    /// Sequence expected by the next non-replay advance.
148    pub sequence: u64,
149    /// Current durable lifecycle.
150    pub status: MutationJobStatus,
151    /// Current convergence phase.
152    pub phase: MutationJobPhase,
153    /// Authoritative keys examined across committed advances.
154    pub keys_scanned_total: u64,
155    /// Rows changed across committed advances.
156    pub rows_updated_total: u64,
157    /// Verify passes restarted because stable convergence was not proven.
158    pub verify_restarts_total: u64,
159}
160
161/// Identity and expected sequence for one idempotent bounded advance.
162#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
163pub struct MutationJobAdvanceRequest {
164    /// Target mutation job.
165    pub job_id: MutationJobId,
166    /// Exact sequence observed before issuing this request.
167    pub expected_sequence: u64,
168    /// Stable request identity reused after a lost reply.
169    pub idempotency_key: MutationJobIdempotencyKey,
170}
171
172impl MutationJobAdvanceRequest {
173    /// Construct one request from already admitted identities.
174    #[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/// Replayable result of one committed bounded advance.
194#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
195pub struct MutationJobAdvanceReceipt {
196    /// Sequence named by the request.
197    pub request_sequence: u64,
198    /// Durable sequence after the advance committed.
199    pub committed_sequence: u64,
200    /// Durable lifecycle after the advance committed.
201    pub status: MutationJobStatus,
202    /// Durable convergence phase after the advance committed.
203    pub phase: MutationJobPhase,
204    /// Authoritative keys examined by this advance.
205    pub keys_scanned: u64,
206    /// Rows changed by this advance.
207    pub rows_updated: u64,
208    /// Authoritative keys examined across committed advances.
209    pub keys_scanned_total: u64,
210    /// Rows changed across committed advances.
211    pub rows_updated_total: u64,
212    /// Verify restarts across committed advances.
213    pub verify_restarts_total: u64,
214}
215
216/// Variable-sized component constrained by the durable record protocol.
217#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
218pub enum MutationJobPayloadKind {
219    /// Canonical accepted mutation intent.
220    Intent,
221    /// Private engine continuation.
222    Continuation,
223    /// Retained replay receipt and request identity.
224    Receipt,
225    /// Complete progress-store record envelope.
226    Record,
227}
228
229/// Typed mutation-job protocol, lifecycle, or persistence failure.
230#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
231pub enum MutationJobError {
232    /// Job identity was all zeroes.
233    InvalidJobId,
234    /// Idempotency key was empty or exceeded its byte bound.
235    InvalidIdempotencyKey,
236    /// A retained job id was reused for different canonical meaning.
237    IdentityConflict,
238    /// The requested job does not exist.
239    NotFound,
240    /// The request did not name the job's current sequence.
241    StaleSequence { expected: u64, actual: u64 },
242    /// Acknowledgement targeted an active job.
243    Active,
244    /// A non-replay advance targeted a completed job.
245    Completed,
246    /// A non-replay advance targeted a job that must restart.
247    RestartRequired(MutationJobRestartReason),
248    /// A bounded persisted component exceeded its engine-owned limit.
249    PayloadTooLarge {
250        kind: MutationJobPayloadKind,
251        limit: u64,
252        observed: u64,
253    },
254    /// Current accepted authority no longer matches the frozen intent.
255    AuthorityMismatch,
256    /// The requested mutation cannot be represented by the fixed-intent engine.
257    IneligibleIntent,
258    /// The shared excluded progress store reached its hard capacity.
259    CapacityExceeded,
260    /// A checked sequence or cumulative counter would overflow.
261    CounterOverflow,
262    /// Retained progress bytes or their state closure were corrupt.
263    CorruptProgressStore,
264    /// Retained progress bytes use an unsupported format version.
265    IncompatibleProgressFormat,
266    /// Commit or recovery evidence cannot prove one exact transition.
267    CommitCorruption,
268    /// Target mutation execution failed before progress committed.
269    TargetMutationFailed(MutationJobTargetFailureReason),
270    /// Target traversal failed before progress committed.
271    TargetQueryFailed,
272    /// An internal database invariant prevented the operation.
273    Internal,
274    /// The enclosing request exhausted aggregate IcyDB work allowance.
275    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    // Identity preparation failures are not evidence of corrupt persisted data
288    // or ineligible syntax. Keep the existing bounded resource error intact.
289    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    /// Prove that cancellation can remove only the exact initial record.
413    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}