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