1use 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 #[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 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 #[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 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 #[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 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 const COMMIT_MEMORY_ID: u8 = 244;
250 const COMMIT_STABLE_KEY: &'static str = "icydb.test.mutation_job.commit.v1";
251 const STARTUP_MEMORY_ID: u8 = 246;
252 const STARTUP_STABLE_KEY: &'static str = "icydb.test.mutation_job.startup.control.v1";
253 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 245;
254 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str = "icydb.test.mutation_job.progress.v1";
255 }
256
257 thread_local! {
258 static STORE_REGISTRY: StoreRegistry = StoreRegistry::new();
259 }
260
261 fn job_id() -> MutationJobId {
262 MutationJobId::try_from_bytes([19; 32]).expect("nonzero mutation job id should admit")
263 }
264
265 fn session() -> DbSession<TestCanister> {
266 let root = RequestExecutionRoot::__new_runtime_root();
267 DbSession::new(&STORE_REGISTRY, &root)
268 }
269
270 #[cfg(feature = "sql")]
271 fn exhausted_session() -> DbSession<TestCanister> {
272 let budget =
273 HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
274 let root = RequestExecutionRoot::new_for_tests(budget);
275 DbSession::new(&STORE_REGISTRY, &root)
276 }
277
278 #[test]
279 fn session_load_and_terminal_acknowledgement_preserve_the_store_contract() {
280 let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
281 .expect("bounded initial record should admit");
282 with_mutation_progress_store::<TestCanister, _>(|store| {
283 store.insert_mutation(&initial).map(|_| ())
284 })
285 .expect("initial record should insert");
286
287 let session = session();
288 assert_eq!(
289 session.mutation_job_state(job_id()),
290 Ok(initial.state().clone())
291 );
292 assert_eq!(
293 session.acknowledge_mutation_job(job_id(), 0),
294 Err(MutationJobError::Active),
295 );
296
297 let request = MutationJobAdvanceRequest::new(
298 job_id(),
299 0,
300 MutationJobIdempotencyKey::new("authority-drift")
301 .expect("bounded idempotency key should admit"),
302 );
303 let (terminal, _) = initial
304 .apply_transition(
305 &request,
306 MutationJobTransition::new(
307 MutationJobStatus::RestartRequired(
308 MutationJobRestartReason::ManagedTimestampRegression,
309 ),
310 MutationJobPhase::Forward,
311 Vec::new(),
312 0,
313 0,
314 0,
315 ),
316 )
317 .expect("terminal restart receipt should admit");
318 with_mutation_progress_store::<TestCanister, _>(|store| store.replace_mutation(&terminal))
319 .expect("terminal record should replace active state");
320
321 assert_eq!(
322 session.acknowledge_mutation_job(job_id(), 0),
323 Err(MutationJobError::StaleSequence {
324 expected: 0,
325 actual: 1,
326 }),
327 );
328 assert_eq!(session.acknowledge_mutation_job(job_id(), 1), Ok(()));
329 assert_eq!(session.acknowledge_mutation_job(job_id(), 1), Ok(()));
330 assert_eq!(
331 session.mutation_job_state(job_id()),
332 Err(MutationJobError::NotFound),
333 );
334 }
335
336 #[cfg(feature = "sql")]
337 fn canonical_intent(
338 authority: u8,
339 scope_value: u64,
340 timestamp: i64,
341 ) -> CanonicalMutationIntent {
342 let scope = Expr::Binary {
343 op: BinaryOp::Eq,
344 left: Box::new(Expr::Field(FieldId::new("collection_id"))),
345 right: Box::new(Expr::Literal(Value::Nat64(scope_value))),
346 };
347 let patch = AcceptedFixedUpdatePatch::from_canonical_fields(vec![(
348 FieldSlot::from_validated_index(1),
349 vec![3, 4, 5],
350 )])
351 .expect("fixed patch should admit");
352 CanonicalMutationIntent::new(
353 [authority; 16],
354 [authority; 32],
355 "journaled".to_string(),
356 "schema::Token".to_string(),
357 7,
358 11,
359 1,
360 [authority; 16],
361 &scope,
362 &patch,
363 Timestamp::from_millis(timestamp),
364 17,
365 )
366 .expect("canonical intent should admit")
367 }
368
369 #[cfg(feature = "sql")]
370 #[test]
371 fn occupied_start_distinguishes_replay_authority_drift_and_identity_conflict() {
372 let retained_intent = canonical_intent(1, 7, 100);
373 let retained = MutationJobRecord::new(
374 job_id(),
375 retained_intent.encode().expect("intent should encode"),
376 vec![9],
377 )
378 .expect("record should admit");
379
380 assert_eq!(
381 resolve_occupied_mutation_job_start(&retained, &canonical_intent(1, 7, 200)),
382 Ok(retained.state().clone()),
383 );
384 assert_eq!(
385 resolve_occupied_mutation_job_start(&retained, &canonical_intent(2, 7, 200)),
386 Err(MutationJobError::AuthorityMismatch),
387 );
388 assert_eq!(
389 resolve_occupied_mutation_job_start(&retained, &canonical_intent(1, 8, 200)),
390 Err(MutationJobError::IdentityConflict),
391 );
392 }
393
394 #[cfg(feature = "sql")]
395 #[test]
396 fn aggregate_budget_exhaustion_does_not_advance_durable_state() {
397 let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
398 .expect("bounded initial record should admit");
399 with_mutation_progress_store::<TestCanister, _>(|store| {
400 store.insert_mutation(&initial).map(|_| ())
401 })
402 .expect("initial record should insert");
403 let request = MutationJobAdvanceRequest::new(
404 job_id(),
405 0,
406 MutationJobIdempotencyKey::new("budget-exhausted")
407 .expect("bounded idempotency key should admit"),
408 );
409
410 assert!(matches!(
411 exhausted_session().advance_trusted_mutation_job(&request),
412 Err(MutationJobError::ExecutionBudgetExceeded {
413 limit: 0,
414 observed: 1,
415 ..
416 })
417 ));
418 assert_eq!(
419 session().mutation_job_state(job_id()),
420 Ok(initial.state().clone()),
421 );
422 }
423
424 #[cfg(feature = "sql")]
425 #[test]
426 fn exact_replay_precedes_advance_budget_accounting() {
427 let initial = MutationJobRecord::new(job_id(), vec![1, 2], vec![3])
428 .expect("bounded initial record should admit");
429 let request = MutationJobAdvanceRequest::new(
430 job_id(),
431 0,
432 MutationJobIdempotencyKey::new("lost-response")
433 .expect("bounded idempotency key should admit"),
434 );
435 let (advanced, receipt) = initial
436 .apply_transition(
437 &request,
438 MutationJobTransition::new(
439 MutationJobStatus::Active,
440 MutationJobPhase::Forward,
441 vec![4],
442 1,
443 1,
444 0,
445 ),
446 )
447 .expect("bounded successor should admit");
448 with_mutation_progress_store::<TestCanister, _>(|store| {
449 store.insert_mutation(&advanced).map(|_| ())
450 })
451 .expect("advanced record should insert");
452
453 assert_eq!(
454 exhausted_session().advance_trusted_mutation_job(&request),
455 Ok(receipt),
456 );
457 assert_eq!(
458 session().mutation_job_state(job_id()),
459 Ok(advanced.state().clone()),
460 );
461 }
462}