Skip to main content

icydb_core/db/session/
mutation_job.rs

1//! Module: session::mutation_job
2//! Responsibility: charged mutation-job start, state load, phase dispatch, and terminal acknowledgement.
3//! Does not own: SQL lowering, Forward/Verify execution, or authorization.
4//! Boundary: trusted session API -> excluded mutation progress record.
5
6use crate::{
7    db::{
8        DbSession, MutationJobError, MutationJobId, MutationJobState, ProgressJobInventory,
9        executor::budget::{ExecutionBudgetExceeded, HardExecutionContext},
10        integrity::with_mutation_progress_store,
11    },
12    metrics::sink::{MetricsEvent, MutationJobLifecycleEvent, record},
13    traits::CanisterKind,
14};
15use icydb_diagnostic_code::{
16    DiagnosticExecutionBudgetResource, DiagnosticExecutionBudgetScope, DiagnosticExecutionLane,
17};
18
19#[cfg(feature = "sql")]
20use crate::db::{
21    MutationJobAdvanceReceipt, MutationJobAdvanceRequest, MutationJobPhase,
22    executor::budget::with_mutation_execution_budget,
23    integrity::InsertMutationJobResult,
24    mutation_job::{CanonicalMutationIntent, MutationJobRecord},
25    session::sql::validate_current_initial_mutation_job_continuation,
26};
27
28#[cfg(feature = "sql")]
29const MUTATION_JOB_START_SHAPE: u64 = 0x6d75_7461_7465_0100;
30const MUTATION_JOB_LOAD_SHAPE: u64 = 0x6d75_7461_7465_0101;
31const MUTATION_JOB_ACKNOWLEDGE_SHAPE: u64 = 0x6d75_7461_7465_0102;
32#[cfg(feature = "sql")]
33const MUTATION_JOB_ADVANCE_SHAPE: u64 = 0x6d75_7461_7465_0103;
34#[cfg(feature = "sql")]
35const MUTATION_JOB_CANCEL_UNADVANCED_SHAPE: u64 = 0x6d75_7461_7465_0104;
36const PROGRESS_JOB_INVENTORY_SHAPE: u64 = 0x6d75_7461_7465_0105;
37
38impl<C: CanisterKind> DbSession<C> {
39    /// Start one durable trusted fixed SQL mutation job.
40    ///
41    /// SQL is parsed and admitted exactly once for a new identity. The
42    /// catalog-native intent and initial engine checkpoint are durably retained
43    /// before this method returns, and no target row is read or mutated.
44    /// Repeating the same canonical request returns the retained state without
45    /// replacing its statement timestamp or resetting progress.
46    #[cfg(feature = "sql")]
47    pub fn start_trusted_sql_mutation_job(
48        &self,
49        job_id: MutationJobId,
50        sql: &str,
51    ) -> Result<MutationJobState, MutationJobError> {
52        self.with_metrics(|| {
53            job_id.validate()?;
54            self.charge_mutation_job_operation(
55                DiagnosticExecutionLane::Mutation,
56                MUTATION_JOB_START_SHAPE,
57            )?;
58            let prepared = self.prepare_mutation_job_start(job_id, sql)?;
59            let submitted_intent = CanonicalMutationIntent::decode(&prepared.canonical_intent)?;
60            let submitted = MutationJobRecord::new(
61                job_id,
62                prepared.canonical_intent,
63                prepared.engine_continuation,
64            )?;
65            with_mutation_progress_store::<C, _>(|store| {
66                match store.insert_mutation(&submitted)? {
67                    InsertMutationJobResult::Inserted => {
68                        record_mutation_job_lifecycle(MutationJobLifecycleEvent::StartInserted);
69                        Ok(submitted.state().clone())
70                    }
71                    InsertMutationJobResult::Occupied(retained) => {
72                        resolve_occupied_mutation_job_start(&retained, &submitted_intent).inspect(
73                            |_| {
74                                record_mutation_job_lifecycle(
75                                    MutationJobLifecycleEvent::StartExactReplay,
76                                );
77                            },
78                        )
79                    }
80                }
81            })
82        })
83    }
84
85    /// Load the bounded public state for one retained mutation job.
86    pub fn mutation_job_state(
87        &self,
88        job_id: MutationJobId,
89    ) -> Result<MutationJobState, MutationJobError> {
90        self.with_metrics(|| {
91            self.charge_mutation_job_operation(
92                DiagnosticExecutionLane::TrustedRead,
93                MUTATION_JOB_LOAD_SHAPE,
94            )?;
95            with_mutation_progress_store::<C, _>(|store| store.load_mutation(job_id)).map(
96                |record| {
97                    record_mutation_job_lifecycle(MutationJobLifecycleEvent::StateLoaded);
98                    record.state().clone()
99                },
100            )
101        })
102    }
103
104    /// Advance one durable mutation job through one bounded engine-owned step.
105    ///
106    /// The request carries only job identity, expected sequence, and a replay
107    /// key. SQL and continuation bytes remain private IcyDB custody.
108    #[cfg(feature = "sql")]
109    pub fn advance_trusted_mutation_job(
110        &self,
111        request: &MutationJobAdvanceRequest,
112    ) -> Result<MutationJobAdvanceReceipt, MutationJobError> {
113        self.with_metrics(|| {
114            let retained =
115                with_mutation_progress_store::<C, _>(|store| store.load_mutation(request.job_id))?;
116            if let Some(receipt) = retained.exact_replay(request)? {
117                record_mutation_job_lifecycle(MutationJobLifecycleEvent::AdvanceExactReplay);
118                return Ok(receipt.clone());
119            }
120            self.charge_mutation_job_operation(
121                DiagnosticExecutionLane::Mutation,
122                MUTATION_JOB_ADVANCE_SHAPE,
123            )?;
124            retained.ensure_can_advance(request)?;
125            let context = HardExecutionContext::new(
126                DiagnosticExecutionBudgetScope::Execution,
127                DiagnosticExecutionLane::Mutation,
128                MUTATION_JOB_ADVANCE_SHAPE,
129            );
130            with_mutation_execution_budget(
131                context,
132                || match retained.state().phase {
133                    MutationJobPhase::Forward => {
134                        self.advance_mutation_job_forward(&retained, request)
135                    }
136                    MutationJobPhase::Verify => {
137                        self.advance_mutation_job_verify(&retained, request)
138                    }
139                },
140                mutation_job_execution_internal_error,
141            )
142        })
143    }
144
145    /// Remove one terminal mutation job after its result has been consumed.
146    ///
147    /// Repeating acknowledgement after a lost response succeeds when the job
148    /// is already absent. Active jobs and stale terminal sequences fail closed.
149    pub fn acknowledge_mutation_job(
150        &self,
151        job_id: MutationJobId,
152        expected_terminal_sequence: u64,
153    ) -> Result<(), MutationJobError> {
154        self.with_metrics(|| {
155            self.charge_mutation_job_operation(
156                DiagnosticExecutionLane::Mutation,
157                MUTATION_JOB_ACKNOWLEDGE_SHAPE,
158            )?;
159            with_mutation_progress_store::<C, _>(|store| {
160                store.acknowledge_mutation(job_id, expected_terminal_sequence)
161            })
162            .inspect(|()| {
163                record_mutation_job_lifecycle(MutationJobLifecycleEvent::TerminalAcknowledged);
164            })
165        })
166    }
167
168    /// Idempotently remove one exact initial mutation-job record.
169    ///
170    /// Cancellation is available only before any page or receipt exists. A
171    /// logical restart must use a fresh [`MutationJobId`]; absent-record
172    /// success never makes an old identity reusable.
173    #[cfg(feature = "sql")]
174    pub fn cancel_unadvanced_mutation_job(
175        &self,
176        job_id: MutationJobId,
177        expected_sequence: u64,
178    ) -> Result<(), MutationJobError> {
179        self.with_metrics(|| {
180            job_id.validate()?;
181            self.charge_mutation_job_operation(
182                DiagnosticExecutionLane::Mutation,
183                MUTATION_JOB_CANCEL_UNADVANCED_SHAPE,
184            )?;
185            with_mutation_progress_store::<C, _>(|store| {
186                store.cancel_unadvanced_mutation(
187                    job_id,
188                    expected_sequence,
189                    validate_current_initial_mutation_job_continuation,
190                )
191            })
192            .inspect(|()| {
193                record_mutation_job_lifecycle(MutationJobLifecycleEvent::CancelUnadvanced);
194            })
195        })
196    }
197
198    /// Return one complete fail-closed inventory of shared retained progress.
199    ///
200    /// Callers remain responsible for authorization. The result exposes only
201    /// family, job identity, bounded lifecycle, sequence, and capacity facts.
202    pub fn progress_job_inventory(&self) -> Result<ProgressJobInventory, MutationJobError> {
203        self.with_metrics(|| {
204            self.charge_mutation_job_operation(
205                DiagnosticExecutionLane::TrustedRead,
206                PROGRESS_JOB_INVENTORY_SHAPE,
207            )?;
208            with_mutation_progress_store::<C, _>(|store| store.inventory()).inspect(|inventory| {
209                record_mutation_job_lifecycle(MutationJobLifecycleEvent::InventoryLoaded);
210                record(MetricsEvent::MutationJobCapacity {
211                    retained_count: inventory.retained_count,
212                    hard_limit: inventory.hard_limit,
213                    reserved_integrity_headroom: inventory.reserved_integrity_headroom,
214                    integrity_count: inventory.integrity_count,
215                    resumable_count: inventory.resumable_count,
216                    mutation_count: inventory.mutation_count,
217                    retained_record_bytes: inventory.retained_record_bytes,
218                });
219            })
220        })
221    }
222
223    fn charge_mutation_job_operation(
224        &self,
225        lane: DiagnosticExecutionLane,
226        shape: u64,
227    ) -> Result<(), MutationJobError> {
228        self.db
229            .request_execution_scope()
230            .charge(
231                HardExecutionContext::new(DiagnosticExecutionBudgetScope::Execution, lane, shape),
232                DiagnosticExecutionBudgetResource::QueryExecutions,
233                1,
234            )
235            .map_err(mutation_job_execution_budget_error)
236    }
237}
238
239fn record_mutation_job_lifecycle(event: MutationJobLifecycleEvent) {
240    record(MetricsEvent::MutationJobLifecycle { event });
241}
242
243#[cfg(feature = "sql")]
244fn resolve_occupied_mutation_job_start(
245    retained: &MutationJobRecord,
246    submitted_intent: &CanonicalMutationIntent,
247) -> Result<MutationJobState, MutationJobError> {
248    let retained_intent = CanonicalMutationIntent::decode(retained.canonical_intent())?;
249    if retained_intent.same_start_request(submitted_intent) {
250        return Ok(retained.state().clone());
251    }
252    if !retained_intent.same_authority(submitted_intent) {
253        return Err(MutationJobError::AuthorityMismatch);
254    }
255    Err(MutationJobError::IdentityConflict)
256}
257
258const fn mutation_job_execution_budget_error(error: ExecutionBudgetExceeded) -> MutationJobError {
259    MutationJobError::ExecutionBudgetExceeded {
260        resource: error.resource().raw(),
261        limit: error.limit(),
262        observed: error.observed(),
263        scope: error.scope().raw(),
264        lane: error.lane().raw(),
265        normalized_shape_fingerprint_prefix: error.normalized_shape_fingerprint_prefix(),
266    }
267}
268
269#[cfg(feature = "sql")]
270fn mutation_job_execution_internal_error(_error: crate::error::InternalError) -> MutationJobError {
271    MutationJobError::Internal
272}
273
274#[cfg(test)]
275mod tests {
276    use super::*;
277    use crate::{
278        db::{
279            MutationJobAdvanceRequest, MutationJobIdempotencyKey, MutationJobPhase,
280            MutationJobRestartReason, MutationJobStatus, RequestExecutionRoot, StoreRegistry,
281            mutation_job::{MutationJobRecord, MutationJobTransition},
282        },
283        traits::Path,
284    };
285    #[cfg(feature = "sql")]
286    use crate::{
287        db::{
288            data::{AcceptedFixedUpdatePatch, FieldSlot},
289            executor::budget::{HardExecutionBudget, HardExecutionFailureHeadroom},
290            query::plan::expr::{BinaryOp, Expr, FieldId},
291        },
292        types::Timestamp,
293        value::Value,
294    };
295
296    struct TestCanister;
297
298    impl Path for TestCanister {
299        const PATH: &'static str = "db::session::mutation_job::tests::Canister";
300    }
301
302    impl CanisterKind for TestCanister {
303        const COMMIT_MEMORY_ID: u8 = 244;
304        const COMMIT_STABLE_KEY: &'static str = "icydb.test.mutation_job.commit.v1";
305        const STARTUP_MEMORY_ID: u8 = 246;
306        const STARTUP_STABLE_KEY: &'static str = "icydb.test.mutation_job.startup.control.v1";
307        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 245;
308        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str = "icydb.test.mutation_job.progress.v1";
309    }
310
311    thread_local! {
312        static STORE_REGISTRY: StoreRegistry = StoreRegistry::new();
313    }
314
315    fn job_id() -> MutationJobId {
316        MutationJobId::try_from_bytes([19; 32]).expect("nonzero mutation job id should admit")
317    }
318
319    fn session() -> DbSession<TestCanister> {
320        let root = RequestExecutionRoot::__new_runtime_root();
321        DbSession::new(&STORE_REGISTRY, &root)
322    }
323
324    #[cfg(feature = "sql")]
325    fn exhausted_session() -> DbSession<TestCanister> {
326        let budget =
327            HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
328        let root = RequestExecutionRoot::new_for_tests(budget);
329        DbSession::new(&STORE_REGISTRY, &root)
330    }
331
332    #[test]
333    fn session_load_and_terminal_acknowledgement_preserve_the_store_contract() {
334        let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
335            .expect("bounded initial record should admit");
336        with_mutation_progress_store::<TestCanister, _>(|store| {
337            store.insert_mutation(&initial).map(|_| ())
338        })
339        .expect("initial record should insert");
340
341        let session = session();
342        assert_eq!(
343            session.mutation_job_state(job_id()),
344            Ok(initial.state().clone())
345        );
346        assert_eq!(
347            session.acknowledge_mutation_job(job_id(), 0),
348            Err(MutationJobError::Active),
349        );
350
351        let request = MutationJobAdvanceRequest::new(
352            job_id(),
353            0,
354            MutationJobIdempotencyKey::new("authority-drift")
355                .expect("bounded idempotency key should admit"),
356        );
357        let (terminal, _) = initial
358            .apply_transition(
359                &request,
360                MutationJobTransition::new(
361                    MutationJobStatus::RestartRequired(
362                        MutationJobRestartReason::ManagedTimestampRegression,
363                    ),
364                    MutationJobPhase::Forward,
365                    Vec::new(),
366                    0,
367                    0,
368                    0,
369                ),
370            )
371            .expect("terminal restart receipt should admit");
372        with_mutation_progress_store::<TestCanister, _>(|store| store.replace_mutation(&terminal))
373            .expect("terminal record should replace active state");
374
375        assert_eq!(
376            session.acknowledge_mutation_job(job_id(), 0),
377            Err(MutationJobError::StaleSequence {
378                expected: 0,
379                actual: 1,
380            }),
381        );
382        assert_eq!(session.acknowledge_mutation_job(job_id(), 1), Ok(()));
383        assert_eq!(session.acknowledge_mutation_job(job_id(), 1), Ok(()));
384        assert_eq!(
385            session.mutation_job_state(job_id()),
386            Err(MutationJobError::NotFound),
387        );
388    }
389
390    #[cfg(feature = "sql")]
391    fn canonical_intent(
392        authority: u8,
393        scope_value: u64,
394        timestamp: i64,
395    ) -> CanonicalMutationIntent {
396        let scope = Expr::Binary {
397            op: BinaryOp::Eq,
398            left: Box::new(Expr::Field(FieldId::new("collection_id"))),
399            right: Box::new(Expr::Literal(Value::Nat64(scope_value))),
400        };
401        let patch = AcceptedFixedUpdatePatch::from_canonical_fields(vec![(
402            FieldSlot::from_validated_index(1),
403            vec![3, 4, 5],
404        )])
405        .expect("fixed patch should admit");
406        CanonicalMutationIntent::new(
407            [authority; 16],
408            [authority; 32],
409            "journaled".to_string(),
410            "schema::Token".to_string(),
411            7,
412            11,
413            1,
414            [authority; 16],
415            &scope,
416            &patch,
417            Timestamp::from_millis(timestamp),
418            17,
419        )
420        .expect("canonical intent should admit")
421    }
422
423    #[cfg(feature = "sql")]
424    #[test]
425    fn occupied_start_distinguishes_replay_authority_drift_and_identity_conflict() {
426        let retained_intent = canonical_intent(1, 7, 100);
427        let retained = MutationJobRecord::new(
428            job_id(),
429            retained_intent.encode().expect("intent should encode"),
430            vec![9],
431        )
432        .expect("record should admit");
433
434        assert_eq!(
435            resolve_occupied_mutation_job_start(&retained, &canonical_intent(1, 7, 200)),
436            Ok(retained.state().clone()),
437        );
438        assert_eq!(
439            resolve_occupied_mutation_job_start(&retained, &canonical_intent(2, 7, 200)),
440            Err(MutationJobError::AuthorityMismatch),
441        );
442        assert_eq!(
443            resolve_occupied_mutation_job_start(&retained, &canonical_intent(1, 8, 200)),
444            Err(MutationJobError::IdentityConflict),
445        );
446    }
447
448    #[cfg(feature = "sql")]
449    #[test]
450    fn aggregate_budget_exhaustion_does_not_advance_durable_state() {
451        let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
452            .expect("bounded initial record should admit");
453        with_mutation_progress_store::<TestCanister, _>(|store| {
454            store.insert_mutation(&initial).map(|_| ())
455        })
456        .expect("initial record should insert");
457        let request = MutationJobAdvanceRequest::new(
458            job_id(),
459            0,
460            MutationJobIdempotencyKey::new("budget-exhausted")
461                .expect("bounded idempotency key should admit"),
462        );
463
464        assert!(matches!(
465            exhausted_session().advance_trusted_mutation_job(&request),
466            Err(MutationJobError::ExecutionBudgetExceeded {
467                limit: 0,
468                observed: 1,
469                ..
470            })
471        ));
472        assert_eq!(
473            session().mutation_job_state(job_id()),
474            Ok(initial.state().clone()),
475        );
476    }
477
478    #[cfg(feature = "sql")]
479    #[test]
480    fn exact_replay_precedes_advance_budget_accounting() {
481        let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
482            .expect("bounded initial record should admit");
483        let request = MutationJobAdvanceRequest::new(
484            job_id(),
485            0,
486            MutationJobIdempotencyKey::new("lost-response")
487                .expect("bounded idempotency key should admit"),
488        );
489        let (advanced, receipt) = initial
490            .apply_transition(
491                &request,
492                MutationJobTransition::new(
493                    MutationJobStatus::Active,
494                    MutationJobPhase::Forward,
495                    vec![4],
496                    1,
497                    1,
498                    0,
499                ),
500            )
501            .expect("bounded successor should admit");
502        with_mutation_progress_store::<TestCanister, _>(|store| {
503            store.insert_mutation(&advanced).map(|_| ())
504        })
505        .expect("advanced record should insert");
506
507        assert_eq!(
508            exhausted_session().advance_trusted_mutation_job(&request),
509            Ok(receipt),
510        );
511        assert_eq!(
512            session().mutation_job_state(job_id()),
513            Ok(advanced.state().clone()),
514        );
515    }
516}