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 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}