1use 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 #[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 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 #[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 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 #[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 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}