1use super::AcceptedSchemaCatalogContext;
8use crate::{
9 db::{
10 DbSession, DynamicMutation, DynamicMutationResult, DynamicStructuralPatch,
11 DynamicTypedBindingError, DynamicTypedEntityBinding, DynamicTypedFieldBindingRequest,
12 DynamicTypedFieldType, DynamicTypedMutation, DynamicTypedStructuralPatch, DynamicWriteCell,
13 commit::{CommitRowOp, database_incarnation_id},
14 data::{
15 AcceptedMutationIntentPatch, AcceptedPreKeyInsert, DecodedDataStoreKey, FieldSlot,
16 RawRow, StructuralRowContract, StructuralSlotReader,
17 canonical_row_from_raw_row_with_accepted_decode_contract,
18 resolve_existing_replace_structural_patch_with_accepted_contract,
19 resolve_insert_structural_patch_with_accepted_contract,
20 resolve_update_structural_patch_with_accepted_contract,
21 },
22 executor::{
23 AcceptedMutationConstraintScheduler, commit_structural_row_ops_with_window_for_path,
24 mutation_key_exists_error,
25 },
26 schema::{
27 AcceptedFieldKind, AcceptedIdentityAllocation, AcceptedRowLayoutRuntimeContract,
28 FieldId, FieldInsertGeneration, IdentityStatementCursor, lower_field_type,
29 output_value_from_runtime,
30 },
31 write_context::{AcceptedWriteContext, MutationMode},
32 },
33 error::{InternalError, MutationDiagnosticContext},
34 metrics::sink::{MetricsEvent, SaveMutationKind, record},
35 traits::CanisterKind,
36 types::{CurrentTimestamp, Timestamp},
37 value::{InputValue, Value},
38};
39use icydb_schema::{EntitySourceKey, FieldSourceKey, FieldType, TypeSourceKey};
40
41#[derive(Clone, Debug, Eq, PartialEq)]
42struct AcceptedIdentityInsertField {
43 field_id: FieldId,
44 field_slot: usize,
45 accepted_kind: AcceptedFieldKind,
46}
47
48pub(in crate::db::session) enum AcceptedStructuralMutationTarget {
51 ResolveFromAfterImage,
52 Expected(Box<DecodedDataStoreKey>),
53}
54
55impl AcceptedStructuralMutationTarget {
56 pub(in crate::db::session) fn expected(key: DecodedDataStoreKey) -> Self {
57 Self::Expected(Box::new(key))
58 }
59}
60
61pub(in crate::db::session) enum AcceptedStructuralMutation {
64 Save {
65 mode: MutationMode,
66 target: AcceptedStructuralMutationTarget,
67 patch: AcceptedMutationIntentPatch,
68 },
69 Delete {
70 key: Box<DecodedDataStoreKey>,
71 },
72}
73
74impl AcceptedStructuralMutation {
75 pub(in crate::db::session) const fn save(
76 mode: MutationMode,
77 target: AcceptedStructuralMutationTarget,
78 patch: AcceptedMutationIntentPatch,
79 ) -> Self {
80 Self::Save {
81 mode,
82 target,
83 patch,
84 }
85 }
86
87 pub(in crate::db::session) fn delete(key: DecodedDataStoreKey) -> Self {
88 Self::Delete { key: Box::new(key) }
89 }
90}
91
92const MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS: usize = 4_096;
93const MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES: usize = 16 * 1024 * 1024;
94const MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES: usize = 1024 * 1024;
95
96fn add_structural_mutation_staged_bytes(
97 total: &mut usize,
98 lengths: impl IntoIterator<Item = usize>,
99) -> Result<(), InternalError> {
100 for length in lengths {
101 *total = total.checked_add(length).ok_or_else(|| {
102 InternalError::mutation_batch_staged_bytes_exceeded(
103 None,
104 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
105 )
106 })?;
107 if *total > MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES {
108 return Err(InternalError::mutation_batch_staged_bytes_exceeded(
109 Some(*total),
110 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
111 ));
112 }
113 }
114 Ok(())
115}
116
117fn validate_structural_mutation_result_bytes(encoded_bytes: usize) -> Result<(), InternalError> {
118 if encoded_bytes > MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES {
119 return Err(InternalError::mutation_batch_result_bytes_exceeded(
120 encoded_bytes,
121 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES,
122 ));
123 }
124 Ok(())
125}
126
127pub(in crate::db::session) struct AcceptedStructuralMutationRow {
129 values: Vec<Value>,
130 logical_changed: bool,
131}
132
133impl AcceptedStructuralMutationRow {
134 #[cfg(feature = "sql")]
135 pub(in crate::db::session) fn into_values(self) -> Vec<Value> {
136 self.values
137 }
138
139 pub(in crate::db::session) const fn logical_changed(&self) -> bool {
140 self.logical_changed
141 }
142}
143
144const fn dynamic_mutation_mode(request: &DynamicMutation) -> Option<MutationMode> {
145 match request {
146 DynamicMutation::Insert { .. } => Some(MutationMode::Insert),
147 DynamicMutation::Update { .. } => Some(MutationMode::Update),
148 DynamicMutation::Replace { .. } => Some(MutationMode::Replace),
149 DynamicMutation::Delete { .. } => None,
150 }
151}
152
153const fn dynamic_typed_mutation_mode(request: &DynamicTypedMutation) -> MutationMode {
154 match request {
155 DynamicTypedMutation::Insert { .. } => MutationMode::Insert,
156 DynamicTypedMutation::Update { .. } => MutationMode::Update,
157 DynamicTypedMutation::Replace { .. } => MutationMode::Replace,
158 }
159}
160
161const fn diagnostic_mutation_operation(
162 mode: MutationMode,
163) -> icydb_diagnostic_code::DiagnosticMutationOperation {
164 match mode {
165 MutationMode::Insert => icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
166 MutationMode::Replace => icydb_diagnostic_code::DiagnosticMutationOperation::Replace,
167 MutationMode::Update => icydb_diagnostic_code::DiagnosticMutationOperation::Update,
168 }
169}
170
171const fn mutation_diagnostic_context(
172 entity_tag: crate::types::EntityTag,
173 mode: MutationMode,
174 batch_position: u32,
175) -> MutationDiagnosticContext {
176 MutationDiagnosticContext::new(
177 entity_tag.value(),
178 diagnostic_mutation_operation(mode),
179 batch_position,
180 )
181}
182
183const fn dynamic_write_context(operation_timestamp: Timestamp) -> AcceptedWriteContext {
184 AcceptedWriteContext::new(operation_timestamp)
185}
186
187fn insert_key_exists_after_generation(identity_generated: bool) -> InternalError {
188 if identity_generated {
189 InternalError::identity_state_corruption()
190 } else {
191 mutation_key_exists_error()
192 }
193}
194
195fn dynamic_key(
196 entity_tag: crate::types::EntityTag,
197 key: &InputValue,
198) -> Result<DecodedDataStoreKey, InternalError> {
199 let value = key
200 .clone()
201 .try_into_runtime_non_enum()
202 .ok_or_else(InternalError::executor_unsupported)?;
203 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
204}
205
206fn lower_dynamic_patch(
207 entity_path: &str,
208 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
209 patch: &DynamicStructuralPatch,
210 mode: MutationMode,
211 mutation_context: MutationDiagnosticContext,
212) -> Result<AcceptedMutationIntentPatch, InternalError> {
213 let mut lowered = AcceptedMutationIntentPatch::new();
214 for (field_name, cell) in patch.fields() {
215 let slot = descriptor
216 .field_slot_index_by_name(field_name)
217 .ok_or_else(|| {
218 InternalError::mutation_structural_field_unknown(entity_path, field_name)
219 })?;
220 let field = descriptor
221 .field_for_slot_index(slot)
222 .ok_or_else(InternalError::executor_invariant)?;
223 if !matches!(cell, DynamicWriteCell::Omitted)
224 && (field.write_policy().insert_generation().is_some()
225 || field.write_policy().write_management().is_some())
226 {
227 return Err(InternalError::mutation_database_owned_field_explicit(
228 mutation_context,
229 field.field_id().get(),
230 ));
231 }
232 let slot = FieldSlot::from_validated_index(slot);
233 lowered = match cell {
234 DynamicWriteCell::Omitted => lowered,
235 DynamicWriteCell::Default => match mode {
236 MutationMode::Insert | MutationMode::Replace => {
237 lowered.set_explicit_insert_default(slot)
238 }
239 MutationMode::Update => lowered.set_explicit_update_default(slot),
240 },
241 DynamicWriteCell::Null => lowered.set_authored(slot, InputValue::Null),
242 DynamicWriteCell::Value(value) => lowered.set_authored(slot, value.clone()),
243 };
244 }
245 Ok(lowered)
246}
247
248fn lower_dynamic_mutation_intent(
249 entity_tag: crate::types::EntityTag,
250 entity_path: &str,
251 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
252 request: &DynamicMutation,
253 batch_position: u32,
254) -> Result<(AcceptedStructuralMutation, Option<SaveMutationKind>), InternalError> {
255 match request {
256 DynamicMutation::Insert { patch, .. } => {
257 let mode = MutationMode::Insert;
258 Ok((
259 AcceptedStructuralMutation::save(
260 mode,
261 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
262 lower_dynamic_patch(
263 entity_path,
264 descriptor,
265 patch,
266 mode,
267 mutation_diagnostic_context(entity_tag, mode, batch_position),
268 )?,
269 ),
270 Some(SaveMutationKind::Insert),
271 ))
272 }
273 DynamicMutation::Update { key, patch, .. }
274 | DynamicMutation::Replace { key, patch, .. } => {
275 let mode =
276 dynamic_mutation_mode(request).ok_or_else(InternalError::executor_invariant)?;
277 let kind = match mode {
278 MutationMode::Insert => SaveMutationKind::Insert,
279 MutationMode::Replace => SaveMutationKind::Replace,
280 MutationMode::Update => SaveMutationKind::Update,
281 };
282 Ok((
283 AcceptedStructuralMutation::save(
284 mode,
285 AcceptedStructuralMutationTarget::expected(dynamic_key(entity_tag, key)?),
286 lower_dynamic_patch(
287 entity_path,
288 descriptor,
289 patch,
290 mode,
291 mutation_diagnostic_context(entity_tag, mode, batch_position),
292 )?,
293 ),
294 Some(kind),
295 ))
296 }
297 DynamicMutation::Delete { key, .. } => Ok((
298 AcceptedStructuralMutation::delete(dynamic_key(entity_tag, key)?),
299 None,
300 )),
301 }
302}
303
304fn lower_typed_patch(
305 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
306 patch: &DynamicTypedStructuralPatch,
307 mode: MutationMode,
308 mutation_context: MutationDiagnosticContext,
309) -> Result<AcceptedMutationIntentPatch, InternalError> {
310 let mut lowered = AcceptedMutationIntentPatch::new();
311 for (field_id, slot, cell) in patch.fields() {
312 let slot_index = usize::from(*slot);
313 let field = descriptor
314 .field_for_slot_index(slot_index)
315 .ok_or_else(InternalError::store_invariant)?;
316 if field.field_id().get() != *field_id {
317 return Err(InternalError::store_invariant());
318 }
319 if !matches!(cell, DynamicWriteCell::Omitted)
320 && (field.write_policy().insert_generation().is_some()
321 || field.write_policy().write_management().is_some())
322 {
323 return Err(InternalError::mutation_database_owned_field_explicit(
324 mutation_context,
325 field.field_id().get(),
326 ));
327 }
328 let slot = FieldSlot::from_validated_index(slot_index);
329 lowered = match cell {
330 DynamicWriteCell::Omitted => lowered,
331 DynamicWriteCell::Default => match mode {
332 MutationMode::Insert | MutationMode::Replace => {
333 lowered.set_explicit_insert_default(slot)
334 }
335 MutationMode::Update => lowered.set_explicit_update_default(slot),
336 },
337 DynamicWriteCell::Null => lowered.set_authored(slot, InputValue::Null),
338 DynamicWriteCell::Value(value) => lowered.set_authored(slot, value.clone()),
339 };
340 }
341 Ok(lowered)
342}
343
344fn preserve_dynamic_replacement_identity(
345 key: &DecodedDataStoreKey,
346 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
347 mut patch: AcceptedMutationIntentPatch,
348) -> Result<AcceptedMutationIntentPatch, InternalError> {
349 let primary_key_slots = descriptor.primary_key_slot_indices();
350 let runtime_key = key.primary_key_runtime_value();
351 let components = match runtime_key {
352 Value::List(values) if primary_key_slots.len() > 1 => values,
353 value if primary_key_slots.len() == 1 => vec![value],
354 _ => return Err(InternalError::executor_invariant()),
355 };
356 if components.len() != primary_key_slots.len() {
357 return Err(InternalError::executor_invariant());
358 }
359
360 for (slot, value) in primary_key_slots.iter().copied().zip(components) {
361 let _ = descriptor
362 .field_for_slot_index(slot)
363 .ok_or_else(InternalError::executor_invariant)?;
364 let has_explicit_intent = patch
365 .entries()
366 .iter()
367 .any(|entry| entry.slot().index() == slot);
368 if has_explicit_intent {
369 continue;
370 }
371 let value = InputValue::try_from_runtime_non_enum(&value)
372 .ok_or_else(InternalError::executor_invariant)?;
373 patch =
374 patch.set_preserved_replacement_identity(FieldSlot::from_validated_index(slot), value);
375 }
376
377 Ok(patch)
378}
379
380fn accepted_identity_insert_field(
384 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
385) -> Result<Option<AcceptedIdentityInsertField>, InternalError> {
386 let mut identity = None;
387 for field in descriptor.fields() {
388 if field.write_policy().insert_generation() != Some(FieldInsertGeneration::Identity) {
389 continue;
390 }
391 let field_slot = usize::from(field.slot().get());
392 if identity
393 .replace(AcceptedIdentityInsertField {
394 field_id: field.field_id(),
395 field_slot,
396 accepted_kind: field.kind().clone(),
397 })
398 .is_some()
399 || descriptor.primary_key_slot_indices() != [field_slot]
400 {
401 return Err(InternalError::identity_corruption());
402 }
403 }
404 Ok(identity)
405}
406
407fn checked_pre_key_candidate_count(count: usize) -> Result<u32, InternalError> {
408 u32::try_from(count).map_err(|_| InternalError::identity_candidate_count_exhausted())
409}
410
411fn validate_identity_materialization(
412 entity_tag: crate::types::EntityTag,
413 identity_field: &AcceptedIdentityInsertField,
414 candidate: &AcceptedPreKeyInsert,
415 allocation: &AcceptedIdentityAllocation,
416 data_key: &DecodedDataStoreKey,
417 reader: &StructuralSlotReader<'_>,
418) -> Result<(), InternalError> {
419 let owner = allocation.owner();
420 let slot_value = reader.required_cached_value(identity_field.field_slot)?;
421 if candidate.entity_tag() != entity_tag
422 || candidate.input_ordinal() != allocation.input_ordinal()
423 || owner.entity_tag() != entity_tag
424 || owner.field_id() != identity_field.field_id
425 || allocation.field_slot() != identity_field.field_slot
426 || slot_value != allocation.value()
427 || data_key.primary_key_runtime_value() != *allocation.value()
428 {
429 return Err(InternalError::identity_corruption());
430 }
431 Ok(())
432}
433
434fn data_key_from_row(
435 entity_tag: crate::types::EntityTag,
436 contract: &StructuralRowContract,
437 row: &RawRow,
438) -> Result<DecodedDataStoreKey, InternalError> {
439 let reader =
440 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(row, contract)?;
441 let values = contract
442 .primary_key_slot_indices()
443 .iter()
444 .map(|slot| reader.required_cached_value(*slot).cloned())
445 .collect::<Result<Vec<_>, _>>()?;
446 let value = match values.as_slice() {
447 [value] => value.clone(),
448 _ => Value::List(values),
449 };
450 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
451}
452
453#[cfg(feature = "sql")]
454pub(in crate::db::session) fn structural_data_key_from_runtime_values(
455 entity_tag: crate::types::EntityTag,
456 values: Vec<Value>,
457) -> Result<DecodedDataStoreKey, InternalError> {
458 let value = match values.as_slice() {
459 [value] => value.clone(),
460 _ => Value::List(values),
461 };
462 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
463}
464
465fn validated_existing_row(
466 store: crate::db::registry::StoreHandle,
467 data_key: &DecodedDataStoreKey,
468 contract: &StructuralRowContract,
469) -> Result<Option<RawRow>, InternalError> {
470 let raw_key = data_key.to_raw()?;
471 let row = store.with_data(|data| data.get(&raw_key));
472 if let Some(row) = row.as_ref() {
473 let reader =
474 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(row, contract)?;
475 reader.validate_primary_key(data_key)?;
476 }
477 Ok(row)
478}
479
480fn prepare_dynamic_mutation_result(
481 catalog: &AcceptedSchemaCatalogContext,
482 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
483 rows: Vec<AcceptedStructuralMutationRow>,
484 enforce_mixed_batch_result_bound: bool,
485) -> Result<DynamicMutationResult, InternalError> {
486 let affected_rows = rows.iter().try_fold(0_u32, |total, row| {
487 total
488 .checked_add(u32::from(row.logical_changed()))
489 .ok_or_else(InternalError::executor_invariant)
490 })?;
491 let columns = descriptor
492 .fields()
493 .iter()
494 .map(|field| field.name().to_string())
495 .collect();
496 let rows = rows
497 .into_iter()
498 .map(|row| {
499 row.values
500 .iter()
501 .map(|value| {
502 output_value_from_runtime(catalog.enum_catalog(), value)
503 .map_err(|_| InternalError::store_invariant())
504 })
505 .collect::<Result<Vec<_>, _>>()
506 })
507 .collect::<Result<Vec<_>, _>>()?;
508 let result = DynamicMutationResult {
509 entity: catalog.snapshot().entity_name().to_string(),
510 columns,
511 rows,
512 affected_rows,
513 };
514 if enforce_mixed_batch_result_bound {
515 let encoded =
516 candid::encode_one(&result).map_err(|_| InternalError::executor_invariant())?;
517 validate_structural_mutation_result_bytes(encoded.len())?;
518 }
519 Ok(result)
520}
521
522fn dynamic_typed_field_type(
523 field_type: DynamicTypedFieldType,
524) -> Result<FieldType, DynamicTypedBindingError> {
525 match field_type {
526 DynamicTypedFieldType::Scalar(scalar) => Ok(FieldType::Scalar(scalar)),
527 DynamicTypedFieldType::List(item) => {
528 Ok(FieldType::List(Box::new(dynamic_typed_field_type(*item)?)))
529 }
530 DynamicTypedFieldType::Named(source_key) => TypeSourceKey::try_new(source_key)
531 .map(FieldType::Named)
532 .map_err(|_| DynamicTypedBindingError::FieldUnavailable),
533 }
534}
535
536fn typed_adapter_field_kind_matches(
537 accepted: &AcceptedFieldKind,
538 expected: &AcceptedFieldKind,
539) -> bool {
540 if accepted == expected {
541 return true;
542 }
543 match (accepted, expected) {
544 (AcceptedFieldKind::Relation { key_kind, .. }, expected) => {
545 typed_adapter_field_kind_matches(key_kind, expected)
546 }
547 (AcceptedFieldKind::List(accepted), AcceptedFieldKind::List(expected)) => {
548 typed_adapter_field_kind_matches(accepted, expected)
549 }
550 _ => false,
551 }
552}
553
554impl<C: CanisterKind> DbSession<C> {
555 pub fn issue_typed_entity_binding(
557 &self,
558 entity_source_key: &str,
559 field_requests: &[DynamicTypedFieldBindingRequest],
560 ) -> Result<DynamicTypedEntityBinding, DynamicTypedBindingError> {
561 let entity_source = EntitySourceKey::try_new(entity_source_key)
562 .map_err(|_| DynamicTypedBindingError::FieldUnavailable)?;
563 let field_requests = field_requests
564 .iter()
565 .map(|request| {
566 Ok((
567 FieldSourceKey::try_new(request.source_key.clone())
568 .map_err(|_| DynamicTypedBindingError::FieldUnavailable)?,
569 dynamic_typed_field_type(request.field_type.clone())?,
570 request.nullable,
571 ))
572 })
573 .collect::<Result<Vec<_>, DynamicTypedBindingError>>()?;
574 let catalog = self
575 .find_accepted_schema_catalog_context_for_entity_source_key(entity_source.as_str())?
576 .ok_or(DynamicTypedBindingError::FieldUnavailable)?;
577 let identity = catalog.identity();
578 if identity.entity_path() != entity_source.as_str() {
579 return Err(InternalError::store_invariant().into());
580 }
581 let store = self.db.recovered_store(identity.store_path())?;
582 let bundle = store
583 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
584 .ok_or_else(InternalError::store_invariant)?;
585 let entity_tag = identity.entity_tag();
586 if bundle.source_bindings().entity(&entity_source) != Some(entity_tag)
587 || bundle.revision() != catalog.revision()
588 {
589 return Err(InternalError::store_invariant().into());
590 }
591 let snapshot = bundle
592 .entity_snapshots()
593 .get(&entity_tag)
594 .ok_or_else(InternalError::store_invariant)?;
595 let row_contract = catalog.inspection_plan().row_contract();
596 let mut fields = Vec::with_capacity(field_requests.len());
597 for (source, field_type, nullable) in &field_requests {
598 let field_id = bundle
599 .source_bindings()
600 .field(entity_tag, source)
601 .ok_or(DynamicTypedBindingError::FieldUnavailable)?;
602 let field = snapshot
603 .fields()
604 .iter()
605 .find(|field| field.id() == field_id)
606 .ok_or_else(InternalError::store_invariant)?;
607 let runtime_field =
608 row_contract.required_accepted_field_contract(usize::from(field.slot().get()))?;
609 if runtime_field.field_id() != field_id {
610 return Err(InternalError::store_invariant().into());
611 }
612 let expected_kind = lower_field_type(field_type, bundle.source_bindings())
613 .map_err(|_| DynamicTypedBindingError::IncompatibleField)?;
614 if field.nullable() != *nullable
615 || !typed_adapter_field_kind_matches(field.kind(), &expected_kind)
616 {
617 return Err(DynamicTypedBindingError::IncompatibleField);
618 }
619 fields.push((
620 source.as_str().to_string(),
621 field_id.get(),
622 field.slot().get(),
623 field.name().to_string(),
624 ));
625 }
626 let adapter_names = bundle.typed_adapter_names()?;
627
628 DynamicTypedEntityBinding::new(
629 database_incarnation_id()?.to_bytes(),
630 entity_source.as_str().to_string(),
631 snapshot.entity_name().to_string(),
632 entity_tag.value(),
633 catalog.revision().get(),
634 catalog.fingerprint(),
635 row_contract.current_layout_version().get(),
636 fields,
637 adapter_names.named_types,
638 adapter_names.enum_variants,
639 adapter_names.composite_fields,
640 )
641 .map_err(Into::into)
642 }
643
644 pub(in crate::db::session) fn current_typed_entity_binding_catalog(
645 &self,
646 binding: &DynamicTypedEntityBinding,
647 ) -> Result<Option<AcceptedSchemaCatalogContext>, InternalError> {
648 if database_incarnation_id()?.to_bytes() != binding.database_incarnation {
649 return Ok(None);
650 }
651 let Some(catalog) = self.find_accepted_schema_catalog_context_for_entity_source_key(
652 binding.entity_source.as_str(),
653 )?
654 else {
655 return Ok(None);
656 };
657 let row_contract = catalog.inspection_plan().row_contract();
658 let identity = catalog.identity();
659 if identity.entity_path() != binding.entity_source.as_str()
660 || identity.entity_tag().value() != binding.entity_tag
661 || catalog.revision().get() != binding.accepted_revision
662 || catalog.fingerprint() != binding.accepted_fingerprint
663 || row_contract.current_layout_version().get() != binding.entity_generation
664 {
665 return Ok(None);
666 }
667 let entity_source = EntitySourceKey::try_new(binding.entity_source.clone())
668 .map_err(|_| InternalError::store_invariant())?;
669 let store = self.db.recovered_store(identity.store_path())?;
670 let bundle = store
671 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
672 .ok_or_else(InternalError::store_invariant)?;
673 if bundle.revision() != catalog.revision()
674 || bundle.source_bindings().entity(&entity_source) != Some(identity.entity_tag())
675 {
676 return Ok(None);
677 }
678 let snapshot = bundle
679 .entity_snapshots()
680 .get(&identity.entity_tag())
681 .ok_or_else(InternalError::store_invariant)?;
682 for (source_key, expected_field_id, expected_slot) in binding.field_identity_bindings() {
683 let source = FieldSourceKey::try_new(source_key)
684 .map_err(|_| InternalError::store_invariant())?;
685 let Some(field_id) = bundle
686 .source_bindings()
687 .field(identity.entity_tag(), &source)
688 else {
689 return Ok(None);
690 };
691 let Some(field) = snapshot
692 .fields()
693 .iter()
694 .find(|field| field.id() == field_id)
695 else {
696 return Err(InternalError::store_invariant());
697 };
698 if field_id.get() != expected_field_id || field.slot().get() != expected_slot {
699 return Ok(None);
700 }
701 }
702 Ok(Some(catalog))
703 }
704
705 pub fn typed_entity_binding_is_current(
707 &self,
708 binding: &DynamicTypedEntityBinding,
709 ) -> Result<bool, InternalError> {
710 self.current_typed_entity_binding_catalog(binding)
711 .map(|catalog| catalog.is_some())
712 }
713
714 #[cfg(feature = "sql")]
717 pub(in crate::db::session) fn execute_accepted_structural_delete_batch(
718 &self,
719 catalog: &AcceptedSchemaCatalogContext,
720 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
721 keys: Vec<DecodedDataStoreKey>,
722 precommit_validation: impl FnOnce(&[Vec<Value>]) -> Result<(), InternalError>,
723 ) -> Result<Vec<Vec<Value>>, InternalError> {
724 let mutations = keys
725 .into_iter()
726 .map(AcceptedStructuralMutation::delete)
727 .collect();
728 self.execute_accepted_structural_mutation_batch_inner(
729 catalog,
730 descriptor,
731 mutations,
732 Timestamp::now(),
733 false,
734 |rows| {
735 let rows = rows
736 .into_iter()
737 .map(AcceptedStructuralMutationRow::into_values)
738 .collect::<Vec<_>>();
739 precommit_validation(rows.as_slice())?;
740 Ok(rows)
741 },
742 )
743 }
744
745 pub(in crate::db::session) fn execute_accepted_structural_save_batch<T>(
753 &self,
754 catalog: &AcceptedSchemaCatalogContext,
755 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
756 mutations: Vec<AcceptedStructuralMutation>,
757 operation_timestamp: Timestamp,
758 precommit_preparation: impl FnOnce(
759 Vec<AcceptedStructuralMutationRow>,
760 ) -> Result<T, InternalError>,
761 ) -> Result<T, InternalError> {
762 self.execute_accepted_structural_mutation_batch_inner(
763 catalog,
764 descriptor,
765 mutations,
766 operation_timestamp,
767 false,
768 precommit_preparation,
769 )
770 }
771
772 #[cfg(feature = "sql")]
774 pub(in crate::db::session) fn execute_accepted_structural_update_prefix(
775 &self,
776 catalog: &AcceptedSchemaCatalogContext,
777 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
778 mutations: Vec<AcceptedStructuralMutation>,
779 operation_timestamp: Timestamp,
780 ) -> Result<usize, InternalError> {
781 self.execute_accepted_structural_mutation_batch_inner(
782 catalog,
783 descriptor,
784 mutations,
785 operation_timestamp,
786 true,
787 |rows| Ok(rows.len()),
788 )
789 }
790
791 #[expect(
792 clippy::too_many_lines,
793 reason = "one phased owner keeps accepted authority, mutation context, precommit preparation, output capture, and commit staging inseparable"
794 )]
795 fn execute_accepted_structural_mutation_batch_inner<T>(
796 &self,
797 catalog: &AcceptedSchemaCatalogContext,
798 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
799 mutations: Vec<AcceptedStructuralMutation>,
800 operation_timestamp: Timestamp,
801 largest_journaled_prefix: bool,
802 precommit_preparation: impl FnOnce(
803 Vec<AcceptedStructuralMutationRow>,
804 ) -> Result<T, InternalError>,
805 ) -> Result<T, InternalError> {
806 let identity = catalog.identity();
807 let entity_path = identity.entity_path();
808 let store_path = identity.store_path();
809 let row_decode_contract =
810 descriptor.row_decode_contract(catalog.value_catalog_handle().clone());
811 let row_contract = StructuralRowContract::from_accepted_decode_contract(
812 entity_path,
813 row_decode_contract.clone(),
814 );
815 let store = self.db.recovered_store(store_path)?;
816 let write_context = dynamic_write_context(operation_timestamp);
817 let identity_field = accepted_identity_insert_field(descriptor)?;
818 let identity_incarnation = identity_field
819 .as_ref()
820 .map(|_| database_incarnation_id())
821 .transpose()?;
822 let mutation_count = mutations.len();
823 if mutation_count > MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS {
824 return Err(InternalError::mutation_batch_too_many_items(
825 mutation_count,
826 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
827 ));
828 }
829 let identity_candidate_count = mutations
830 .iter()
831 .filter(|mutation| {
832 matches!(
833 mutation,
834 AcceptedStructuralMutation::Save {
835 mode: MutationMode::Insert,
836 target: AcceptedStructuralMutationTarget::ResolveFromAfterImage,
837 ..
838 }
839 )
840 })
841 .count();
842 let _ = checked_pre_key_candidate_count(identity_candidate_count)?;
843 let mut identity_cursor: Option<IdentityStatementCursor> = None;
844 let mut identity_insert_ordinal = 0_u32;
845 let mut scheduler = AcceptedMutationConstraintScheduler::new(
846 entity_path,
847 identity.entity_tag(),
848 row_decode_contract.clone(),
849 catalog.fingerprint(),
850 catalog.fingerprint_method_version(),
851 catalog.accepted_row_constraints(),
852 mutation_count,
853 );
854 let mut output = Vec::with_capacity(mutation_count);
855 let mut staged_bytes = 0_usize;
856
857 for (input_index, mutation) in mutations.into_iter().enumerate() {
858 let batch_input_ordinal = u32::try_from(input_index).map_err(|_| {
859 InternalError::mutation_batch_too_many_items(
860 mutation_count,
861 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
862 )
863 })?;
864 let AcceptedStructuralMutation::Save {
865 mode,
866 target,
867 patch: authored_patch,
868 } = mutation
869 else {
870 let AcceptedStructuralMutation::Delete { key } = mutation else {
871 return Err(InternalError::executor_invariant());
872 };
873 let before = validated_existing_row(store, &key, &row_contract)?
874 .ok_or_else(|| InternalError::store_not_found(&key))?;
875 let raw_key = key.to_raw()?;
876 let canonical_before = canonical_row_from_raw_row_with_accepted_decode_contract(
877 entity_path,
878 row_decode_contract.clone(),
879 &before,
880 )?;
881 add_structural_mutation_staged_bytes(
882 &mut staged_bytes,
883 [
884 raw_key.as_bytes().len(),
885 canonical_before.as_raw_row().as_bytes().len(),
886 ],
887 )?;
888 scheduler.schedule_delete(
889 CommitRowOp::new(
890 entity_path,
891 raw_key,
892 Some(canonical_before.as_raw_row().as_bytes().to_vec()),
893 None,
894 catalog.fingerprint(),
895 ),
896 batch_input_ordinal,
897 )?;
898 let reader = StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
899 canonical_before.as_raw_row(),
900 &row_contract,
901 )?;
902 let mut values = Vec::with_capacity(descriptor.fields().len());
903 for field in descriptor.fields() {
904 values.push(
905 reader
906 .required_cached_value(usize::from(field.slot().get()))?
907 .clone(),
908 );
909 }
910 output.push(AcceptedStructuralMutationRow {
911 values,
912 logical_changed: true,
913 });
914 continue;
915 };
916 let mutation_context =
917 mutation_diagnostic_context(identity.entity_tag(), mode, batch_input_ordinal);
918 let (expected_key, pre_key_insert, mut keyed_patch) = match target {
919 AcceptedStructuralMutationTarget::ResolveFromAfterImage => {
920 let candidate_ordinal =
921 if identity_field.is_some() && matches!(mode, MutationMode::Insert) {
922 identity_insert_ordinal
923 } else {
924 batch_input_ordinal
925 };
926 (
927 None,
928 Some(AcceptedPreKeyInsert::new(
929 identity.entity_tag(),
930 authored_patch,
931 candidate_ordinal,
932 )),
933 None,
934 )
935 }
936 AcceptedStructuralMutationTarget::Expected(key) => {
937 (Some(*key), None, Some(authored_patch))
938 }
939 };
940 if matches!(mode, MutationMode::Replace)
941 && let Some(key) = expected_key.as_ref()
942 {
943 let patch = keyed_patch
944 .take()
945 .ok_or_else(InternalError::executor_invariant)?;
946 keyed_patch = Some(preserve_dynamic_replacement_identity(
947 key, descriptor, patch,
948 )?);
949 }
950 let patch = pre_key_insert
951 .as_ref()
952 .map(AcceptedPreKeyInsert::fields)
953 .or(keyed_patch.as_ref())
954 .ok_or_else(InternalError::executor_invariant)?;
955 let before = expected_key
956 .as_ref()
957 .map(|key| validated_existing_row(store, key, &row_contract))
958 .transpose()?
959 .flatten();
960 match mode {
961 MutationMode::Insert if before.is_some() => {
962 return Err(mutation_key_exists_error());
963 }
964 MutationMode::Update if before.is_none() => {
965 let key = expected_key
966 .as_ref()
967 .ok_or_else(InternalError::executor_invariant)?;
968 return Err(InternalError::store_not_found(key));
969 }
970 MutationMode::Insert | MutationMode::Replace | MutationMode::Update => {}
971 }
972
973 let identity_allocation = if let Some(identity_field) = identity_field.as_ref()
974 && matches!(mode, MutationMode::Insert)
975 && before.is_none()
976 {
977 let candidate = pre_key_insert.as_ref().ok_or_else(|| {
978 InternalError::mutation_database_owned_field_explicit(
979 mutation_context,
980 identity_field.field_id.get(),
981 )
982 })?;
983 if identity_cursor.is_none() {
984 let incarnation = identity_incarnation
985 .ok_or_else(InternalError::identity_state_corruption)?;
986 identity_cursor = Some(store.with_schema(|schema_store| {
987 schema_store.identity_statement_cursor(
988 incarnation,
989 identity.entity_tag(),
990 identity_field.field_id,
991 &identity_field.accepted_kind,
992 )
993 })?);
994 }
995 let allocation = identity_cursor
996 .as_mut()
997 .ok_or_else(InternalError::identity_state_corruption)?
998 .allocate(identity_field.field_slot, candidate.input_ordinal())?;
999 identity_insert_ordinal = identity_insert_ordinal
1000 .checked_add(1)
1001 .ok_or_else(InternalError::identity_candidate_count_exhausted)?;
1002 Some(allocation)
1003 } else if let Some(identity_field) = identity_field.as_ref()
1004 && matches!(mode, MutationMode::Replace)
1005 && before.is_none()
1006 {
1007 return Err(InternalError::mutation_database_owned_field_explicit(
1008 mutation_context,
1009 identity_field.field_id.get(),
1010 ));
1011 } else {
1012 None
1013 };
1014
1015 let resolved = match (mode, before.as_ref()) {
1016 (MutationMode::Insert | MutationMode::Replace, None) => {
1017 resolve_insert_structural_patch_with_accepted_contract(
1018 entity_path,
1019 row_decode_contract.clone(),
1020 catalog.fingerprint(),
1021 catalog.accepted_row_constraints(),
1022 patch,
1023 write_context,
1024 mutation_context,
1025 identity_allocation.as_ref(),
1026 )?
1027 }
1028 (MutationMode::Update, Some(before)) => {
1029 resolve_update_structural_patch_with_accepted_contract(
1030 entity_path,
1031 row_decode_contract.clone(),
1032 catalog.fingerprint(),
1033 catalog.accepted_row_constraints(),
1034 before,
1035 patch,
1036 write_context,
1037 mutation_context,
1038 )?
1039 }
1040 (MutationMode::Replace, Some(before)) => {
1041 resolve_existing_replace_structural_patch_with_accepted_contract(
1042 entity_path,
1043 row_decode_contract.clone(),
1044 catalog.fingerprint(),
1045 catalog.accepted_row_constraints(),
1046 before,
1047 patch,
1048 write_context,
1049 mutation_context,
1050 )?
1051 }
1052 (MutationMode::Insert, Some(_)) | (MutationMode::Update, None) => {
1053 return Err(InternalError::executor_invariant());
1054 }
1055 };
1056 let (after, provenance) = resolved.into_parts();
1057 let reader = StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
1058 after.as_raw_row(),
1059 &row_contract,
1060 )?;
1061 let data_key = match expected_key {
1062 Some(key) => {
1063 reader.validate_primary_key(&key)?;
1064 key
1065 }
1066 None => {
1067 data_key_from_row(identity.entity_tag(), &row_contract, after.as_raw_row())?
1068 }
1069 };
1070 if let Some(allocation) = identity_allocation.as_ref() {
1071 validate_identity_materialization(
1072 identity.entity_tag(),
1073 identity_field
1074 .as_ref()
1075 .ok_or_else(InternalError::identity_corruption)?,
1076 pre_key_insert
1077 .as_ref()
1078 .ok_or_else(InternalError::identity_corruption)?,
1079 allocation,
1080 &data_key,
1081 &reader,
1082 )?;
1083 }
1084 if matches!(mode, MutationMode::Insert)
1085 && validated_existing_row(store, &data_key, &row_contract)?.is_some()
1086 {
1087 return Err(insert_key_exists_after_generation(
1088 identity_allocation.is_some(),
1089 ));
1090 }
1091 let raw_key = data_key.to_raw()?;
1092 let canonical_before = before
1093 .as_ref()
1094 .map(|before| {
1095 canonical_row_from_raw_row_with_accepted_decode_contract(
1096 entity_path,
1097 row_decode_contract.clone(),
1098 before,
1099 )
1100 })
1101 .transpose()?;
1102 let logical_changed = canonical_before.as_ref().is_none_or(|before| {
1103 before.as_raw_row().as_bytes() != after.as_raw_row().as_bytes()
1104 });
1105 let physical_changed = before
1106 .as_ref()
1107 .is_none_or(|before| before.as_bytes() != after.as_raw_row().as_bytes());
1108 add_structural_mutation_staged_bytes(
1109 &mut staged_bytes,
1110 [
1111 raw_key.as_bytes().len(),
1112 canonical_before
1113 .as_ref()
1114 .map_or(0, |before| before.as_raw_row().as_bytes().len()),
1115 after.as_raw_row().as_bytes().len(),
1116 ],
1117 )?;
1118 let row_op = physical_changed.then(|| {
1119 CommitRowOp::new(
1120 entity_path,
1121 raw_key.clone(),
1122 canonical_before
1123 .as_ref()
1124 .map(|before| before.as_raw_row().as_bytes().to_vec()),
1125 Some(after.as_raw_row().as_bytes().to_vec()),
1126 catalog.fingerprint(),
1127 )
1128 });
1129 scheduler.schedule_save_after_image(
1130 mode,
1131 &data_key,
1132 after.as_raw_row(),
1133 provenance.as_slice(),
1134 row_op,
1135 batch_input_ordinal,
1136 )?;
1137 if physical_changed {
1138 #[cfg(feature = "sql")]
1139 if largest_journaled_prefix
1140 && !crate::db::commit::journaled_row_ops_fit_commit_window(scheduler.rows())
1141 {
1142 scheduler.pop_last_save_row()?;
1143 if output.is_empty() {
1144 return Err(InternalError::query_sql_write_boundary(
1145 icydb_diagnostic_code::SqlWriteBoundaryCode::ResumableUpdateSingleRowResourceExceeded,
1146 ));
1147 }
1148 break;
1149 }
1150 }
1151
1152 let mut values = Vec::with_capacity(descriptor.fields().len());
1153 for field in descriptor.fields() {
1154 values.push(
1155 reader
1156 .required_cached_value(usize::from(field.slot().get()))?
1157 .clone(),
1158 );
1159 }
1160 output.push(AcceptedStructuralMutationRow {
1161 values,
1162 logical_changed,
1163 });
1164 }
1165
1166 #[cfg(not(feature = "sql"))]
1167 let _ = largest_journaled_prefix;
1168
1169 let batch = scheduler.finish();
1170 let prepared = precommit_preparation(output)?;
1171 let identity_ranges = identity_cursor
1172 .map(IdentityStatementCursor::into_range_advance)
1173 .transpose()?
1174 .into_iter()
1175 .flatten()
1176 .collect::<Vec<_>>();
1177 if batch.is_empty() && !identity_ranges.is_empty() {
1178 return Err(InternalError::identity_corruption());
1179 }
1180 if !batch.is_empty() {
1181 commit_structural_row_ops_with_window_for_path(
1182 &self.db,
1183 entity_path,
1184 batch,
1185 identity_ranges,
1186 "accepted_structural_batch_apply",
1187 )?;
1188 }
1189 Ok(prepared)
1190 }
1191
1192 fn execute_one_accepted_save_mutation(
1193 &self,
1194 catalog: &AcceptedSchemaCatalogContext,
1195 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
1196 mode: MutationMode,
1197 target: AcceptedStructuralMutationTarget,
1198 patch: AcceptedMutationIntentPatch,
1199 ) -> Result<DynamicMutationResult, InternalError> {
1200 let identity = catalog.identity();
1201 let entity_path = identity.entity_path();
1202 let result = self.execute_accepted_structural_save_batch(
1203 catalog,
1204 descriptor,
1205 vec![AcceptedStructuralMutation::save(mode, target, patch)],
1206 Timestamp::now(),
1207 |rows| prepare_dynamic_mutation_result(catalog, descriptor, rows, false),
1208 )?;
1209 record(MetricsEvent::SaveMutation {
1210 entity_path: entity_path.into(),
1211 kind: match mode {
1212 MutationMode::Insert => SaveMutationKind::Insert,
1213 MutationMode::Replace => SaveMutationKind::Replace,
1214 MutationMode::Update => SaveMutationKind::Update,
1215 },
1216 rows_touched: u64::from(result.affected_rows),
1217 });
1218 Ok(result)
1219 }
1220
1221 pub fn execute_trusted_dynamic_mutation(
1228 &self,
1229 request: &DynamicMutation,
1230 ) -> Result<DynamicMutationResult, InternalError> {
1231 self.execute_trusted_dynamic_mutation_batch_with_result_policy(vec![request.clone()], false)
1232 }
1233
1234 pub fn execute_trusted_dynamic_mutation_batch(
1240 &self,
1241 requests: Vec<DynamicMutation>,
1242 ) -> Result<DynamicMutationResult, InternalError> {
1243 self.execute_trusted_dynamic_mutation_batch_with_result_policy(requests, true)
1244 }
1245
1246 fn execute_trusted_dynamic_mutation_batch_with_result_policy(
1247 &self,
1248 requests: Vec<DynamicMutation>,
1249 enforce_mixed_batch_result_bound: bool,
1250 ) -> Result<DynamicMutationResult, InternalError> {
1251 if requests.is_empty() {
1252 return Err(InternalError::mutation_batch_empty());
1253 }
1254 if requests.len() > MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS {
1255 return Err(InternalError::mutation_batch_too_many_items(
1256 requests.len(),
1257 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1258 ));
1259 }
1260 let first = requests
1261 .first()
1262 .ok_or_else(InternalError::mutation_batch_empty)?;
1263 if first.entity().is_empty() {
1264 return Err(InternalError::executor_unsupported());
1265 }
1266 let catalog = self.accepted_schema_catalog_context_for_entity_name(Some(first.entity()))?;
1267 let accepted_identity = catalog.identity();
1268 let descriptor =
1269 AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())?;
1270 let mut mutations = Vec::with_capacity(requests.len());
1271 let mut save_kinds = Vec::with_capacity(requests.len());
1272
1273 for (batch_position, request) in requests.iter().enumerate() {
1274 let batch_position = u32::try_from(batch_position).map_err(|_| {
1275 InternalError::mutation_batch_too_many_items(
1276 requests.len(),
1277 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1278 )
1279 })?;
1280 if request.entity().is_empty() {
1281 return Err(InternalError::executor_unsupported());
1282 }
1283 let item_catalog =
1284 self.accepted_schema_catalog_context_for_entity_name(Some(request.entity()))?;
1285 if item_catalog.identity() != accepted_identity {
1286 return Err(InternalError::mutation_batch_entity_mismatch(
1287 batch_position,
1288 accepted_identity.entity_tag().value(),
1289 item_catalog.identity().entity_tag().value(),
1290 ));
1291 }
1292 let (mutation, save_kind) = lower_dynamic_mutation_intent(
1293 accepted_identity.entity_tag(),
1294 accepted_identity.entity_path(),
1295 &descriptor,
1296 request,
1297 batch_position,
1298 )?;
1299 mutations.push(mutation);
1300 save_kinds.push(save_kind);
1301 }
1302
1303 let entity_path = accepted_identity.entity_path_handle();
1304 let (result, metrics) = self.execute_accepted_structural_mutation_batch_inner(
1305 &catalog,
1306 &descriptor,
1307 mutations,
1308 Timestamp::now(),
1309 false,
1310 |rows| {
1311 if rows.len() != save_kinds.len() {
1312 return Err(InternalError::executor_invariant());
1313 }
1314 let metrics = rows
1315 .iter()
1316 .zip(save_kinds)
1317 .filter_map(|(row, kind)| kind.map(|kind| (kind, row.logical_changed())))
1318 .collect::<Vec<_>>();
1319 let result = prepare_dynamic_mutation_result(
1320 &catalog,
1321 &descriptor,
1322 rows,
1323 enforce_mixed_batch_result_bound,
1324 )?;
1325 Ok((result, metrics))
1326 },
1327 )?;
1328 for (kind, logical_changed) in metrics {
1329 record(MetricsEvent::SaveMutation {
1330 entity_path: entity_path.clone(),
1331 kind,
1332 rows_touched: u64::from(logical_changed),
1333 });
1334 }
1335 Ok(result)
1336 }
1337
1338 #[doc(hidden)]
1341 pub fn execute_trusted_typed_mutation(
1342 &self,
1343 binding: &DynamicTypedEntityBinding,
1344 request: &DynamicTypedMutation,
1345 ) -> Result<Option<DynamicMutationResult>, InternalError> {
1346 let Some(catalog) = self.current_typed_entity_binding_catalog(binding)? else {
1347 return Ok(None);
1348 };
1349 let identity = catalog.identity();
1350 let descriptor =
1351 AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())?;
1352 let mode = dynamic_typed_mutation_mode(request);
1353 let (target, patch) = match request {
1354 DynamicTypedMutation::Insert { patch } => (
1355 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
1356 patch,
1357 ),
1358 DynamicTypedMutation::Update { key, patch }
1359 | DynamicTypedMutation::Replace { key, patch } => (
1360 AcceptedStructuralMutationTarget::expected(dynamic_key(
1361 identity.entity_tag(),
1362 key,
1363 )?),
1364 patch,
1365 ),
1366 };
1367 if !patch.is_bound_to(binding) {
1368 return Ok(None);
1369 }
1370 let patch = lower_typed_patch(
1371 &descriptor,
1372 patch,
1373 mode,
1374 mutation_diagnostic_context(identity.entity_tag(), mode, 0),
1375 )?;
1376 self.execute_one_accepted_save_mutation(&catalog, &descriptor, mode, target, patch)
1377 .map(Some)
1378 }
1379
1380 pub fn execute_trusted_dynamic_insert_batch(
1386 &self,
1387 entity: &str,
1388 patches: Vec<DynamicStructuralPatch>,
1389 ) -> Result<DynamicMutationResult, InternalError> {
1390 let mutations = patches
1391 .into_iter()
1392 .map(|patch| DynamicMutation::Insert {
1393 entity: entity.to_string(),
1394 patch,
1395 })
1396 .collect();
1397 self.execute_trusted_dynamic_mutation_batch_with_result_policy(mutations, false)
1398 }
1399}
1400
1401#[cfg(test)]
1402mod typed_adapter_tests {
1403 use super::{
1404 AcceptedFieldKind, DbSession, DynamicTypedBindingError, DynamicTypedFieldBindingRequest,
1405 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell, dynamic_typed_field_type,
1406 typed_adapter_field_kind_matches,
1407 };
1408 use crate::{
1409 db::{
1410 data::DataStore,
1411 index::IndexStore,
1412 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
1413 schema::{
1414 AcceptedSchemaRevision, FieldId, FieldStorageDecode, LeafCodec,
1415 PersistedFieldSnapshot, PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot,
1416 SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
1417 accepted_schema_candidate_with_field_bindings_for_tests,
1418 },
1419 },
1420 traits::{CanisterKind, Path},
1421 types::EntityTag,
1422 value::InputValue,
1423 };
1424 use icydb_schema::{EntitySourceKey, FieldSourceKey, ScalarType};
1425 use std::{cell::RefCell, collections::BTreeMap};
1426
1427 const STORE_PATH: &str = "session::write::typed_adapter_tests::Store";
1428 const ENTITY_SOURCE: &str = "session::write::typed_adapter_tests::Entity";
1429 const OTHER_ENTITY_SOURCE: &str = "session::write::typed_adapter_tests::OtherEntity";
1430 const ID_SOURCE: &str = "session::write::typed_adapter_tests::Entity::id";
1431 const VALUE_SOURCE: &str = "session::write::typed_adapter_tests::Entity::value";
1432 const REPLACEMENT_SOURCE: &str =
1433 "session::write::typed_adapter_tests::Entity::replacement_value";
1434 const OTHER_ID_SOURCE: &str = "session::write::typed_adapter_tests::OtherEntity::id";
1435
1436 struct TestCanister;
1437
1438 impl Path for TestCanister {
1439 const PATH: &'static str = "session::write::typed_adapter_tests::Canister";
1440 }
1441
1442 impl CanisterKind for TestCanister {
1443 const COMMIT_MEMORY_ID: u8 = 41;
1444 const COMMIT_STABLE_KEY: &'static str = "icydb.typed_adapter_tests.commit.v1";
1445 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 42;
1446 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
1447 "icydb.typed_adapter_tests.integrity.progress.v1";
1448 }
1449
1450 thread_local! {
1451 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
1452 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
1453 static SCHEMA_STORE: RefCell<SchemaStore> =
1454 const { RefCell::new(SchemaStore::init_heap()) };
1455 static STORE_REGISTRY: StoreRegistry = {
1456 let mut registry = StoreRegistry::new();
1457 registry.register_store(
1458 STORE_PATH,
1459 &DATA_STORE,
1460 &INDEX_STORE,
1461 &SCHEMA_STORE,
1462 StoreAllocationIdentities::absent(),
1463 StoreRuntimeStorageCapabilities::heap(),
1464 ).expect("typed adapter test store should register");
1465 registry
1466 };
1467 }
1468
1469 fn nat64_field(id: u32, name: &str, slot: u16) -> PersistedFieldSnapshot {
1470 PersistedFieldSnapshot::new_initial(
1471 FieldId::new(id),
1472 name.to_string(),
1473 SchemaFieldSlot::new(slot),
1474 AcceptedFieldKind::Nat64,
1475 Vec::new(),
1476 false,
1477 SchemaInsertDefault::None,
1478 FieldStorageDecode::ByKind,
1479 LeafCodec::Scalar(ScalarCodec::Nat64),
1480 )
1481 }
1482
1483 fn snapshot(
1484 entity_source: &str,
1485 entity_name: &str,
1486 fields: Vec<PersistedFieldSnapshot>,
1487 ) -> PersistedSchemaSnapshot {
1488 let layout = SchemaRowLayout::initial(
1489 fields
1490 .iter()
1491 .map(|field| (field.id(), field.slot()))
1492 .collect(),
1493 );
1494 PersistedSchemaSnapshot::new(
1495 SchemaVersion::initial(),
1496 entity_source.to_string(),
1497 entity_name.to_string(),
1498 FieldId::new(1),
1499 layout,
1500 fields,
1501 )
1502 }
1503
1504 fn field_source(source: &str) -> FieldSourceKey {
1505 FieldSourceKey::try_new(source).expect("typed field source should admit")
1506 }
1507
1508 fn entity_source(source: &str) -> EntitySourceKey {
1509 EntitySourceKey::try_new(source).expect("typed entity source should admit")
1510 }
1511
1512 fn publish(
1513 session: &DbSession<TestCanister>,
1514 expected: AcceptedSchemaRevision,
1515 revision: AcceptedSchemaRevision,
1516 snapshots: BTreeMap<EntityTag, PersistedSchemaSnapshot>,
1517 fields: BTreeMap<(EntityTag, FieldSourceKey), FieldId>,
1518 ) {
1519 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
1520 STORE_PATH, revision, snapshots, fields,
1521 );
1522 let store = session
1523 .db
1524 .store_handle(STORE_PATH)
1525 .expect("typed adapter test store should resolve");
1526 crate::db::commit::publish_accepted_schema_candidate(
1527 STORE_PATH, store, expected, &candidate,
1528 )
1529 .expect("typed binding candidate should publish");
1530 }
1531
1532 fn request(source: &str) -> DynamicTypedFieldBindingRequest {
1533 DynamicTypedFieldBindingRequest::new(
1534 source.to_string(),
1535 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
1536 false,
1537 )
1538 }
1539
1540 fn assert_query_diagnostic(
1541 error: crate::db::QueryError,
1542 code: icydb_diagnostic_code::DiagnosticCode,
1543 origin: icydb_diagnostic_code::ErrorOrigin,
1544 detail: icydb_diagnostic_code::DiagnosticDetail,
1545 ) {
1546 let diagnostic = error.diagnostic();
1547 assert_eq!(diagnostic.code(), code);
1548 assert_eq!(diagnostic.origin(), origin);
1549 assert_eq!(diagnostic.detail(), Some(&detail));
1550 }
1551
1552 #[test]
1553 fn typed_adapter_kind_matching_is_exact_but_accepts_relation_key_wrappers() {
1554 let relation = AcceptedFieldKind::Relation {
1555 target_path: "test::Target".to_string(),
1556 target_entity_name: "Target".to_string(),
1557 target_entity_tag: EntityTag::new(7),
1558 target_store_path: "test::Store".to_string(),
1559 key_kind: Box::new(AcceptedFieldKind::Nat64),
1560 };
1561
1562 assert!(typed_adapter_field_kind_matches(
1563 &relation,
1564 &AcceptedFieldKind::Nat64,
1565 ));
1566 assert!(typed_adapter_field_kind_matches(
1567 &AcceptedFieldKind::List(Box::new(relation)),
1568 &AcceptedFieldKind::List(Box::new(AcceptedFieldKind::Nat64)),
1569 ));
1570 assert!(!typed_adapter_field_kind_matches(
1571 &AcceptedFieldKind::Nat64,
1572 &AcceptedFieldKind::Nat32,
1573 ));
1574 }
1575
1576 #[test]
1577 fn typed_adapter_field_contract_rejects_invalid_named_source_identity() {
1578 assert!(matches!(
1579 dynamic_typed_field_type(DynamicTypedFieldType::Named(String::new())),
1580 Err(DynamicTypedBindingError::FieldUnavailable),
1581 ));
1582 assert!(matches!(
1583 dynamic_typed_field_type(DynamicTypedFieldType::Scalar(ScalarType::Nat16)),
1584 Ok(icydb_schema::FieldType::Scalar(ScalarType::Nat16)),
1585 ));
1586 }
1587
1588 #[expect(clippy::too_many_lines)]
1591 #[test]
1592 fn typed_binding_uses_accepted_ids_and_slots_across_renames_and_name_reuse() {
1593 let entity_tag = EntityTag::new(91);
1594 let other_entity_tag = EntityTag::new(92);
1595 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
1596 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
1597 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
1598
1599 let session = DbSession::<TestCanister>::new(
1600 &STORE_REGISTRY,
1601 &crate::db::RequestExecutionRoot::__new_runtime_root(),
1602 );
1603 session
1604 .db
1605 .ensure_recovered_state()
1606 .expect("typed adapter test database should initialize");
1607 publish(
1608 &session,
1609 AcceptedSchemaRevision::NONE,
1610 AcceptedSchemaRevision::INITIAL,
1611 BTreeMap::from([(
1612 entity_tag,
1613 snapshot(
1614 ENTITY_SOURCE,
1615 "Entity",
1616 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
1617 ),
1618 )]),
1619 BTreeMap::from([
1620 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1621 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1622 ]),
1623 );
1624
1625 let initial_catalog = session
1626 .find_accepted_schema_catalog_context_for_entity_source_key(ENTITY_SOURCE)
1627 .expect("initial source catalog lookup should inspect")
1628 .expect("initial source catalog should exist");
1629 assert_eq!(initial_catalog.identity().entity_tag(), entity_tag);
1630 let initial = session
1631 .issue_typed_entity_binding(
1632 entity_source(ENTITY_SOURCE).as_str(),
1633 &[request(ID_SOURCE), request(VALUE_SOURCE)],
1634 )
1635 .expect("initial typed binding should issue");
1636 assert_eq!(initial.field_slot(ID_SOURCE), Some(0));
1637 assert_eq!(initial.field_slot(VALUE_SOURCE), Some(1));
1638 assert_eq!(initial.output_field_slot("value"), Some(1));
1639 let initial_patch = initial
1640 .bind_write_fields(vec![(
1641 VALUE_SOURCE.to_string(),
1642 DynamicWriteCell::Value(InputValue::Nat64(7)),
1643 )])
1644 .expect("source-bound patch should lower");
1645 assert_eq!(
1646 initial_patch.fields(),
1647 &[(2, 1, DynamicWriteCell::Value(InputValue::Nat64(7)))]
1648 );
1649
1650 publish(
1651 &session,
1652 AcceptedSchemaRevision::INITIAL,
1653 AcceptedSchemaRevision::new(2),
1654 BTreeMap::from([
1655 (
1656 entity_tag,
1657 snapshot(
1658 ENTITY_SOURCE,
1659 "RenamedEntity",
1660 vec![
1661 nat64_field(1, "id", 0),
1662 nat64_field(2, "renamed_value", 1),
1663 nat64_field(3, "value", 2),
1664 ],
1665 ),
1666 ),
1667 (
1668 other_entity_tag,
1669 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
1670 ),
1671 ]),
1672 BTreeMap::from([
1673 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1674 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1675 (
1676 (entity_tag, field_source(REPLACEMENT_SOURCE)),
1677 FieldId::new(3),
1678 ),
1679 (
1680 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
1681 FieldId::new(1),
1682 ),
1683 ]),
1684 );
1685
1686 let stale_authority = session
1687 .ensure_accepted_schema_authority_is_current_for_store_path(
1688 STORE_PATH,
1689 initial_catalog.value_catalog_handle().authority(),
1690 )
1691 .expect_err("the initial accepted authority must be stale after revision two");
1692 assert_eq!(
1693 stale_authority.diagnostic_facts(),
1694 vec![
1695 (
1696 icydb_diagnostic_code::DiagnosticFactTag::ExpectedRevision,
1697 AcceptedSchemaRevision::INITIAL.get(),
1698 ),
1699 (
1700 icydb_diagnostic_code::DiagnosticFactTag::CurrentRevision,
1701 AcceptedSchemaRevision::new(2).get(),
1702 ),
1703 ],
1704 );
1705
1706 assert!(
1707 !session
1708 .typed_entity_binding_is_current(&initial)
1709 .expect("renamed binding currentness should inspect")
1710 );
1711 let renamed = session
1712 .issue_typed_entity_binding(ENTITY_SOURCE, &[request(ID_SOURCE), request(VALUE_SOURCE)])
1713 .expect("renamed source-bound adapter should rebind");
1714 assert_eq!(renamed.entity(), "RenamedEntity");
1715 assert_eq!(renamed.field_slot(VALUE_SOURCE), Some(1));
1716 assert_eq!(renamed.output_field_slot("renamed_value"), Some(1));
1717 assert_eq!(renamed.output_field_slot("value"), None);
1718
1719 publish(
1720 &session,
1721 AcceptedSchemaRevision::new(2),
1722 AcceptedSchemaRevision::new(3),
1723 BTreeMap::from([
1724 (
1725 entity_tag,
1726 snapshot(
1727 ENTITY_SOURCE,
1728 "RenamedEntity",
1729 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
1730 ),
1731 ),
1732 (
1733 other_entity_tag,
1734 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
1735 ),
1736 ]),
1737 BTreeMap::from([
1738 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1739 (
1740 (entity_tag, field_source(REPLACEMENT_SOURCE)),
1741 FieldId::new(2),
1742 ),
1743 (
1744 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
1745 FieldId::new(1),
1746 ),
1747 ]),
1748 );
1749
1750 assert!(matches!(
1751 session.issue_typed_entity_binding(
1752 ENTITY_SOURCE,
1753 &[request(ID_SOURCE), request(VALUE_SOURCE)],
1754 ),
1755 Err(DynamicTypedBindingError::FieldUnavailable),
1756 ));
1757 assert!(
1758 !session
1759 .typed_entity_binding_is_current(&renamed)
1760 .expect("removed source binding should become stale")
1761 );
1762
1763 let replacement = session
1764 .issue_typed_entity_binding(
1765 ENTITY_SOURCE,
1766 &[request(ID_SOURCE), request(REPLACEMENT_SOURCE)],
1767 )
1768 .expect("explicit replacement source should bind");
1769 assert!(
1770 session
1771 .execute_trusted_typed_mutation(
1772 &replacement,
1773 &DynamicTypedMutation::Insert {
1774 patch: initial_patch
1775 },
1776 )
1777 .expect("cross-binding patch should fail closed")
1778 .is_none()
1779 );
1780 let patch = replacement
1781 .bind_write_fields(vec![
1782 (
1783 ID_SOURCE.to_string(),
1784 DynamicWriteCell::Value(InputValue::Nat64(1)),
1785 ),
1786 (
1787 REPLACEMENT_SOURCE.to_string(),
1788 DynamicWriteCell::Value(InputValue::Nat64(9)),
1789 ),
1790 ])
1791 .expect("replacement source write should bind by accepted IDs and slots");
1792 let result = session
1793 .execute_trusted_typed_mutation(&replacement, &DynamicTypedMutation::Insert { patch })
1794 .expect("typed insert should use the accepted mutation pipeline")
1795 .expect("replacement binding should remain current");
1796 assert_eq!(result.entity, "RenamedEntity");
1797 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
1798 assert_eq!(
1799 result.rows,
1800 vec![vec![
1801 crate::value::OutputValue::Nat64(1),
1802 crate::value::OutputValue::Nat64(9)
1803 ]]
1804 );
1805 assert_eq!(result.affected_rows, 1);
1806
1807 let second_patch = replacement
1808 .bind_write_fields(vec![
1809 (
1810 ID_SOURCE.to_string(),
1811 DynamicWriteCell::Value(InputValue::Nat64(2)),
1812 ),
1813 (
1814 REPLACEMENT_SOURCE.to_string(),
1815 DynamicWriteCell::Value(InputValue::Nat64(10)),
1816 ),
1817 ])
1818 .expect("second source-bound patch should lower");
1819 session
1820 .execute_trusted_typed_mutation(
1821 &replacement,
1822 &DynamicTypedMutation::Insert {
1823 patch: second_patch,
1824 },
1825 )
1826 .expect("second typed insert should use the accepted mutation pipeline")
1827 .expect("replacement binding should remain current");
1828
1829 {
1830 let query = crate::db::DynamicQuery::new("RenamedEntity")
1831 .select(["id", "value"])
1832 .order_by(crate::db::asc("id"))
1833 .limit(1);
1834 let result = session
1835 .execute_trusted_live_page(&query, None)
1836 .expect("SQL-free dynamic execution should use accepted authority");
1837 assert_eq!(result.entity, "RenamedEntity");
1838 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
1839 assert_eq!(
1840 result.rows,
1841 vec![vec![
1842 crate::value::OutputValue::Nat64(1),
1843 crate::value::OutputValue::Nat64(9)
1844 ]]
1845 );
1846 assert_eq!(result.row_count, 1);
1847 assert_query_diagnostic(
1848 session
1849 .execute_trusted_live_page(&query.cursor("00"), None)
1850 .expect_err("scalar execution must reject grouped cursor state"),
1851 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1852 icydb_diagnostic_code::ErrorOrigin::Query,
1853 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1854 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1855 },
1856 );
1857 assert_query_diagnostic(
1858 session
1859 .execute_public_dynamic_grouped_query(
1860 &crate::db::DynamicQuery::new("RenamedEntity").grouped_limits(1, 1024),
1861 )
1862 .expect_err("grouped execution must reject scalar query state"),
1863 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1864 icydb_diagnostic_code::ErrorOrigin::Query,
1865 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1866 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1867 },
1868 );
1869
1870 let grouped_query = crate::db::DynamicQuery::new("RenamedEntity")
1871 .filter(crate::db::FieldRef::new("id").eq(1_u64))
1872 .group_by("value")
1873 .aggregate(crate::db::count())
1874 .grouped_limits(1, 1024)
1875 .limit(1);
1876 let grouped = session
1877 .execute_public_dynamic_grouped_query(&grouped_query)
1878 .expect("SQL-free grouped execution should use accepted authority");
1879 let typed_grouped = session
1880 .execute_public_dynamic_grouped_query_for_typed_binding(
1881 &replacement,
1882 &grouped_query,
1883 )
1884 .expect("typed grouped execution should inspect accepted authority")
1885 .expect("replacement binding should remain current");
1886 assert_eq!(typed_grouped, grouped);
1887 assert!(
1888 session
1889 .execute_public_dynamic_grouped_query_for_typed_binding(
1890 &renamed,
1891 &grouped_query,
1892 )
1893 .expect("stale grouped binding should inspect accepted authority")
1894 .is_none(),
1895 "stale typed grouped bindings must fail closed before execution"
1896 );
1897 assert_eq!(grouped.entity, "RenamedEntity");
1898 assert_eq!(grouped.row_count, 1);
1899 assert_eq!(grouped.rows.len(), 1);
1900 assert_eq!(
1901 grouped.rows[0].group_key(),
1902 &[crate::value::OutputValue::Nat64(9)]
1903 );
1904 assert_eq!(
1905 grouped.rows[0].aggregate_values(),
1906 &[crate::value::OutputValue::Nat64(1)]
1907 );
1908 assert_eq!(grouped.next_cursor, None);
1909
1910 assert_query_diagnostic(
1911 session
1912 .execute_public_dynamic_grouped_query(&grouped_query.clone().select(["value"]))
1913 .expect_err("grouped output must reject scalar selection"),
1914 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1915 icydb_diagnostic_code::ErrorOrigin::Query,
1916 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1917 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1918 },
1919 );
1920 assert_query_diagnostic(
1921 session
1922 .execute_public_dynamic_grouped_query(
1923 &crate::db::DynamicQuery::new("RenamedEntity")
1924 .group_by("value")
1925 .aggregate(crate::db::count()),
1926 )
1927 .expect_err("public grouped execution must require explicit limits"),
1928 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1929 icydb_diagnostic_code::ErrorOrigin::Query,
1930 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1931 reason:
1932 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
1933 },
1934 );
1935 assert_query_diagnostic(
1936 session
1937 .execute_trusted_dynamic_grouped_query(
1938 &crate::db::DynamicQuery::new("RenamedEntity")
1939 .group_by("value")
1940 .aggregate(crate::db::count())
1941 .grouped_limits(0, 1024),
1942 )
1943 .expect_err("trusted grouped execution must reject zero limits"),
1944 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1945 icydb_diagnostic_code::ErrorOrigin::Query,
1946 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1947 reason:
1948 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
1949 },
1950 );
1951 assert_query_diagnostic(
1952 session
1953 .execute_public_dynamic_grouped_query(&grouped_query.grouped_limits(101, 1024))
1954 .expect_err("public grouped execution must enforce its group budget"),
1955 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1956 icydb_diagnostic_code::ErrorOrigin::Query,
1957 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1958 reason:
1959 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryExceedsBudget,
1960 },
1961 );
1962
1963 let paged_query = crate::db::DynamicQuery::new("RenamedEntity")
1964 .group_by("value")
1965 .aggregate(crate::db::count())
1966 .grouped_limits(2, 2 * 1024)
1967 .limit(1);
1968 assert_query_diagnostic(
1969 session
1970 .execute_public_dynamic_grouped_query(&paged_query)
1971 .expect_err("public grouped execution must reject an unbounded full scan"),
1972 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1973 icydb_diagnostic_code::ErrorOrigin::Query,
1974 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1975 reason:
1976 icydb_diagnostic_code::QueryReadAdmissionCode::UnboundedFullScanRejected,
1977 },
1978 );
1979 let first_page = session
1980 .execute_trusted_dynamic_grouped_query(&paged_query)
1981 .expect("SQL-free grouped first page should execute");
1982 assert_eq!(first_page.row_count, 1);
1983 assert_eq!(
1984 first_page.rows[0].group_key(),
1985 &[crate::value::OutputValue::Nat64(9)]
1986 );
1987 let cursor = first_page
1988 .next_cursor
1989 .expect("first grouped page should return a continuation cursor");
1990 assert_query_diagnostic(
1991 session
1992 .execute_trusted_dynamic_grouped_query(
1993 &paged_query.clone().cursor(format!("{cursor}0")),
1994 )
1995 .expect_err("tampered grouped cursor must fail closed"),
1996 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
1997 icydb_diagnostic_code::ErrorOrigin::Cursor,
1998 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1999 kind: icydb_diagnostic_code::QueryErrorKind::InvalidContinuationCursor,
2000 },
2001 );
2002 let second_page = session
2003 .execute_trusted_dynamic_grouped_query(&paged_query.cursor(cursor))
2004 .expect("SQL-free grouped continuation should execute");
2005 assert_eq!(second_page.row_count, 1);
2006 assert_eq!(
2007 second_page.rows[0].group_key(),
2008 &[crate::value::OutputValue::Nat64(10)]
2009 );
2010 assert_eq!(second_page.next_cursor, None);
2011 }
2012 }
2013}
2014
2015#[cfg(test)]
2016mod mixed_relation_batch_tests {
2017 use super::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell};
2018 use crate::{
2019 db::{
2020 DynamicQuery, asc,
2021 data::DataStore,
2022 desc,
2023 index::IndexStore,
2024 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
2025 schema::{
2026 AcceptedConstraintCatalog, AcceptedFieldKind, AcceptedSchemaRevision, FieldId,
2027 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2028 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2029 PersistedRelationEdgeSnapshot, PersistedSchemaSnapshot, RelationId, ScalarCodec,
2030 SchemaFieldSlot, SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore,
2031 SchemaVersion, accepted_schema_candidate_with_field_bindings_for_tests,
2032 },
2033 },
2034 error::ErrorClass,
2035 traits::{CanisterKind, Path},
2036 types::EntityTag,
2037 value::{InputValue, OutputValue},
2038 };
2039 use icydb_schema::FieldSourceKey;
2040 use std::{cell::RefCell, collections::BTreeMap};
2041
2042 const STORE_PATH: &str = "session::write::mixed_relation_batch_tests::Store";
2043 const ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node";
2044 const ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::id";
2045 const PARENT_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::parent_id";
2046 const CODE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::code";
2047 const ENTITY_NAME: &str = "MixedRelationNode";
2048 const ENTITY_TAG: EntityTag = EntityTag::new(94);
2049 const OTHER_ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other";
2050 const OTHER_ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::id";
2051 const OTHER_VALUE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::value";
2052 const OTHER_ENTITY_NAME: &str = "MixedRelationOther";
2053 const OTHER_ENTITY_TAG: EntityTag = EntityTag::new(95);
2054
2055 struct TestCanister;
2056
2057 impl Path for TestCanister {
2058 const PATH: &'static str = "session::write::mixed_relation_batch_tests::Canister";
2059 }
2060
2061 impl CanisterKind for TestCanister {
2062 const COMMIT_MEMORY_ID: u8 = 47;
2063 const COMMIT_STABLE_KEY: &'static str = "icydb.mixed_relation_batch_tests.commit.v1";
2064 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 48;
2065 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2066 "icydb.mixed_relation_batch_tests.integrity.progress.v1";
2067 }
2068
2069 thread_local! {
2070 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2071 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2072 static SCHEMA_STORE: RefCell<SchemaStore> =
2073 const { RefCell::new(SchemaStore::init_heap()) };
2074 static STORE_REGISTRY: StoreRegistry = {
2075 let mut registry = StoreRegistry::new();
2076 registry.register_store(
2077 STORE_PATH,
2078 &DATA_STORE,
2079 &INDEX_STORE,
2080 &SCHEMA_STORE,
2081 StoreAllocationIdentities::absent(),
2082 StoreRuntimeStorageCapabilities::heap(),
2083 ).expect("mixed relation test store should register");
2084 registry
2085 };
2086 }
2087
2088 fn source_key(source: &str) -> FieldSourceKey {
2089 FieldSourceKey::try_new(source).expect("mixed relation field source should admit")
2090 }
2091
2092 fn relation_snapshot() -> PersistedSchemaSnapshot {
2093 let fields = vec![
2094 PersistedFieldSnapshot::new_initial(
2095 FieldId::new(1),
2096 "id".to_string(),
2097 SchemaFieldSlot::new(0),
2098 AcceptedFieldKind::Nat64,
2099 Vec::new(),
2100 false,
2101 SchemaInsertDefault::None,
2102 FieldStorageDecode::ByKind,
2103 LeafCodec::Scalar(ScalarCodec::Nat64),
2104 ),
2105 PersistedFieldSnapshot::new_initial(
2106 FieldId::new(2),
2107 "parent_id".to_string(),
2108 SchemaFieldSlot::new(1),
2109 AcceptedFieldKind::Nat64,
2110 Vec::new(),
2111 true,
2112 SchemaInsertDefault::None,
2113 FieldStorageDecode::ByKind,
2114 LeafCodec::Scalar(ScalarCodec::Nat64),
2115 ),
2116 PersistedFieldSnapshot::new_initial(
2117 FieldId::new(3),
2118 "code".to_string(),
2119 SchemaFieldSlot::new(2),
2120 AcceptedFieldKind::Nat64,
2121 Vec::new(),
2122 false,
2123 SchemaInsertDefault::None,
2124 FieldStorageDecode::ByKind,
2125 LeafCodec::Scalar(ScalarCodec::Nat64),
2126 ),
2127 ];
2128 let relation = PersistedRelationEdgeSnapshot::new(
2129 RelationId::new(1).expect("mixed relation identity should be non-zero"),
2130 "parent".to_string(),
2131 ENTITY_SOURCE.to_string(),
2132 vec![FieldId::new(2)],
2133 );
2134 let snapshot = PersistedSchemaSnapshot::new_with_indexes(
2135 SchemaVersion::initial(),
2136 ENTITY_SOURCE.to_string(),
2137 ENTITY_NAME.to_string(),
2138 FieldId::new(1),
2139 SchemaRowLayout::initial(
2140 fields
2141 .iter()
2142 .map(|field| (field.id(), field.slot()))
2143 .collect(),
2144 ),
2145 fields,
2146 vec![PersistedIndexSnapshot::new(
2147 SchemaIndexId::new(1).expect("mixed unique index identity should be non-zero"),
2148 1,
2149 "by_code".to_string(),
2150 STORE_PATH.to_string(),
2151 true,
2152 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2153 FieldId::new(3),
2154 SchemaFieldSlot::new(2),
2155 vec!["code".to_string()],
2156 AcceptedFieldKind::Nat64,
2157 false,
2158 )]),
2159 None,
2160 )],
2161 )
2162 .with_relations(vec![relation]);
2163 let constraints = AcceptedConstraintCatalog::initial(
2164 snapshot.fields(),
2165 snapshot.indexes(),
2166 snapshot.relations(),
2167 )
2168 .expect("mixed relation constraints should close");
2169 snapshot.with_constraint_catalog(constraints)
2170 }
2171
2172 fn other_snapshot() -> PersistedSchemaSnapshot {
2173 let fields = vec![
2174 PersistedFieldSnapshot::new_initial(
2175 FieldId::new(1),
2176 "id".to_string(),
2177 SchemaFieldSlot::new(0),
2178 AcceptedFieldKind::Nat64,
2179 Vec::new(),
2180 false,
2181 SchemaInsertDefault::None,
2182 FieldStorageDecode::ByKind,
2183 LeafCodec::Scalar(ScalarCodec::Nat64),
2184 ),
2185 PersistedFieldSnapshot::new_initial(
2186 FieldId::new(2),
2187 "value".to_string(),
2188 SchemaFieldSlot::new(1),
2189 AcceptedFieldKind::Nat64,
2190 Vec::new(),
2191 false,
2192 SchemaInsertDefault::None,
2193 FieldStorageDecode::ByKind,
2194 LeafCodec::Scalar(ScalarCodec::Nat64),
2195 ),
2196 ];
2197 PersistedSchemaSnapshot::new(
2198 SchemaVersion::initial(),
2199 OTHER_ENTITY_SOURCE.to_string(),
2200 OTHER_ENTITY_NAME.to_string(),
2201 FieldId::new(1),
2202 SchemaRowLayout::initial(
2203 fields
2204 .iter()
2205 .map(|field| (field.id(), field.slot()))
2206 .collect(),
2207 ),
2208 fields,
2209 )
2210 }
2211
2212 fn initialize() -> DbSession<TestCanister> {
2213 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2214 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2215 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2216 let session = DbSession::<TestCanister>::new(
2217 &STORE_REGISTRY,
2218 &crate::db::RequestExecutionRoot::__new_runtime_root(),
2219 );
2220 session
2221 .db
2222 .ensure_recovered_state()
2223 .expect("mixed relation database should initialize");
2224 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2225 STORE_PATH,
2226 AcceptedSchemaRevision::INITIAL,
2227 BTreeMap::from([
2228 (ENTITY_TAG, relation_snapshot()),
2229 (OTHER_ENTITY_TAG, other_snapshot()),
2230 ]),
2231 BTreeMap::from([
2232 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2233 ((ENTITY_TAG, source_key(PARENT_SOURCE)), FieldId::new(2)),
2234 ((ENTITY_TAG, source_key(CODE_SOURCE)), FieldId::new(3)),
2235 (
2236 (OTHER_ENTITY_TAG, source_key(OTHER_ID_SOURCE)),
2237 FieldId::new(1),
2238 ),
2239 (
2240 (OTHER_ENTITY_TAG, source_key(OTHER_VALUE_SOURCE)),
2241 FieldId::new(2),
2242 ),
2243 ]),
2244 );
2245 let store = session
2246 .db
2247 .store_handle(STORE_PATH)
2248 .expect("mixed relation store should resolve");
2249 crate::db::commit::publish_accepted_schema_candidate(
2250 STORE_PATH,
2251 store,
2252 AcceptedSchemaRevision::NONE,
2253 &candidate,
2254 )
2255 .expect("mixed relation candidate should publish");
2256 session
2257 }
2258
2259 fn patch(id: Option<u64>, parent: Option<u64>, code: Option<u64>) -> DynamicStructuralPatch {
2260 let mut fields = Vec::new();
2261 if let Some(id) = id {
2262 fields.push((
2263 "id".to_string(),
2264 DynamicWriteCell::Value(InputValue::Nat64(id)),
2265 ));
2266 }
2267 fields.push((
2268 "parent_id".to_string(),
2269 parent.map_or(DynamicWriteCell::Null, |parent| {
2270 DynamicWriteCell::Value(InputValue::Nat64(parent))
2271 }),
2272 ));
2273 if let Some(code) = code {
2274 fields.push((
2275 "code".to_string(),
2276 DynamicWriteCell::Value(InputValue::Nat64(code)),
2277 ));
2278 }
2279 DynamicStructuralPatch::new(fields)
2280 }
2281
2282 fn insert(id: u64, parent: Option<u64>) -> DynamicMutation {
2283 insert_with_code(id, parent, id)
2284 }
2285
2286 fn insert_with_code(id: u64, parent: Option<u64>, code: u64) -> DynamicMutation {
2287 DynamicMutation::Insert {
2288 entity: ENTITY_NAME.to_string(),
2289 patch: patch(Some(id), parent, Some(code)),
2290 }
2291 }
2292
2293 fn update_parent(id: u64, parent: Option<u64>) -> DynamicMutation {
2294 DynamicMutation::Update {
2295 entity: ENTITY_NAME.to_string(),
2296 key: InputValue::Nat64(id),
2297 patch: patch(None, parent, None),
2298 }
2299 }
2300
2301 fn update_code(id: u64, code: u64) -> DynamicMutation {
2302 DynamicMutation::Update {
2303 entity: ENTITY_NAME.to_string(),
2304 key: InputValue::Nat64(id),
2305 patch: DynamicStructuralPatch::new(vec![(
2306 "code".to_string(),
2307 DynamicWriteCell::Value(InputValue::Nat64(code)),
2308 )]),
2309 }
2310 }
2311
2312 fn delete(id: u64) -> DynamicMutation {
2313 DynamicMutation::Delete {
2314 entity: ENTITY_NAME.to_string(),
2315 key: InputValue::Nat64(id),
2316 }
2317 }
2318
2319 fn expected_row(id: u64, parent: Option<u64>) -> Vec<OutputValue> {
2320 expected_row_with_code(id, parent, id)
2321 }
2322
2323 fn expected_row_with_code(id: u64, parent: Option<u64>, code: u64) -> Vec<OutputValue> {
2324 vec![
2325 OutputValue::Nat64(id),
2326 parent.map_or(OutputValue::Null, OutputValue::Nat64),
2327 OutputValue::Nat64(code),
2328 ]
2329 }
2330
2331 fn other_patch(id: Option<u64>, value: u64) -> DynamicStructuralPatch {
2332 let mut fields = Vec::new();
2333 if let Some(id) = id {
2334 fields.push((
2335 "id".to_string(),
2336 DynamicWriteCell::Value(InputValue::Nat64(id)),
2337 ));
2338 }
2339 fields.push((
2340 "value".to_string(),
2341 DynamicWriteCell::Value(InputValue::Nat64(value)),
2342 ));
2343 DynamicStructuralPatch::new(fields)
2344 }
2345
2346 fn assert_relation_violation(error: &crate::error::InternalError) {
2347 assert!(error.diagnostic_facts().contains(&(
2348 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
2349 icydb_diagnostic_code::DiagnosticConstraintKind::Relation.raw(),
2350 )));
2351 }
2352
2353 #[test]
2354 fn live_pages_resume_mixed_projection_from_authenticated_hidden_order_values() {
2355 let session = initialize();
2356 session
2357 .execute_trusted_dynamic_mutation_batch(vec![
2358 insert_with_code(1, None, 10),
2359 insert_with_code(2, Some(1), 20),
2360 insert_with_code(3, None, 30),
2361 ])
2362 .expect("live-page rows should insert");
2363 let query = DynamicQuery::new(ENTITY_NAME)
2364 .select(["id"])
2365 .order_by(desc("code"));
2366
2367 let first = session
2368 .execute_public_live_page(&query, None)
2369 .expect("initial live page should execute");
2370 assert_eq!(
2371 first.rows,
2372 vec![vec![OutputValue::Nat64(3)], vec![OutputValue::Nat64(2)]]
2373 );
2374 let cursor = first
2375 .continuation
2376 .as_deref()
2377 .expect("unreturned matching row should produce continuation");
2378 let second = session
2379 .execute_public_live_page(&query, Some(cursor))
2380 .expect("authenticated live continuation should resume");
2381 assert_eq!(second.rows, vec![vec![OutputValue::Nat64(1)]]);
2382 assert_eq!(second.continuation, None);
2383
2384 let total_limit = session
2385 .execute_public_live_page(&query.clone().limit(2), None)
2386 .expect("total live-page limit should execute");
2387 assert_eq!(
2388 total_limit.rows,
2389 vec![vec![OutputValue::Nat64(3)], vec![OutputValue::Nat64(2)]],
2390 );
2391 assert_eq!(
2392 total_limit.continuation, None,
2393 "query LIMIT is a total traversal window rather than a page size",
2394 );
2395
2396 let three_row_window = query.clone().limit(3);
2397 let limited_first = session
2398 .execute_public_live_page(&three_row_window, None)
2399 .expect("first total-window page should execute");
2400 let limited_cursor = limited_first
2401 .continuation
2402 .as_deref()
2403 .expect("a partially consumed total window should continue");
2404 let limited_second = session
2405 .execute_public_live_page(&three_row_window, Some(limited_cursor))
2406 .expect("remaining total window should preserve the plan signature");
2407 assert_eq!(limited_second.rows, vec![vec![OutputValue::Nat64(1)]]);
2408 assert_eq!(limited_second.continuation, None);
2409
2410 let mixed_order = DynamicQuery::new(ENTITY_NAME)
2411 .select(["id"])
2412 .order_by(desc("parent_id"))
2413 .order_by(asc("id"));
2414 let mixed_first = session
2415 .execute_trusted_live_page(&mixed_order, None)
2416 .expect("mixed-direction nullable order should execute");
2417 assert_eq!(
2418 mixed_first.rows,
2419 vec![vec![OutputValue::Nat64(2)], vec![OutputValue::Nat64(1)]],
2420 );
2421 let mixed_cursor = mixed_first
2422 .continuation
2423 .as_deref()
2424 .expect("duplicate null order values should retain continuation");
2425 let mixed_second = session
2426 .execute_trusted_live_page(&mixed_order, Some(mixed_cursor))
2427 .expect("mixed-direction nullable order should resume");
2428 assert_eq!(mixed_second.rows, vec![vec![OutputValue::Nat64(3)]]);
2429 assert_eq!(mixed_second.continuation, None);
2430
2431 let mismatched_window = session
2432 .execute_public_live_page(&query.clone().limit(3), Some(cursor))
2433 .expect_err("a changed total limit must invalidate the continuation");
2434 assert_eq!(
2435 mismatched_window.diagnostic_code(),
2436 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
2437 );
2438
2439 let mut tampered = cursor.as_bytes().to_vec();
2440 let last = tampered.len().saturating_sub(1);
2441 tampered[last] = if tampered[last] == b'0' { b'1' } else { b'0' };
2442 let tampered = String::from_utf8(tampered).expect("hex cursor should remain UTF-8");
2443 let error = session
2444 .execute_public_live_page(&query, Some(tampered.as_str()))
2445 .expect_err("tampered cursor must fail closed");
2446 assert_eq!(
2447 error.diagnostic_code(),
2448 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
2449 );
2450 }
2451
2452 #[test]
2453 fn accepted_relation_edges_drive_catalog_and_describe_introspection() {
2454 let session = initialize();
2455 let entities = session
2456 .show_entities()
2457 .expect("accepted entity catalog should resolve");
2458 let source = entities
2459 .iter()
2460 .find(|entity| entity.entity_name() == ENTITY_NAME)
2461 .expect("relation source should be listed");
2462 assert_eq!(source.relations(), 1);
2463
2464 let description = session
2465 .try_describe_entity_by_name(ENTITY_NAME)
2466 .expect("accepted relation source should describe");
2467 let [relation] = description.relations() else {
2468 panic!("accepted relation edge should produce one relation row");
2469 };
2470 assert_eq!(relation.field(), "parent_id");
2471 assert_eq!(relation.target_path(), ENTITY_SOURCE);
2472 assert_eq!(relation.target_entity_name(), ENTITY_NAME);
2473 assert_eq!(relation.target_store_path(), STORE_PATH);
2474 assert_eq!(
2475 relation.cardinality(),
2476 crate::db::EntityRelationCardinality::Single,
2477 );
2478 }
2479
2480 #[test]
2481 fn mixed_relation_validation_uses_the_complete_final_row_overlay() {
2482 let session = initialize();
2483 session
2484 .execute_trusted_dynamic_mutation_batch(vec![insert(1, None), insert(2, Some(1))])
2485 .expect("the initial relation should commit");
2486
2487 let blocked = session
2488 .execute_trusted_dynamic_mutation(&delete(1))
2489 .expect_err("an unaffected committed source must block target deletion");
2490 assert_relation_violation(&blocked);
2491
2492 let deleted = session
2493 .execute_trusted_dynamic_mutation_batch(vec![delete(2), delete(1)])
2494 .expect("a source and its target should delete atomically");
2495 assert_eq!(
2496 deleted.rows,
2497 vec![expected_row(2, Some(1)), expected_row(1, None)],
2498 );
2499
2500 session
2501 .execute_trusted_dynamic_mutation_batch(vec![insert(3, None), insert(4, Some(3))])
2502 .expect("the update-away fixture should commit");
2503 let updated_away = session
2504 .execute_trusted_dynamic_mutation_batch(vec![update_parent(4, None), delete(3)])
2505 .expect("an updated final source may release a deleted target");
2506 assert_eq!(
2507 updated_away.rows,
2508 vec![expected_row(4, None), expected_row(3, None)],
2509 );
2510
2511 session
2512 .execute_trusted_dynamic_mutation_batch(vec![insert(5, None), insert(6, Some(5))])
2513 .expect("the retained-reference fixture should commit");
2514 let retained = session
2515 .execute_trusted_dynamic_mutation_batch(vec![update_parent(6, Some(5)), delete(5)])
2516 .expect_err("a final updated source must still block target deletion");
2517 assert_relation_violation(&retained);
2518
2519 session
2520 .execute_trusted_dynamic_mutation(&insert(7, None))
2521 .expect("the inserted-reference fixture target should commit");
2522 let inserted_reference = session
2523 .execute_trusted_dynamic_mutation_batch(vec![insert(8, Some(7)), delete(7)])
2524 .expect_err("a final inserted source must not reference a deleted target");
2525 assert_relation_violation(&inserted_reference);
2526
2527 let inserted_target = session
2528 .execute_trusted_dynamic_mutation_batch(vec![insert(10, Some(9)), insert(9, None)])
2529 .expect("an inserted relation should see its batch-final target");
2530 assert_eq!(
2531 inserted_target.rows,
2532 vec![expected_row(10, Some(9)), expected_row(9, None)],
2533 );
2534
2535 session
2536 .execute_trusted_dynamic_mutation(&insert(11, None))
2537 .expect("the updated-reference fixture source should commit");
2538 let updated_target = session
2539 .execute_trusted_dynamic_mutation_batch(vec![
2540 update_parent(11, Some(12)),
2541 insert(12, None),
2542 ])
2543 .expect("an updated relation should see its batch-final target");
2544 assert_eq!(
2545 updated_target.rows,
2546 vec![expected_row(11, Some(12)), expected_row(12, None)],
2547 );
2548 }
2549
2550 #[test]
2551 fn mixed_batch_rejects_cross_entity_missing_and_collision_then_honors_replace() {
2552 let session = initialize();
2553 session
2554 .execute_trusted_dynamic_mutation(&insert(1, None))
2555 .expect("the primary mixed fixture row should commit");
2556 session
2557 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
2558 entity: OTHER_ENTITY_NAME.to_string(),
2559 patch: other_patch(Some(1), 10),
2560 })
2561 .expect("the secondary mixed fixture row should commit");
2562
2563 let mixed_entity = session
2564 .execute_trusted_dynamic_mutation_batch(vec![
2565 update_code(1, 11),
2566 DynamicMutation::Update {
2567 entity: OTHER_ENTITY_NAME.to_string(),
2568 key: InputValue::Nat64(1),
2569 patch: other_patch(None, 11),
2570 },
2571 ])
2572 .expect_err("one atomic batch must not cross accepted entities");
2573 assert_eq!(mixed_entity.class(), ErrorClass::Conflict);
2574 assert_eq!(
2575 mixed_entity.diagnostic_facts(),
2576 vec![
2577 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
2578 (
2579 icydb_diagnostic_code::DiagnosticFactTag::ExpectedEntityTag,
2580 ENTITY_TAG.value(),
2581 ),
2582 (
2583 icydb_diagnostic_code::DiagnosticFactTag::ActualEntityTag,
2584 OTHER_ENTITY_TAG.value(),
2585 ),
2586 ],
2587 );
2588
2589 let missing = session
2590 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 12), delete(99)])
2591 .expect_err("a late missing delete must reject the earlier staged update");
2592 assert_eq!(missing.class(), ErrorClass::NotFound);
2593
2594 session
2595 .execute_trusted_dynamic_mutation(&insert(2, None))
2596 .expect("the collision fixture should commit");
2597 let collision = session
2598 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 13), insert(2, None)])
2599 .expect_err("an insert collision must reject the earlier staged update");
2600 assert_eq!(collision.class(), ErrorClass::Conflict);
2601 let failures_unchanged = session
2602 .execute_trusted_dynamic_mutation(&update_code(1, 1))
2603 .expect("failed batches must preserve the original unique value");
2604 assert_eq!(failures_unchanged.affected_rows, 0);
2605
2606 let replaced = session
2607 .execute_trusted_dynamic_mutation_batch(vec![
2608 update_code(1, 14),
2609 DynamicMutation::Replace {
2610 entity: ENTITY_NAME.to_string(),
2611 key: InputValue::Nat64(99),
2612 patch: patch(None, None, Some(99)),
2613 },
2614 ])
2615 .expect("ordinary caller-key replace should insert its absent final row");
2616 assert_eq!(
2617 replaced.rows,
2618 vec![
2619 expected_row_with_code(1, None, 14),
2620 expected_row_with_code(99, None, 99),
2621 ],
2622 );
2623
2624 let unchanged = session
2625 .execute_trusted_dynamic_mutation(&update_code(1, 14))
2626 .expect("the successful mixed replace must publish its preceding update");
2627 assert_eq!(unchanged.affected_rows, 0);
2628 let other_unchanged = session
2629 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
2630 entity: OTHER_ENTITY_NAME.to_string(),
2631 key: InputValue::Nat64(1),
2632 patch: other_patch(None, 10),
2633 })
2634 .expect("cross-entity rejection must preserve the secondary row");
2635 assert_eq!(other_unchanged.affected_rows, 0);
2636 }
2637
2638 #[test]
2639 fn mixed_batch_unique_swap_and_delete_release_use_the_final_overlay() {
2640 let session = initialize();
2641 session
2642 .execute_trusted_dynamic_mutation_batch(vec![
2643 insert_with_code(1, None, 10),
2644 insert_with_code(2, None, 20),
2645 ])
2646 .expect("the unique-overlay fixture should commit");
2647
2648 let swapped = session
2649 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 20), update_code(2, 10)])
2650 .expect("two final rows should atomically swap unique memberships");
2651 assert_eq!(
2652 swapped.rows,
2653 vec![
2654 expected_row_with_code(1, None, 20),
2655 expected_row_with_code(2, None, 10),
2656 ],
2657 );
2658
2659 let released = session
2660 .execute_trusted_dynamic_mutation_batch(vec![delete(1), insert_with_code(3, None, 20)])
2661 .expect("a delete should release unique membership to a final inserted row");
2662 assert_eq!(
2663 released.rows,
2664 vec![
2665 expected_row_with_code(1, None, 20),
2666 expected_row_with_code(3, None, 20),
2667 ],
2668 );
2669 }
2670}
2671
2672#[cfg(test)]
2673mod identity_pre_key_tests {
2674 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2675 use super::DynamicTypedEntityBinding;
2676 use super::{
2677 AcceptedMutationIntentPatch, AcceptedRowLayoutRuntimeContract, AcceptedStructuralMutation,
2678 AcceptedStructuralMutationTarget, DbSession, DynamicMutation, DynamicStructuralPatch,
2679 DynamicTypedFieldBindingRequest, DynamicTypedFieldType, DynamicTypedMutation,
2680 DynamicWriteCell, FieldSlot, MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
2681 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
2682 add_structural_mutation_staged_bytes, checked_pre_key_candidate_count,
2683 insert_key_exists_after_generation, validate_structural_mutation_result_bytes,
2684 };
2685 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2686 use crate::db::data::DecodedDataStoreKey;
2687 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2688 use crate::db::executor::budget::{
2689 HardExecutionBudget, HardExecutionContext, HardExecutionFailureHeadroom,
2690 with_query_execution_budget_for_tests,
2691 };
2692 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2693 use crate::db::{
2694 CompareProofAndAdvanceError, DynamicQuery, ExhaustiveReadError, RawDataStoreKey,
2695 ReadSetRevisionError, ResumableJobAdvance, ResumableJobAdvanceRequest,
2696 ResumableJobAdvanceStatus, ResumableJobError, ResumableJobId, ResumableJobIdempotencyKey,
2697 ResumableJobStatus, asc,
2698 };
2699 use crate::{
2700 db::{
2701 commit::{database_incarnation_id, forget_recovered_domain_for_tests},
2702 data::DataStore,
2703 executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
2704 index::IndexStore,
2705 integrity::{
2706 PhysicalUnitCheckpoint, QuickIntegrityStatus, RowInspectionLimits,
2707 execute_quick_integrity, execute_row_integrity_page,
2708 },
2709 journal::JournalTailStore,
2710 registry::{
2711 StoreAllocationIdentities, StoreAllocationIdentity, StoreRegistry,
2712 StoreRuntimeStorageCapabilities,
2713 },
2714 schema::{
2715 AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
2716 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2717 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2718 PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
2719 SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
2720 accepted_schema_candidate_with_field_bindings_for_tests,
2721 },
2722 write_context::MutationMode,
2723 },
2724 error::{ErrorClass, ErrorOrigin, InternalError},
2725 testing::test_memory,
2726 traits::{CanisterKind, Path},
2727 types::{EntityTag, Timestamp},
2728 value::{InputValue, OutputValue, Value},
2729 };
2730 use icydb_schema::{FieldSourceKey, ScalarType};
2731 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2732 use std::cell::Cell;
2733 use std::{cell::RefCell, collections::BTreeMap, time::Instant};
2734
2735 const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
2736 const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
2737 const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
2738 const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
2739 const ENTITY_NAME: &str = "IdentityRow";
2740 const ENTITY_TAG: EntityTag = EntityTag::new(93);
2741 const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
2742 const UNRELATED_STORE_PATH: &str = "session::write::identity_pre_key_tests::UnrelatedStore";
2743
2744 struct TestCanister;
2745
2746 impl Path for TestCanister {
2747 const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
2748 }
2749
2750 impl CanisterKind for TestCanister {
2751 const COMMIT_MEMORY_ID: u8 = 45;
2752 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
2753 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
2754 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2755 "icydb.identity_pre_key_tests.integrity.progress.v1";
2756 }
2757
2758 thread_local! {
2759 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2760 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2761 static SCHEMA_STORE: RefCell<SchemaStore> =
2762 const { RefCell::new(SchemaStore::init_heap()) };
2763 static UNRELATED_DATA_STORE: RefCell<DataStore> =
2764 const { RefCell::new(DataStore::init_heap()) };
2765 static UNRELATED_INDEX_STORE: RefCell<IndexStore> =
2766 const { RefCell::new(IndexStore::init_heap()) };
2767 static UNRELATED_SCHEMA_STORE: RefCell<SchemaStore> =
2768 const { RefCell::new(SchemaStore::init_heap()) };
2769 static STORE_REGISTRY: StoreRegistry = {
2770 let mut registry = StoreRegistry::new();
2771 registry.register_store(
2772 STORE_PATH,
2773 &DATA_STORE,
2774 &INDEX_STORE,
2775 &SCHEMA_STORE,
2776 StoreAllocationIdentities::absent(),
2777 StoreRuntimeStorageCapabilities::heap(),
2778 ).expect("identity pre-key test store should register");
2779 registry.register_store(
2780 UNRELATED_STORE_PATH,
2781 &UNRELATED_DATA_STORE,
2782 &UNRELATED_INDEX_STORE,
2783 &UNRELATED_SCHEMA_STORE,
2784 StoreAllocationIdentities::absent(),
2785 StoreRuntimeStorageCapabilities::heap(),
2786 ).expect("unrelated identity test store should register");
2787 registry
2788 };
2789 static JOURNALED_DATA_STORE: RefCell<DataStore> =
2790 RefCell::new(DataStore::init_journaled(test_memory(186)));
2791 static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
2792 RefCell::new(IndexStore::init_journaled(test_memory(187)));
2793 static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
2794 RefCell::new(SchemaStore::init_journaled(test_memory(188)));
2795 static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
2796 RefCell::new(JournalTailStore::init(test_memory(189)));
2797 static JOURNALED_STORE_REGISTRY: StoreRegistry = {
2798 let mut registry = StoreRegistry::new();
2799 registry.register_journaled_store(
2800 JOURNALED_STORE_PATH,
2801 &JOURNALED_DATA_STORE,
2802 &JOURNALED_INDEX_STORE,
2803 &JOURNALED_SCHEMA_STORE,
2804 &JOURNALED_TAIL_STORE,
2805 StoreAllocationIdentities::new_journaled(
2806 StoreAllocationIdentity::new(186, "icydb.test.identity-range.data.v1"),
2807 StoreAllocationIdentity::new(187, "icydb.test.identity-range.index.v1"),
2808 StoreAllocationIdentity::new(188, "icydb.test.identity-range.schema.v1"),
2809 StoreAllocationIdentity::new(189, "icydb.test.identity-range.journal.v1"),
2810 ),
2811 StoreRuntimeStorageCapabilities::journaled(),
2812 ).expect("identity range journaled store should register");
2813 registry
2814 };
2815 }
2816
2817 struct JournaledTestCanister;
2818
2819 impl Path for JournaledTestCanister {
2820 const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
2821 }
2822
2823 impl CanisterKind for JournaledTestCanister {
2824 const COMMIT_MEMORY_ID: u8 = 190;
2825 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
2826 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
2827 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2828 "icydb.identity_range_tests.integrity.progress.v1";
2829 }
2830
2831 fn source_key(source: &str) -> FieldSourceKey {
2832 FieldSourceKey::try_new(source).expect("identity test field source should admit")
2833 }
2834
2835 fn identity_snapshot(store_path: &str) -> PersistedSchemaSnapshot {
2836 let fields = vec![
2837 PersistedFieldSnapshot::new_initial_with_write_policy(
2838 FieldId::new(1),
2839 "id".to_string(),
2840 SchemaFieldSlot::new(0),
2841 AcceptedFieldKind::Nat64,
2842 Vec::new(),
2843 false,
2844 SchemaInsertDefault::None,
2845 SchemaFieldWritePolicy::from_model_policies(
2846 Some(FieldInsertGeneration::Identity),
2847 None,
2848 ),
2849 FieldStorageDecode::ByKind,
2850 LeafCodec::Scalar(ScalarCodec::Nat64),
2851 ),
2852 PersistedFieldSnapshot::new_initial(
2853 FieldId::new(2),
2854 "payload".to_string(),
2855 SchemaFieldSlot::new(1),
2856 AcceptedFieldKind::Nat64,
2857 Vec::new(),
2858 false,
2859 SchemaInsertDefault::None,
2860 FieldStorageDecode::ByKind,
2861 LeafCodec::Scalar(ScalarCodec::Nat64),
2862 ),
2863 ];
2864 PersistedSchemaSnapshot::new_with_indexes(
2865 SchemaVersion::initial(),
2866 ENTITY_SOURCE.to_string(),
2867 ENTITY_NAME.to_string(),
2868 FieldId::new(1),
2869 SchemaRowLayout::initial(
2870 fields
2871 .iter()
2872 .map(|field| (field.id(), field.slot()))
2873 .collect(),
2874 ),
2875 fields,
2876 vec![PersistedIndexSnapshot::new(
2877 SchemaIndexId::new(1).expect("identity test index ID should admit"),
2878 1,
2879 "by_payload".to_string(),
2880 store_path.to_string(),
2881 false,
2882 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2883 FieldId::new(2),
2884 SchemaFieldSlot::new(1),
2885 vec!["payload".to_string()],
2886 AcceptedFieldKind::Nat64,
2887 false,
2888 )]),
2889 None,
2890 )],
2891 )
2892 }
2893
2894 fn initialize() -> DbSession<TestCanister> {
2895 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2896 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2897 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2898 UNRELATED_DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2899 UNRELATED_INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2900 UNRELATED_SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2901 let session = DbSession::<TestCanister>::new(
2902 &STORE_REGISTRY,
2903 &crate::db::RequestExecutionRoot::__new_runtime_root(),
2904 );
2905 session
2906 .db
2907 .ensure_recovered_state()
2908 .expect("identity pre-key test database should initialize");
2909 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2910 STORE_PATH,
2911 AcceptedSchemaRevision::INITIAL,
2912 BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH))]),
2913 BTreeMap::from([
2914 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2915 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2916 ]),
2917 );
2918 let store = session
2919 .db
2920 .store_handle(STORE_PATH)
2921 .expect("identity pre-key test store should resolve");
2922 crate::db::commit::publish_accepted_schema_candidate(
2923 STORE_PATH,
2924 store,
2925 AcceptedSchemaRevision::NONE,
2926 &candidate,
2927 )
2928 .expect("identity candidate should publish with explicit zero state");
2929 session
2930 }
2931
2932 fn initialize_journaled_with_root() -> (
2933 DbSession<JournaledTestCanister>,
2934 crate::db::RequestExecutionRoot,
2935 ) {
2936 let root = crate::db::RequestExecutionRoot::__new_runtime_root();
2937 let session = DbSession::<JournaledTestCanister>::new(&JOURNALED_STORE_REGISTRY, &root);
2938 session
2939 .db
2940 .ensure_recovered_state()
2941 .expect("journaled identity database should initialize");
2942 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2943 JOURNALED_STORE_PATH,
2944 AcceptedSchemaRevision::INITIAL,
2945 BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH))]),
2946 BTreeMap::from([
2947 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2948 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2949 ]),
2950 );
2951 let store = session
2952 .db
2953 .store_handle(JOURNALED_STORE_PATH)
2954 .expect("journaled identity store should resolve");
2955 crate::db::commit::publish_accepted_schema_candidate(
2956 JOURNALED_STORE_PATH,
2957 store,
2958 AcceptedSchemaRevision::NONE,
2959 &candidate,
2960 )
2961 .expect("journaled identity candidate should publish");
2962 (session, root)
2963 }
2964
2965 fn initialize_journaled() -> DbSession<JournaledTestCanister> {
2966 initialize_journaled_with_root().0
2967 }
2968
2969 fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
2970 AcceptedMutationIntentPatch::new()
2971 .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
2972 }
2973
2974 fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
2975 DynamicStructuralPatch::new(vec![(
2976 "payload".to_string(),
2977 DynamicWriteCell::Value(InputValue::Nat64(value)),
2978 )])
2979 }
2980
2981 fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
2982 vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
2983 }
2984
2985 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2986 fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
2987 session
2988 .issue_typed_entity_binding(
2989 ENTITY_SOURCE,
2990 &[
2991 DynamicTypedFieldBindingRequest::new(
2992 ID_SOURCE.to_string(),
2993 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2994 false,
2995 ),
2996 DynamicTypedFieldBindingRequest::new(
2997 PAYLOAD_SOURCE.to_string(),
2998 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2999 false,
3000 ),
3001 ],
3002 )
3003 .expect("exact-key test binding should issue")
3004 }
3005
3006 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3007 fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
3008 let output = session
3009 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3010 entity: ENTITY_NAME.to_string(),
3011 patch: dynamic_payload_patch(payload),
3012 })
3013 .expect("exact-key fixture insert should commit");
3014 match output.rows.as_slice() {
3015 [row] => match row.as_slice() {
3016 [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
3017 if *actual_payload == payload =>
3018 {
3019 *id
3020 }
3021 _ => panic!("exact-key fixture should return its identity and payload"),
3022 },
3023 _ => panic!("exact-key fixture insert should return one row"),
3024 }
3025 }
3026
3027 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3028 fn identity_row_stored_bytes<C: CanisterKind>(
3029 session: &DbSession<C>,
3030 store_path: &'static str,
3031 key: u64,
3032 ) -> u64 {
3033 let data_key = DecodedDataStoreKey::try_from_structural_key(ENTITY_TAG, &Value::Nat64(key))
3034 .expect("identity row key should encode");
3035 let raw_key = data_key.to_raw().expect("identity raw key should encode");
3036 let store = session
3037 .db
3038 .recovered_store(store_path)
3039 .expect("identity store should resolve");
3040 store.with_data(|data_store| {
3041 u64::try_from(
3042 data_store
3043 .get(&raw_key)
3044 .expect("inserted identity row should exist")
3045 .len(),
3046 )
3047 .expect("bounded row length should fit u64")
3048 })
3049 }
3050
3051 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3052 fn with_stored_bytes_limit<T>(
3053 limit: u64,
3054 shape_fingerprint_prefix: u64,
3055 operation: impl FnOnce() -> Result<T, crate::db::query::intent::QueryError>,
3056 ) -> Result<T, crate::db::query::intent::QueryError> {
3057 let budget = HardExecutionBudget::uniform_for_tests(
3058 u64::MAX,
3059 HardExecutionFailureHeadroom::new(500, 256),
3060 )
3061 .with_limit_for_tests(
3062 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::StoredBytesRead,
3063 limit,
3064 );
3065 let context = HardExecutionContext::new(
3066 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3067 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3068 shape_fingerprint_prefix,
3069 );
3070
3071 with_query_execution_budget_for_tests(budget, context, operation)
3072 }
3073
3074 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3075 fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
3076 let first = insert_exact_key_fixture(session, 41);
3077 let second = insert_exact_key_fixture(session, 42);
3078 let missing = u64::MAX;
3079 let binding = exact_key_binding(session);
3080 let gets_before = DataStore::current_get_call_count();
3081 let result = session
3082 .execute_public_exact_key_batch_for_typed_binding(
3083 &binding,
3084 &[second, missing, first, second],
3085 )
3086 .expect("exact-key batch should execute")
3087 .expect("exact-key binding should remain current");
3088
3089 assert_eq!(result.positions, vec![0, 1, 2, 0]);
3090 assert_eq!(
3091 result.distinct_rows,
3092 vec![
3093 Some(expected_dynamic_row(second, 42)),
3094 None,
3095 Some(expected_dynamic_row(first, 41)),
3096 ],
3097 );
3098 assert_eq!(
3099 DataStore::current_get_call_count().saturating_sub(gets_before),
3100 3,
3101 "four input positions with one duplicate must perform three physical reads",
3102 );
3103 }
3104
3105 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3106 #[test]
3107 fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
3108 assert_exact_key_batch(&initialize());
3109 assert_exact_key_batch(&initialize_journaled());
3110 }
3111
3112 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3113 fn assert_primary_range_materialization_fetches_once<C: CanisterKind>(
3114 session: &DbSession<C>,
3115 store_path: &'static str,
3116 ) {
3117 let key = insert_exact_key_fixture(session, 41);
3118 let stored_bytes = identity_row_stored_bytes(session, store_path, key);
3119
3120 let scalar = DynamicQuery::new(ENTITY_NAME)
3121 .select(["id", "payload"])
3122 .order_by(asc("id"))
3123 .limit(1);
3124 let gets_before = DataStore::current_get_call_count();
3125 let scalar_page = with_stored_bytes_limit(stored_bytes, 0x7072_696d_6172_792d, || {
3126 session.execute_trusted_live_page(&scalar, None)
3127 })
3128 .expect("one scalar primary-range row should fit one payload-read allowance");
3129 assert_eq!(scalar_page.row_count, 1);
3130 assert_eq!(
3131 DataStore::current_get_call_count().saturating_sub(gets_before),
3132 1,
3133 "scalar primary traversal should fetch its emitted row exactly once",
3134 );
3135
3136 let grouped = DynamicQuery::new(ENTITY_NAME)
3137 .group_by("payload")
3138 .aggregate(crate::db::count())
3139 .grouped_limits(10, 16 * 1_024)
3140 .limit(1);
3141 let gets_before = DataStore::current_get_call_count();
3142 let grouped_page = with_stored_bytes_limit(stored_bytes, 0x6772_6f75_7065_642d, || {
3143 session.execute_trusted_dynamic_grouped_query(&grouped)
3144 })
3145 .expect("one grouped primary-range row should fit one payload-read allowance");
3146 assert_eq!(grouped_page.row_count, 1);
3147 assert_eq!(
3148 DataStore::current_get_call_count().saturating_sub(gets_before),
3149 1,
3150 "grouped primary traversal should fetch its source row exactly once",
3151 );
3152 }
3153
3154 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3155 #[test]
3156 fn row_materialization_fetches_each_required_payload_at_most_once() {
3157 assert_primary_range_materialization_fetches_once(&initialize(), STORE_PATH);
3158 assert_primary_range_materialization_fetches_once(
3159 &initialize_journaled(),
3160 JOURNALED_STORE_PATH,
3161 );
3162 }
3163
3164 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3165 #[test]
3166 fn exhaustive_pages_require_and_recompare_the_complete_source_proof() {
3167 let session = initialize();
3168 let first = insert_exact_key_fixture(&session, 41);
3169 let second = insert_exact_key_fixture(&session, 42);
3170 let third = insert_exact_key_fixture(&session, 43);
3171 let query = DynamicQuery::new(ENTITY_NAME)
3172 .select(["id", "payload"])
3173 .order_by(asc("id"));
3174
3175 let page = session
3176 .execute_trusted_exhaustive_page(&query, None, None)
3177 .expect("initial exhaustive page should capture its source proof");
3178 assert_eq!(
3179 page.rows,
3180 vec![
3181 expected_dynamic_row(first, 41),
3182 expected_dynamic_row(second, 42),
3183 ],
3184 );
3185 let continuation = page
3186 .continuation
3187 .as_deref()
3188 .expect("unreturned row should retain exhaustive continuation");
3189 assert!(matches!(
3190 session.execute_trusted_exhaustive_page(&query, Some(continuation), None),
3191 Err(ExhaustiveReadError::Revision(
3192 ReadSetRevisionError::ResumeProofRequired
3193 )),
3194 ));
3195 let resumed = session
3196 .execute_trusted_exhaustive_page(&query, Some(continuation), Some(&page.proof))
3197 .expect("unchanged proof should resume exhaustive traversal");
3198 assert_eq!(resumed.rows, vec![expected_dynamic_row(third, 43)]);
3199 assert_eq!(resumed.continuation, None);
3200
3201 let stale_page = session
3202 .execute_trusted_exhaustive_page(&query, None, None)
3203 .expect("fresh exhaustive page should capture current revision");
3204 let stale_continuation = stale_page
3205 .continuation
3206 .as_deref()
3207 .expect("fresh three-row traversal should retain continuation");
3208 let _ = insert_exact_key_fixture(&session, 44);
3209 assert!(matches!(
3210 session.execute_trusted_exhaustive_page(
3211 &query,
3212 Some(stale_continuation),
3213 Some(&stale_page.proof),
3214 ),
3215 Err(ExhaustiveReadError::Revision(
3216 ReadSetRevisionError::StoreDataChanged { .. }
3217 )),
3218 ));
3219 }
3220
3221 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3222 #[test]
3223 fn heap_sources_cannot_back_durable_resumable_jobs() {
3224 let session = initialize();
3225 let proof = session
3226 .capture_read_set_revision_proof(&[ENTITY_NAME])
3227 .expect("heap source proof should capture for one-call exhaustive reads");
3228 let job_id = ResumableJobId::try_from_bytes([70; 32])
3229 .expect("nonzero heap test job identity should admit");
3230
3231 assert!(matches!(
3232 session.start_resumable_job(job_id, proof, Vec::new()),
3233 Err(ResumableJobError::SourceProof(
3234 ReadSetRevisionError::DurableStoreRequired { .. }
3235 )),
3236 ));
3237 }
3238
3239 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3240 #[test]
3241 fn proof_and_progress_controls_charge_one_shared_request_scope() {
3242 let (session, root) = initialize_journaled_with_root();
3243 let resource = icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions;
3244 let before = root.observed(resource);
3245 let proof = session
3246 .capture_read_set_revision_proof(&[ENTITY_NAME])
3247 .expect("proof capture should use the retained request scope");
3248 let job_id = ResumableJobId::try_from_bytes([75; 32])
3249 .expect("nonzero accounting job identity should admit");
3250 session
3251 .start_resumable_job(job_id, proof, Vec::new())
3252 .expect("job start should use the same retained request scope");
3253 let _ = session
3254 .resumable_job_state(job_id)
3255 .expect("job load should use the same retained request scope");
3256
3257 assert_eq!(root.observed(resource).saturating_sub(before), 3);
3258 }
3259
3260 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3261 #[test]
3262 fn source_proofs_ignore_unrelated_stores_but_bind_access_state_changes() {
3263 let session = initialize();
3264 let proof = session
3265 .capture_read_set_revision_proof(&[ENTITY_NAME])
3266 .expect("source proof should cover only the entity's physical store");
3267 let shared_store_proof = session
3268 .capture_read_set_revision_proof(&[ENTITY_NAME, ENTITY_NAME])
3269 .expect("entities sharing one physical source should deduplicate");
3270 assert_eq!(shared_store_proof, proof);
3271 assert_eq!(shared_store_proof.stores().len(), 1);
3272 let unrelated = session
3273 .db
3274 .store_handle(UNRELATED_STORE_PATH)
3275 .expect("unrelated registered store should resolve");
3276 unrelated.with_data_mut(|store| {
3277 let _ = store.remove(&RawDataStoreKey::from_persisted_bytes(vec![1]));
3278 });
3279 session
3280 .verify_read_set_revision_proof(&proof)
3281 .expect("a nonparticipating store mutation must not invalidate the proof");
3282
3283 let source = session
3284 .db
3285 .store_handle(STORE_PATH)
3286 .expect("participating source store should resolve");
3287 source
3288 .mark_index_building()
3289 .expect("source access-state transition should advance its revision");
3290 assert!(matches!(
3291 session.verify_read_set_revision_proof(&proof),
3292 Err(ExhaustiveReadError::Revision(
3293 ReadSetRevisionError::StoreAccessChanged { .. }
3294 )),
3295 ));
3296 }
3297
3298 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3299 #[expect(
3300 clippy::too_many_lines,
3301 reason = "one lifecycle test proves successful replay plus pre-page and post-page source invalidation without sharing progress state across tests"
3302 )]
3303 #[test]
3304 fn journaled_job_advance_is_idempotent_and_revision_checked_on_both_sides() {
3305 let session = initialize_journaled();
3306 let proof = session
3307 .capture_read_set_revision_proof(&[ENTITY_NAME])
3308 .expect("journaled source proof should capture");
3309 let job_id =
3310 ResumableJobId::try_from_bytes([71; 32]).expect("nonzero job identity should admit");
3311 session
3312 .start_resumable_job(job_id, proof, vec![0])
3313 .expect("journaled job should start outside its protected source revision");
3314 let request = ResumableJobAdvanceRequest::new(
3315 job_id,
3316 0,
3317 ResumableJobIdempotencyKey::new("page-0")
3318 .expect("bounded idempotency key should admit"),
3319 );
3320 let calls = Cell::new(0_u8);
3321 let receipt = session
3322 .compare_proof_and_advance(&request, |state| {
3323 calls.set(calls.get() + 1);
3324 assert_eq!(state.application_state, vec![0]);
3325 Ok::<_, ()>(
3326 ResumableJobAdvance::new(Some("cursor-1".to_string()), vec![1], vec![9])
3327 .expect("bounded application advance should admit"),
3328 )
3329 })
3330 .expect("unchanged source should advance exactly once");
3331 assert_eq!(calls.get(), 1);
3332 assert_eq!(receipt.status, ResumableJobAdvanceStatus::Advanced);
3333 assert_eq!(receipt.committed_sequence, 1);
3334
3335 let replay = session
3336 .compare_proof_and_advance::<()>(&request, |_| {
3337 panic!("lost-response replay must not execute application work")
3338 })
3339 .expect("same request identity should return its persisted receipt");
3340 assert_eq!(replay, receipt);
3341 let retained = session
3342 .resumable_job_state(job_id)
3343 .expect("advanced state should remain durable");
3344 assert_eq!(retained.sequence, 1);
3345 assert_eq!(retained.application_state, vec![1]);
3346
3347 let _ = insert_exact_key_fixture(&session, 51);
3348 let pre_change_request = ResumableJobAdvanceRequest::new(
3349 job_id,
3350 1,
3351 ResumableJobIdempotencyKey::new("page-1")
3352 .expect("bounded idempotency key should admit"),
3353 );
3354 let pre_change_calls = Cell::new(0_u8);
3355 let invalidated = session
3356 .compare_proof_and_advance::<()>(&pre_change_request, |_| {
3357 pre_change_calls.set(pre_change_calls.get() + 1);
3358 unreachable!("pre-page proof failure must reject before application work")
3359 })
3360 .expect("source drift should persist one replayable invalidation receipt");
3361 assert_eq!(pre_change_calls.get(), 0);
3362 assert_eq!(invalidated.status, ResumableJobAdvanceStatus::Invalidated);
3363 let invalidated_state = session
3364 .resumable_job_state(job_id)
3365 .expect("invalidated job should remain inspectable");
3366 assert_eq!(invalidated_state.status, ResumableJobStatus::Invalidated);
3367 assert_eq!(invalidated_state.continuation, None);
3368 assert_eq!(invalidated_state.application_state, vec![1]);
3369 assert_eq!(
3370 session
3371 .compare_proof_and_advance::<()>(&pre_change_request, |_| {
3372 panic!("invalidation replay must not execute application work")
3373 })
3374 .expect("lost invalidation reply should replay exactly"),
3375 invalidated,
3376 );
3377
3378 let post_proof = session
3379 .capture_read_set_revision_proof(&[ENTITY_NAME])
3380 .expect("post-change journaled proof should capture");
3381 let post_job_id = ResumableJobId::try_from_bytes([72; 32])
3382 .expect("nonzero post-change job identity should admit");
3383 session
3384 .start_resumable_job(post_job_id, post_proof, vec![7])
3385 .expect("post-change journaled job should start");
3386 let post_request = ResumableJobAdvanceRequest::new(
3387 post_job_id,
3388 0,
3389 ResumableJobIdempotencyKey::new("post-page-0")
3390 .expect("bounded idempotency key should admit"),
3391 );
3392 let post_receipt = session
3393 .compare_proof_and_advance::<()>(&post_request, |_| {
3394 let _ = insert_exact_key_fixture(&session, 52);
3395 Ok(ResumableJobAdvance::new(None, vec![8], vec![10])
3396 .expect("bounded post-change candidate should admit"))
3397 })
3398 .expect("post-page drift should discard the candidate and persist invalidation");
3399 assert_eq!(post_receipt.status, ResumableJobAdvanceStatus::Invalidated);
3400 let post_state = session
3401 .resumable_job_state(post_job_id)
3402 .expect("post-page invalidation should remain inspectable");
3403 assert_eq!(post_state.status, ResumableJobStatus::Invalidated);
3404 assert_eq!(post_state.application_state, vec![7]);
3405 session
3406 .acknowledge_resumable_job(post_job_id, post_state.sequence)
3407 .expect("terminal job acknowledgement should remove retained progress");
3408 session
3409 .acknowledge_resumable_job(post_job_id, post_state.sequence)
3410 .expect("lost acknowledgement reply should be safely replayable");
3411 assert_eq!(
3412 session.resumable_job_state(post_job_id),
3413 Err(ResumableJobError::NotFound),
3414 );
3415
3416 let completed_job_id = ResumableJobId::try_from_bytes([74; 32])
3417 .expect("nonzero completed job identity should admit");
3418 let completed_proof = session
3419 .capture_read_set_revision_proof(&[ENTITY_NAME])
3420 .expect("completed-job source proof should capture");
3421 session
3422 .start_resumable_job(completed_job_id, completed_proof, Vec::new())
3423 .expect("completed-job fixture should start");
3424 let completed_request = ResumableJobAdvanceRequest::new(
3425 completed_job_id,
3426 0,
3427 ResumableJobIdempotencyKey::new("complete")
3428 .expect("bounded completion key should admit"),
3429 );
3430 let completed_receipt = session
3431 .compare_proof_and_advance::<()>(&completed_request, |_| {
3432 Ok(ResumableJobAdvance::new(None, vec![99], vec![100])
3433 .expect("bounded terminal advance should admit"))
3434 })
3435 .expect("null continuation should commit terminal completion");
3436 let completed_state = session
3437 .resumable_job_state(completed_job_id)
3438 .expect("completed state should remain replayable before acknowledgement");
3439 assert_eq!(completed_state.status, ResumableJobStatus::Completed);
3440 assert_eq!(
3441 session
3442 .compare_proof_and_advance::<()>(&completed_request, |_| {
3443 panic!("completed request replay must not execute application work")
3444 })
3445 .expect("completed request should replay until acknowledgement"),
3446 completed_receipt,
3447 );
3448 let after_completion = ResumableJobAdvanceRequest::new(
3449 completed_job_id,
3450 1,
3451 ResumableJobIdempotencyKey::new("after-complete")
3452 .expect("bounded post-completion key should admit"),
3453 );
3454 assert!(matches!(
3455 session.compare_proof_and_advance::<()>(&after_completion, |_| {
3456 panic!("completed jobs cannot execute another page")
3457 }),
3458 Err(CompareProofAndAdvanceError::Protocol(
3459 ResumableJobError::Completed
3460 )),
3461 ));
3462 session
3463 .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
3464 .expect("completed job should acknowledge and free capacity");
3465 session
3466 .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
3467 .expect("completion acknowledgement should be idempotent");
3468
3469 let stale_job_id = ResumableJobId::try_from_bytes([73; 32])
3470 .expect("nonzero stale-sequence job identity should admit");
3471 let stale_proof = session
3472 .capture_read_set_revision_proof(&[ENTITY_NAME])
3473 .expect("stale-sequence source proof should capture");
3474 session
3475 .start_resumable_job(stale_job_id, stale_proof, Vec::new())
3476 .expect("stale-sequence job should start");
3477 let stale_request = ResumableJobAdvanceRequest::new(
3478 stale_job_id,
3479 4,
3480 ResumableJobIdempotencyKey::new("stale").expect("bounded idempotency key should admit"),
3481 );
3482 assert!(matches!(
3483 session.compare_proof_and_advance::<()>(&stale_request, |_| {
3484 panic!("stale sequence must reject before application work")
3485 }),
3486 Err(CompareProofAndAdvanceError::Protocol(
3487 ResumableJobError::StaleSequence {
3488 expected: 4,
3489 actual: 0,
3490 }
3491 )),
3492 ));
3493 assert_eq!(
3494 session.acknowledge_resumable_job(stale_job_id, 0),
3495 Err(ResumableJobError::NotTerminal),
3496 );
3497 }
3498
3499 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3500 #[test]
3501 fn exact_key_batch_uses_typed_hard_execution_budget() {
3502 let session = initialize();
3503 let binding = exact_key_binding(&session);
3504 let budget =
3505 HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
3506 let error = session
3507 .execute_exact_key_batch_with_hard_budget_for_tests(&binding, &[u64::MAX], &budget)
3508 .expect_err("zero query budget should reject the exact-key route");
3509
3510 assert!(matches!(
3511 error.diagnostic().detail(),
3512 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3513 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3514 })
3515 ));
3516 let facts = error.diagnostic_facts();
3517 assert_eq!(
3518 &facts[..5],
3519 &[
3520 (
3521 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3522 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions.raw(),
3523 ),
3524 (icydb_diagnostic_code::DiagnosticFactTag::Limit, 0),
3525 (icydb_diagnostic_code::DiagnosticFactTag::Actual, 1),
3526 (
3527 icydb_diagnostic_code::DiagnosticFactTag::ExecutionBudgetScope,
3528 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution.raw(),
3529 ),
3530 (
3531 icydb_diagnostic_code::DiagnosticFactTag::ExecutionLane,
3532 icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead.raw(),
3533 ),
3534 ],
3535 );
3536 assert_eq!(
3537 facts[5].0,
3538 icydb_diagnostic_code::DiagnosticFactTag::QueryShapeFingerprintPrefix,
3539 );
3540 assert_ne!(facts[5].1, 0);
3541 }
3542
3543 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3544 fn assert_planned_query_exhausts(
3545 session: &DbSession<TestCanister>,
3546 query: &crate::db::DynamicQuery,
3547 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3548 ) {
3549 let budget = HardExecutionBudget::uniform_for_tests(
3550 u64::MAX,
3551 HardExecutionFailureHeadroom::new(500, 256),
3552 )
3553 .with_limit_for_tests(resource, 0);
3554 let context = HardExecutionContext::new(
3555 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3556 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3557 0x7068_7973_6963_616c,
3558 );
3559 let error = with_query_execution_budget_for_tests(budget, context, || {
3560 session.execute_trusted_live_page(query, None)
3561 })
3562 .expect_err("the injected zero resource allowance should reject planned execution");
3563
3564 assert!(matches!(
3565 error.diagnostic().detail(),
3566 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3567 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3568 })
3569 ));
3570 assert_eq!(
3571 error.diagnostic_facts()[0],
3572 (
3573 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3574 resource.raw(),
3575 ),
3576 );
3577 }
3578
3579 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3580 fn assert_grouped_query_exhausts(
3581 session: &DbSession<TestCanister>,
3582 query: &crate::db::DynamicQuery,
3583 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3584 ) {
3585 let budget = HardExecutionBudget::uniform_for_tests(
3586 u64::MAX,
3587 HardExecutionFailureHeadroom::new(500, 256),
3588 )
3589 .with_limit_for_tests(resource, 0);
3590 let context = HardExecutionContext::new(
3591 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3592 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3593 0x6772_6f75_7065_642d,
3594 );
3595 let error = with_query_execution_budget_for_tests(budget, context, || {
3596 session.execute_trusted_dynamic_grouped_query(query)
3597 })
3598 .expect_err("the injected zero resource allowance should reject grouped execution");
3599
3600 assert!(matches!(
3601 error.diagnostic().detail(),
3602 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3603 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3604 })
3605 ));
3606 assert_eq!(
3607 error.diagnostic_facts()[0],
3608 (
3609 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3610 resource.raw(),
3611 ),
3612 );
3613 }
3614
3615 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3616 fn assert_sql_query_exhausts(
3617 session: &DbSession<TestCanister>,
3618 sql: &str,
3619 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3620 ) {
3621 let budget = HardExecutionBudget::uniform_for_tests(
3622 u64::MAX,
3623 HardExecutionFailureHeadroom::new(500, 256),
3624 )
3625 .with_limit_for_tests(resource, 0);
3626 let context = HardExecutionContext::new(
3627 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3628 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3629 0x7371_6c2d_736f_7274,
3630 );
3631 let error = with_query_execution_budget_for_tests(budget, context, || {
3632 session.execute_trusted_sql_query(sql)
3633 })
3634 .expect_err("the injected zero resource allowance should reject SQL execution");
3635
3636 assert!(matches!(
3637 error.diagnostic().detail(),
3638 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3639 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3640 })
3641 ));
3642 assert_eq!(
3643 error.diagnostic_facts()[0],
3644 (
3645 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3646 resource.raw(),
3647 ),
3648 );
3649 }
3650
3651 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3652 #[test]
3653 fn planned_read_routes_share_physical_resource_accounting() {
3654 let session = initialize();
3655 let first = insert_exact_key_fixture(&session, 41);
3656 insert_exact_key_fixture(&session, 42);
3657
3658 let fallback = crate::db::DynamicQuery::new(ENTITY_NAME)
3659 .filter(crate::db::FieldRef::new("id").eq(first))
3660 .select(["id", "payload"])
3661 .order_by(crate::db::asc("id"))
3662 .limit(1);
3663 assert_eq!(
3664 session
3665 .execute_trusted_live_page(&fallback, None)
3666 .expect("bounded fallback execution should preserve its result")
3667 .row_count,
3668 1,
3669 );
3670 assert_planned_query_exhausts(
3671 &session,
3672 &fallback,
3673 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::RowsVisited,
3674 );
3675
3676 let covering = crate::db::DynamicQuery::new(ENTITY_NAME)
3677 .filter(crate::db::FieldRef::new("payload").eq(41_u64))
3678 .select(["payload"])
3679 .order_by(crate::db::asc("payload"))
3680 .limit(1);
3681 assert_eq!(
3682 session
3683 .execute_trusted_live_page(&covering, None)
3684 .expect("bounded covering execution should preserve its result")
3685 .row_count,
3686 1,
3687 );
3688 assert_planned_query_exhausts(
3689 &session,
3690 &covering,
3691 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
3692 );
3693
3694 let residual = crate::db::DynamicQuery::new(ENTITY_NAME)
3695 .filter(crate::db::FieldRef::new("payload").eq_field("id"))
3696 .select(["id"])
3697 .order_by(crate::db::asc("id"))
3698 .limit(1);
3699 assert_eq!(
3700 session
3701 .execute_trusted_live_page(&residual, None)
3702 .expect("bounded residual execution should preserve its result")
3703 .row_count,
3704 0,
3705 );
3706 assert_planned_query_exhausts(
3707 &session,
3708 &residual,
3709 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::PredicateExpressionSteps,
3710 );
3711
3712 assert_planned_query_exhausts(
3713 &session,
3714 &fallback,
3715 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::ResultBytes,
3716 );
3717
3718 let grouped = crate::db::DynamicQuery::new(ENTITY_NAME)
3719 .group_by("payload")
3720 .aggregate(crate::db::count())
3721 .order_by(crate::db::asc("payload"))
3722 .grouped_limits(10, 16 * 1_024)
3723 .limit(1);
3724 let grouped_result = session
3725 .execute_trusted_dynamic_grouped_query(&grouped)
3726 .expect("bounded grouped execution should preserve its result");
3727 assert_eq!(grouped_result.row_count, 1);
3728 assert!(grouped_result.next_cursor.is_some());
3729 assert_grouped_query_exhausts(
3730 &session,
3731 &grouped,
3732 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::GroupDistinctEntries,
3733 );
3734 assert_grouped_query_exhausts(
3735 &session,
3736 &grouped,
3737 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::CursorSteps,
3738 );
3739
3740 assert_sql_query_exhausts(
3741 &session,
3742 "SELECT payload, COUNT(*) AS row_count FROM IdentityRow \
3743 GROUP BY payload ORDER BY row_count DESC, payload ASC LIMIT 1",
3744 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::SortEntries,
3745 );
3746 }
3747
3748 fn assert_dynamic_payload(session: &DbSession<TestCanister>, key: u64, expected_payload: u64) {
3749 let unchanged = session
3750 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3751 entity: ENTITY_NAME.to_string(),
3752 key: InputValue::Nat64(key),
3753 patch: dynamic_payload_patch(expected_payload),
3754 })
3755 .expect("the expected row should remain readable through a no-op update");
3756 assert_eq!(unchanged.affected_rows, 0);
3757 assert_eq!(
3758 unchanged.rows,
3759 vec![expected_dynamic_row(key, expected_payload)],
3760 );
3761 }
3762
3763 fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
3764 values
3765 .iter()
3766 .map(|value| {
3767 AcceptedStructuralMutation::save(
3768 MutationMode::Insert,
3769 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
3770 payload_patch(*value),
3771 )
3772 })
3773 .collect()
3774 }
3775
3776 fn assert_identity_boundary(error: &InternalError) {
3777 assert_eq!(error.class(), ErrorClass::Unsupported);
3778 assert_eq!(error.origin(), ErrorOrigin::Identity);
3779 }
3780
3781 #[test]
3782 fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
3783 let generated = insert_key_exists_after_generation(true);
3784 assert_eq!(generated.class(), ErrorClass::Corruption);
3785 assert_eq!(generated.origin(), ErrorOrigin::Identity);
3786
3787 let ordinary = insert_key_exists_after_generation(false);
3788 assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
3789 }
3790
3791 #[cfg(target_pointer_width = "64")]
3792 #[test]
3793 fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
3794 let error = checked_pre_key_candidate_count(
3795 usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
3796 )
3797 .expect_err("candidate counts beyond u32 must reject");
3798 assert_identity_boundary(&error);
3799 }
3800
3801 #[test]
3802 #[expect(
3803 clippy::too_many_lines,
3804 reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
3805 )]
3806 fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
3807 let session = initialize();
3808 let seeded = session
3809 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3810 .expect("seed rows should commit");
3811 assert_eq!(seeded.affected_rows, 1);
3812
3813 let split = session
3814 .execute_trusted_dynamic_mutation_batch(vec![
3815 DynamicMutation::Update {
3816 entity: ENTITY_NAME.to_string(),
3817 key: InputValue::Nat64(1),
3818 patch: dynamic_payload_patch(60),
3819 },
3820 DynamicMutation::Insert {
3821 entity: ENTITY_NAME.to_string(),
3822 patch: dynamic_payload_patch(40),
3823 },
3824 ])
3825 .expect("one holding should split atomically");
3826 assert_eq!(split.affected_rows, 2);
3827 assert_eq!(
3828 split.rows,
3829 vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
3830 "split after-images must retain input order and exact quantity",
3831 );
3832
3833 let rejected_split = session
3834 .execute_trusted_dynamic_mutation_batch(vec![
3835 DynamicMutation::Update {
3836 entity: ENTITY_NAME.to_string(),
3837 key: InputValue::Nat64(1),
3838 patch: dynamic_payload_patch(50),
3839 },
3840 DynamicMutation::Insert {
3841 entity: ENTITY_NAME.to_string(),
3842 patch: DynamicStructuralPatch::new(Vec::new()),
3843 },
3844 ])
3845 .expect_err("an invalid split output must reject the staged source update");
3846 assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
3847 assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
3848 assert_eq!(
3849 rejected_split.diagnostic_facts(),
3850 vec![
3851 (
3852 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3853 ENTITY_TAG.value(),
3854 ),
3855 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
3856 (
3857 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3858 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3859 ),
3860 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3861 ],
3862 );
3863 assert_dynamic_payload(&session, 1, 60);
3864 assert_dynamic_payload(&session, 2, 40);
3865
3866 let transfer = session
3867 .execute_trusted_dynamic_mutation_batch(vec![
3868 DynamicMutation::Update {
3869 entity: ENTITY_NAME.to_string(),
3870 key: InputValue::Nat64(1),
3871 patch: dynamic_payload_patch(70),
3872 },
3873 DynamicMutation::Update {
3874 entity: ENTITY_NAME.to_string(),
3875 key: InputValue::Nat64(2),
3876 patch: dynamic_payload_patch(30),
3877 },
3878 ])
3879 .expect("distinct transfer patches should share one atomic batch");
3880 assert_eq!(
3881 transfer.rows,
3882 vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
3883 "the transfer must preserve the exact total quantity",
3884 );
3885
3886 let merge = session
3887 .execute_trusted_dynamic_mutation_batch(vec![
3888 DynamicMutation::Delete {
3889 entity: ENTITY_NAME.to_string(),
3890 key: InputValue::Nat64(2),
3891 },
3892 DynamicMutation::Update {
3893 entity: ENTITY_NAME.to_string(),
3894 key: InputValue::Nat64(1),
3895 patch: dynamic_payload_patch(100),
3896 },
3897 ])
3898 .expect("two holdings should merge atomically");
3899 assert_eq!(
3900 merge.rows,
3901 vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
3902 "delete before-images and update after-images must retain input order",
3903 );
3904
3905 let resplit = session
3906 .execute_trusted_dynamic_mutation_batch(vec![
3907 DynamicMutation::Update {
3908 entity: ENTITY_NAME.to_string(),
3909 key: InputValue::Nat64(1),
3910 patch: dynamic_payload_patch(60),
3911 },
3912 DynamicMutation::Insert {
3913 entity: ENTITY_NAME.to_string(),
3914 patch: dynamic_payload_patch(40),
3915 },
3916 ])
3917 .expect("the merged holding should split again");
3918 assert_eq!(
3919 resplit.rows,
3920 vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
3921 );
3922
3923 let rejected_merge = session
3924 .execute_trusted_dynamic_mutation_batch(vec![
3925 DynamicMutation::Delete {
3926 entity: ENTITY_NAME.to_string(),
3927 key: InputValue::Nat64(3),
3928 },
3929 DynamicMutation::Update {
3930 entity: ENTITY_NAME.to_string(),
3931 key: InputValue::Nat64(99),
3932 patch: dynamic_payload_patch(100),
3933 },
3934 ])
3935 .expect_err("a late missing merge target must preserve the earlier staged delete");
3936 assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
3937 assert_dynamic_payload(&session, 1, 60);
3938 assert_dynamic_payload(&session, 3, 40);
3939
3940 SCHEMA_STORE.with(|store| {
3941 let cursor = store
3942 .borrow()
3943 .identity_statement_cursor(
3944 database_incarnation_id().expect("database incarnation should remain readable"),
3945 ENTITY_TAG,
3946 FieldId::new(1),
3947 &AcceptedFieldKind::Nat64,
3948 )
3949 .expect("mixed Identity state should remain readable");
3950 assert_eq!(cursor.expected_high_water(), 3);
3951 assert!(!cursor.has_allocations());
3952 });
3953 }
3954
3955 #[test]
3956 fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
3957 let session = initialize();
3958 session
3959 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3960 .expect("the holding fixture should initialize");
3961
3962 let duplicate = session
3963 .execute_trusted_dynamic_mutation_batch(vec![
3964 DynamicMutation::Update {
3965 entity: ENTITY_NAME.to_string(),
3966 key: InputValue::Nat64(1),
3967 patch: dynamic_payload_patch(60),
3968 },
3969 DynamicMutation::Delete {
3970 entity: ENTITY_NAME.to_string(),
3971 key: InputValue::Nat64(1),
3972 },
3973 ])
3974 .expect_err("duplicate targets across operation kinds must reject");
3975 assert!(matches!(
3976 duplicate.diagnostic().detail(),
3977 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3978 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
3979 }),
3980 ));
3981 assert_eq!(
3982 duplicate.diagnostic_facts(),
3983 vec![
3984 (
3985 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3986 ENTITY_TAG.value(),
3987 ),
3988 (
3989 icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
3990 0,
3991 ),
3992 (
3993 icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
3994 1,
3995 ),
3996 ],
3997 );
3998 assert_dynamic_payload(&session, 1, 100);
3999 }
4000
4001 #[test]
4002 fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
4003 let session = initialize();
4004 let empty = session
4005 .execute_trusted_dynamic_mutation_batch(Vec::new())
4006 .expect_err("an empty public batch must reject");
4007 assert!(matches!(
4008 empty.diagnostic().detail(),
4009 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4010 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
4011 }),
4012 ));
4013 assert_eq!(
4014 empty.diagnostic_facts(),
4015 vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
4016 );
4017
4018 let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
4019 .map(|_| DynamicMutation::Delete {
4020 entity: ENTITY_NAME.to_string(),
4021 key: InputValue::Nat64(1),
4022 })
4023 .collect();
4024 let over_bound = session
4025 .execute_trusted_dynamic_mutation_batch(requests)
4026 .expect_err("operation cap plus one must reject before row resolution");
4027 assert!(matches!(
4028 over_bound.diagnostic().detail(),
4029 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4030 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
4031 }),
4032 ));
4033 assert_eq!(
4034 over_bound.diagnostic_facts(),
4035 vec![
4036 (
4037 icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
4038 (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
4039 ),
4040 (
4041 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4042 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
4043 ),
4044 ],
4045 );
4046 }
4047
4048 #[test]
4049 fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
4050 let mut exact = 0;
4051 add_structural_mutation_staged_bytes(
4052 &mut exact,
4053 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
4054 )
4055 .expect("the exact staged-byte boundary should admit");
4056 assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
4057
4058 let error = add_structural_mutation_staged_bytes(&mut exact, [1])
4059 .expect_err("one byte above the staged-byte boundary must reject");
4060 assert!(matches!(
4061 error.diagnostic().detail(),
4062 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4063 boundary:
4064 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
4065 }),
4066 ));
4067 assert_eq!(
4068 error.diagnostic_facts(),
4069 vec![
4070 (
4071 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4072 (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
4073 ),
4074 (
4075 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4076 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
4077 ),
4078 ],
4079 );
4080
4081 validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
4082 .expect("the exact result-byte boundary should admit");
4083 let error = validate_structural_mutation_result_bytes(
4084 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
4085 )
4086 .expect_err("one byte above the result-byte boundary must reject");
4087 assert!(matches!(
4088 error.diagnostic().detail(),
4089 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4090 boundary:
4091 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
4092 }),
4093 ));
4094 assert_eq!(
4095 error.diagnostic_facts(),
4096 vec![
4097 (
4098 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4099 (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
4100 ),
4101 (
4102 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4103 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
4104 ),
4105 ],
4106 );
4107 }
4108
4109 #[expect(
4110 clippy::too_many_lines,
4111 reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
4112 )]
4113 #[test]
4114 fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
4115 let session = initialize();
4116 let catalog = session
4117 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4118 .expect("identity catalog should resolve");
4119 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4120 .expect("identity row layout should build");
4121 let initial_description = session
4122 .try_describe_entity_by_name(ENTITY_NAME)
4123 .expect("accepted Identity description should resolve");
4124 assert_eq!(
4125 initial_description.entity_tag(),
4126 catalog.identity().entity_tag().value()
4127 );
4128 assert_eq!(
4129 initial_description.accepted_schema_fingerprint_method(),
4130 catalog.fingerprint_method_version()
4131 );
4132 assert_eq!(
4133 initial_description.accepted_schema_fingerprint(),
4134 catalog.fingerprint()
4135 );
4136 let initial_identity = initial_description
4137 .identity()
4138 .expect("accepted Identity policy should be described");
4139 assert_eq!(initial_identity.field(), "id");
4140 assert_eq!(initial_identity.generator(), "Identity::next");
4141 assert_eq!(initial_identity.accepted_kind(), "nat64");
4142 assert_eq!(initial_identity.minimum(), 1);
4143 assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
4144 assert_eq!(initial_identity.high_water(), 0);
4145 assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
4146 assert!(!initial_identity.exhausted());
4147
4148 let rejected = session
4149 .execute_accepted_structural_save_batch(
4150 &catalog,
4151 &descriptor,
4152 batch(&[1_000, 2_000]),
4153 Timestamp::from_millis(6),
4154 |_| Err::<(), _>(InternalError::executor_unsupported()),
4155 )
4156 .expect_err("a rejected precommit result must not publish its tentative range");
4157 assert_eq!(rejected.class(), ErrorClass::Unsupported);
4158 assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
4159
4160 let rows = session
4161 .execute_accepted_structural_save_batch(
4162 &catalog,
4163 &descriptor,
4164 batch(&[10, 20, 30]),
4165 Timestamp::from_millis(7),
4166 Ok,
4167 )
4168 .expect("one accepted batch should commit rows and one identity range");
4169 assert_eq!(
4170 rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
4171 vec![
4172 vec![Value::Nat64(1), Value::Nat64(10)],
4173 vec![Value::Nat64(2), Value::Nat64(20)],
4174 vec![Value::Nat64(3), Value::Nat64(30)],
4175 ],
4176 );
4177
4178 let dynamic = session
4179 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4180 entity: ENTITY_NAME.to_string(),
4181 patch: DynamicStructuralPatch::new(vec![(
4182 "payload".to_string(),
4183 DynamicWriteCell::Value(InputValue::Nat64(40)),
4184 )]),
4185 })
4186 .expect("dynamic omission should commit through shared Identity generation");
4187 assert_eq!(dynamic.affected_rows, 1);
4188
4189 for (request, operation) in [
4190 (
4191 DynamicMutation::Insert {
4192 entity: ENTITY_NAME.to_string(),
4193 patch: DynamicStructuralPatch::new(vec![
4194 (
4195 "id".to_string(),
4196 DynamicWriteCell::Value(InputValue::Nat64(41)),
4197 ),
4198 (
4199 "payload".to_string(),
4200 DynamicWriteCell::Value(InputValue::Nat64(42)),
4201 ),
4202 ]),
4203 },
4204 icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
4205 ),
4206 (
4207 DynamicMutation::Update {
4208 entity: ENTITY_NAME.to_string(),
4209 key: InputValue::Nat64(1),
4210 patch: DynamicStructuralPatch::new(vec![(
4211 "id".to_string(),
4212 DynamicWriteCell::Default,
4213 )]),
4214 },
4215 icydb_diagnostic_code::DiagnosticMutationOperation::Update,
4216 ),
4217 ] {
4218 let error = session
4219 .execute_trusted_dynamic_mutation(&request)
4220 .expect_err("structural Identity authorship and regeneration must reject");
4221 assert_eq!(error.class(), ErrorClass::Unsupported);
4222 assert_eq!(error.origin(), ErrorOrigin::Executor);
4223 assert_eq!(
4224 error.diagnostic_facts(),
4225 vec![
4226 (
4227 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4228 ENTITY_TAG.value(),
4229 ),
4230 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
4231 (
4232 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4233 operation.raw(),
4234 ),
4235 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
4236 ],
4237 );
4238 }
4239
4240 let binding = session
4241 .issue_typed_entity_binding(
4242 ENTITY_SOURCE,
4243 &[
4244 DynamicTypedFieldBindingRequest::new(
4245 ID_SOURCE.to_string(),
4246 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4247 false,
4248 ),
4249 DynamicTypedFieldBindingRequest::new(
4250 PAYLOAD_SOURCE.to_string(),
4251 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4252 false,
4253 ),
4254 ],
4255 )
4256 .expect("typed output should bind the Identity field");
4257 let typed_patch = binding
4258 .bind_write_fields(vec![(
4259 PAYLOAD_SOURCE.to_string(),
4260 DynamicWriteCell::Value(InputValue::Nat64(50)),
4261 )])
4262 .expect("typed payload should lower");
4263 let typed = session
4264 .execute_trusted_typed_mutation(
4265 &binding,
4266 &DynamicTypedMutation::Insert { patch: typed_patch },
4267 )
4268 .expect("typed omission should commit through shared Identity generation");
4269 assert_eq!(
4270 typed
4271 .expect("typed insert should return one mutation result")
4272 .affected_rows,
4273 1,
4274 );
4275 let explicit_typed_patch = binding
4276 .bind_write_fields(vec![
4277 (
4278 ID_SOURCE.to_string(),
4279 DynamicWriteCell::Value(InputValue::Nat64(51)),
4280 ),
4281 (
4282 PAYLOAD_SOURCE.to_string(),
4283 DynamicWriteCell::Value(InputValue::Nat64(52)),
4284 ),
4285 ])
4286 .expect("the low-level binding should retain exact authored intent");
4287 let explicit_typed_error = session
4288 .execute_trusted_typed_mutation(
4289 &binding,
4290 &DynamicTypedMutation::Insert {
4291 patch: explicit_typed_patch,
4292 },
4293 )
4294 .expect_err("typed Identity authorship must reject before allocation");
4295 assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
4296 assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
4297 assert_eq!(
4298 explicit_typed_error.diagnostic_facts(),
4299 vec![
4300 (
4301 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4302 ENTITY_TAG.value(),
4303 ),
4304 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
4305 (
4306 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4307 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4308 ),
4309 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
4310 ],
4311 );
4312
4313 let replace_error = session
4314 .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
4315 entity: ENTITY_NAME.to_string(),
4316 key: InputValue::Nat64(99),
4317 patch: DynamicStructuralPatch::new(vec![(
4318 "payload".to_string(),
4319 DynamicWriteCell::Value(InputValue::Nat64(60)),
4320 )]),
4321 })
4322 .expect_err("save-as-insert with a chosen Identity must reject");
4323 assert_eq!(replace_error.class(), ErrorClass::Unsupported);
4324 assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
4325
4326 #[cfg(feature = "sql")]
4327 {
4328 for sql in [
4329 "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
4330 "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
4331 ] {
4332 let _result = session
4333 .execute_trusted_sql_mutation(sql)
4334 .expect("SQL omission and DEFAULT should commit Identity generation");
4335 }
4336
4337 let error = session
4338 .execute_trusted_sql_mutation(
4339 "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
4340 )
4341 .expect_err("an explicit SQL Identity value must reject before allocation");
4342 let diagnostic = error.diagnostic();
4343 assert_eq!(
4344 diagnostic.code(),
4345 icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
4346 );
4347 assert!(matches!(
4348 diagnostic.detail(),
4349 Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
4350 boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
4351 }),
4352 ));
4353 }
4354
4355 let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
4356 assert_eq!(
4357 DATA_STORE.with(|store| store.borrow().len()),
4358 expected_committed
4359 );
4360 SCHEMA_STORE.with(|store| {
4361 let cursor = store
4362 .borrow()
4363 .identity_statement_cursor(
4364 database_incarnation_id().expect("database incarnation should remain readable"),
4365 ENTITY_TAG,
4366 FieldId::new(1),
4367 &AcceptedFieldKind::Nat64,
4368 )
4369 .expect("committed writes must leave active state readable");
4370 assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
4371 assert!(!cursor.has_allocations());
4372 });
4373 let committed_description = session
4374 .try_describe_entity_by_name(ENTITY_NAME)
4375 .expect("committed Identity description should resolve");
4376 let committed_identity = committed_description
4377 .identity()
4378 .expect("accepted Identity policy should remain described");
4379 assert_eq!(
4380 committed_identity.high_water(),
4381 u128::from(expected_committed),
4382 );
4383 assert_eq!(
4384 committed_identity.remaining(),
4385 u128::from(u64::MAX - expected_committed),
4386 );
4387 assert!(!committed_identity.exhausted());
4388 }
4389
4390 #[test]
4391 #[expect(
4392 clippy::too_many_lines,
4393 reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
4394 )]
4395 fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
4396 let session = initialize_journaled();
4397 let catalog = session
4398 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4399 .expect("journaled identity catalog should resolve");
4400 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4401 .expect("journaled identity row layout should build");
4402
4403 for (ordinal, interruption) in [
4404 MutationCommitInterruption::MarkerPersisted,
4405 MutationCommitInterruption::JournalPublished,
4406 MutationCommitInterruption::RowsPublished,
4407 MutationCommitInterruption::StateMaterialized,
4408 ]
4409 .into_iter()
4410 .enumerate()
4411 {
4412 interrupt_next_mutation_commit_for_tests(interruption);
4413 let interrupted = session.execute_accepted_structural_save_batch(
4414 &catalog,
4415 &descriptor,
4416 batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
4417 Timestamp::from_millis(8),
4418 Ok,
4419 );
4420 assert!(
4421 interrupted.is_err(),
4422 "the selected durable boundary should interrupt",
4423 );
4424
4425 let committed = session
4426 .execute_accepted_structural_save_batch(
4427 &catalog,
4428 &descriptor,
4429 batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
4430 Timestamp::from_millis(9),
4431 Ok,
4432 )
4433 .expect("the next mutation must recover before allocating");
4434 let expected_high_water =
4435 u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
4436 assert_eq!(
4437 committed
4438 .into_iter()
4439 .map(|row| row.values)
4440 .collect::<Vec<_>>(),
4441 vec![vec![
4442 Value::Nat64(expected_high_water),
4443 Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
4444 ]],
4445 );
4446 assert_eq!(
4447 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
4448 expected_high_water,
4449 );
4450 JOURNALED_SCHEMA_STORE.with(|store| {
4451 let cursor = store
4452 .borrow()
4453 .identity_statement_cursor(
4454 database_incarnation_id()
4455 .expect("database incarnation should remain readable"),
4456 ENTITY_TAG,
4457 FieldId::new(1),
4458 &AcceptedFieldKind::Nat64,
4459 )
4460 .expect("guarded recovery must leave quiescent active state");
4461 assert_eq!(
4462 cursor.expected_high_water(),
4463 u128::from(expected_high_water),
4464 );
4465 assert!(!cursor.has_allocations());
4466 });
4467 }
4468
4469 for (ordinal, (interruption, deleted_key)) in [
4470 (MutationCommitInterruption::MarkerPersisted, 2),
4471 (MutationCommitInterruption::JournalPublished, 4),
4472 (MutationCommitInterruption::RowPrefixPublished, 6),
4473 (MutationCommitInterruption::RowsPublished, 8),
4474 (MutationCommitInterruption::StateMaterialized, 7),
4475 ]
4476 .into_iter()
4477 .enumerate()
4478 {
4479 let expected_payload =
4480 501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
4481 interrupt_next_mutation_commit_for_tests(interruption);
4482 let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
4483 DynamicMutation::Update {
4484 entity: ENTITY_NAME.to_string(),
4485 key: InputValue::Nat64(1),
4486 patch: dynamic_payload_patch(expected_payload),
4487 },
4488 DynamicMutation::Delete {
4489 entity: ENTITY_NAME.to_string(),
4490 key: InputValue::Nat64(deleted_key),
4491 },
4492 ]);
4493 assert!(
4494 interrupted.is_err(),
4495 "the selected caller-key mixed publication boundary should interrupt",
4496 );
4497 let recovered_update = session
4498 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
4499 entity: ENTITY_NAME.to_string(),
4500 key: InputValue::Nat64(1),
4501 patch: dynamic_payload_patch(expected_payload),
4502 })
4503 .expect("guarded reentry should complete the marker-authorized mixed batch");
4504 assert_eq!(
4505 recovered_update.affected_rows, 0,
4506 "the recovered update must already expose its admitted final image",
4507 );
4508 let recovered_delete = session
4509 .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
4510 entity: ENTITY_NAME.to_string(),
4511 key: InputValue::Nat64(deleted_key),
4512 })
4513 .expect_err("the recovered delete must already be materialized");
4514 assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
4515 JOURNALED_SCHEMA_STORE.with(|store| {
4516 let cursor = store
4517 .borrow()
4518 .identity_statement_cursor(
4519 database_incarnation_id()
4520 .expect("database incarnation should remain readable"),
4521 ENTITY_TAG,
4522 FieldId::new(1),
4523 &AcceptedFieldKind::Nat64,
4524 )
4525 .expect("caller-key recovery must preserve active Identity state");
4526 assert_eq!(cursor.expected_high_water(), 8);
4527 assert!(!cursor.has_allocations());
4528 });
4529 }
4530
4531 forget_recovered_domain_for_tests(&session.db)
4532 .expect("the final journal tail should remain recoverable");
4533 session
4534 .db
4535 .ensure_recovered_state()
4536 .expect("derived rebuild must not allocate another identity");
4537
4538 let quick = execute_quick_integrity(&session.db, catalog.inspection_plan())
4539 .expect("quiescent Identity control inventory should be inspectable");
4540 assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
4541 let row_page = execute_row_integrity_page(
4542 &session.db,
4543 catalog.inspection_plan(),
4544 PhysicalUnitCheckpoint::BeforeFirst,
4545 RowInspectionLimits::standard(),
4546 )
4547 .expect("Identity rows should remain within committed high-water");
4548 assert!(row_page.exhausted());
4549 assert!(row_page.findings().is_empty());
4550
4551 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
4552 assert!(
4553 JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
4554 "derived index rebuild should restore witnesses without allocating identities",
4555 );
4556 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
4557 JOURNALED_SCHEMA_STORE.with(|store| {
4558 let cursor = store
4559 .borrow()
4560 .identity_statement_cursor(
4561 database_incarnation_id().expect("database incarnation should remain readable"),
4562 ENTITY_TAG,
4563 FieldId::new(1),
4564 &AcceptedFieldKind::Nat64,
4565 )
4566 .expect("folded identity state should reopen without allocating");
4567 assert_eq!(cursor.expected_high_water(), 8);
4568 assert!(!cursor.has_allocations());
4569 });
4570 }
4571
4572 #[test]
4573 #[ignore = "release-closeout native timing probe for one marker-authorized Identity recovery"]
4574 fn identity_recovery_closeout_reports_guarded_reentry_time() {
4575 let session = initialize_journaled();
4576 let catalog = session
4577 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4578 .expect("journaled identity catalog should resolve");
4579 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4580 .expect("journaled identity row layout should build");
4581
4582 interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
4583 let interrupted = session.execute_accepted_structural_save_batch(
4584 &catalog,
4585 &descriptor,
4586 batch(&[1]),
4587 Timestamp::from_millis(10),
4588 Ok,
4589 );
4590 assert!(
4591 interrupted.is_err(),
4592 "the selected publication boundary should interrupt",
4593 );
4594
4595 let start = Instant::now();
4596 let committed = session
4597 .execute_accepted_structural_save_batch(
4598 &catalog,
4599 &descriptor,
4600 batch(&[2]),
4601 Timestamp::from_millis(11),
4602 Ok,
4603 )
4604 .expect("guarded reentry should recover before allocation");
4605 let elapsed = start.elapsed();
4606 assert_eq!(
4607 committed
4608 .into_iter()
4609 .map(|row| row.values)
4610 .collect::<Vec<_>>(),
4611 vec![vec![Value::Nat64(2), Value::Nat64(2)]],
4612 );
4613
4614 println!(
4615 "identity recovery closeout: guarded_reentry_nanos={}",
4616 elapsed.as_nanos(),
4617 );
4618 }
4619}
4620
4621#[cfg(test)]
4622mod targeted_rule_mutation_tests {
4623 use super::{
4624 DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
4625 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
4626 };
4627 use crate::{
4628 db::{
4629 data::{DataStore, encode_input_value_for_candidate_field_contract},
4630 index::IndexStore,
4631 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
4632 schema::{
4633 AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
4634 AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
4635 AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
4636 ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
4637 PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
4638 ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
4639 SchemaRowLayout, SchemaStore, SchemaVersion,
4640 accepted_schema_candidate_with_catalogs_for_tests,
4641 build_record_newtype_composite_catalog_for_tests,
4642 empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
4643 },
4644 },
4645 error::InternalError,
4646 traits::{CanisterKind, Path},
4647 types::EntityTag,
4648 value::InputValue,
4649 };
4650 use icydb_schema::{
4651 ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
4652 };
4653 use std::{cell::RefCell, collections::BTreeMap};
4654
4655 const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
4656 const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
4657 const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
4658 const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
4659 const UPDATED_AT_SOURCE: &str =
4660 "session::write::targeted_rule_mutation_tests::Entity::updated_at";
4661 const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
4662 const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
4663 const DEGREE_MEMBER_SOURCE: &str =
4664 "session::write::targeted_rule_mutation_tests::Profile::degree";
4665 const DEGREE_RULE_SOURCE: &str =
4666 "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
4667
4668 struct TestCanister;
4669
4670 impl Path for TestCanister {
4671 const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
4672 }
4673
4674 impl CanisterKind for TestCanister {
4675 const COMMIT_MEMORY_ID: u8 = 43;
4676 const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
4677 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
4678 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
4679 "icydb.targeted_mutation_tests.integrity.progress.v1";
4680 }
4681
4682 thread_local! {
4683 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
4684 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
4685 static SCHEMA_STORE: RefCell<SchemaStore> =
4686 const { RefCell::new(SchemaStore::init_heap()) };
4687 static STORE_REGISTRY: StoreRegistry = {
4688 let mut registry = StoreRegistry::new();
4689 registry.register_store(
4690 STORE_PATH,
4691 &DATA_STORE,
4692 &INDEX_STORE,
4693 &SCHEMA_STORE,
4694 StoreAllocationIdentities::absent(),
4695 StoreRuntimeStorageCapabilities::heap(),
4696 ).expect("targeted mutation test store should register");
4697 registry
4698 };
4699 }
4700
4701 fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
4702 parse(raw.to_string()).expect("test source identity should admit")
4703 }
4704
4705 fn profile_input(degree: u64) -> InputValue {
4706 InputValue::Map(vec![(
4707 InputValue::Text("degree".to_string()),
4708 InputValue::Nat64(degree),
4709 )])
4710 }
4711
4712 fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
4713 DynamicStructuralPatch::new(vec![
4714 (
4715 "id".to_string(),
4716 DynamicWriteCell::Value(InputValue::Nat64(id)),
4717 ),
4718 (
4719 "profile".to_string(),
4720 DynamicWriteCell::Value(profile_input(degree)),
4721 ),
4722 ])
4723 }
4724
4725 fn encoded_value(
4726 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4727 composite_catalog: &AcceptedCompositeCatalog,
4728 name: &str,
4729 kind: &AcceptedFieldKind,
4730 storage_decode: FieldStorageDecode,
4731 leaf_codec: LeafCodec,
4732 value: InputValue,
4733 ) -> Vec<u8> {
4734 let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
4735 encode_input_value_for_candidate_field_contract(
4736 enum_catalog,
4737 composite_catalog,
4738 field,
4739 value,
4740 &mut ValueAdmissionBudget::standard(),
4741 )
4742 .expect("test accepted value should encode")
4743 }
4744
4745 fn nat64_literal(
4746 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4747 composite_catalog: &AcceptedCompositeCatalog,
4748 value: u64,
4749 ) -> AcceptedCheckLiteralV1 {
4750 let kind = AcceptedFieldKind::Nat64;
4751 AcceptedCheckLiteralV1::from_accepted_parts(
4752 kind.clone(),
4753 FieldStorageDecode::ByKind,
4754 LeafCodec::Scalar(ScalarCodec::Nat64),
4755 encoded_value(
4756 enum_catalog,
4757 composite_catalog,
4758 "degree_bound",
4759 &kind,
4760 FieldStorageDecode::ByKind,
4761 LeafCodec::Scalar(ScalarCodec::Nat64),
4762 InputValue::Nat64(value),
4763 ),
4764 )
4765 }
4766
4767 fn targeted_constraint_id(error: &InternalError) -> u32 {
4768 let facts = error.diagnostic_facts();
4769 assert!(facts.contains(&(
4770 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4771 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4772 )));
4773 assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
4774 assert!(facts.contains(&(
4775 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
4776 icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
4777 )));
4778 assert_eq!(
4779 facts
4780 .iter()
4781 .filter(|(tag, _)| matches!(
4782 tag,
4783 icydb_diagnostic_code::DiagnosticFactTag::RootField
4784 | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
4785 ))
4786 .copied()
4787 .collect::<Vec<_>>(),
4788 vec![
4789 (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
4790 (
4791 icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
4792 icydb_diagnostic_code::pack_u32_pair(1, 1),
4793 ),
4794 ]
4795 );
4796 let value = facts
4797 .iter()
4798 .find_map(|(tag, value)| {
4799 (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
4800 })
4801 .expect("targeted mutation should retain its accepted constraint ID");
4802 u32::try_from(value).expect("accepted constraint ID fits u32")
4803 }
4804
4805 #[expect(
4806 clippy::too_many_lines,
4807 reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
4808 )]
4809 #[test]
4810 fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
4811 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
4812 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
4813 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
4814
4815 let entity_tag = EntityTag::new(93);
4816 let enum_catalog = empty_accepted_enum_catalog_for_tests();
4817 let (composite_catalog, profile_type, degree_type, degree_member) =
4818 build_record_newtype_composite_catalog_for_tests(
4819 "tests::TargetedProfile".to_string(),
4820 "degree".to_string(),
4821 "tests::TargetedDegree".to_string(),
4822 AcceptedFieldKind::Nat64,
4823 &enum_catalog,
4824 )
4825 .expect("targeted mutation composites should close");
4826 let profile_kind = AcceptedFieldKind::Composite {
4827 type_id: profile_type,
4828 };
4829 let profile_default = encoded_value(
4830 &enum_catalog,
4831 &composite_catalog,
4832 "profile",
4833 &profile_kind,
4834 FieldStorageDecode::CatalogValue,
4835 LeafCodec::Structural,
4836 profile_input(12),
4837 );
4838 let fields = vec![
4839 PersistedFieldSnapshot::new_initial(
4840 FieldId::new(1),
4841 "id".to_string(),
4842 SchemaFieldSlot::new(0),
4843 AcceptedFieldKind::Nat64,
4844 Vec::new(),
4845 false,
4846 SchemaInsertDefault::None,
4847 FieldStorageDecode::ByKind,
4848 LeafCodec::Scalar(ScalarCodec::Nat64),
4849 ),
4850 PersistedFieldSnapshot::new_initial(
4851 FieldId::new(2),
4852 "profile".to_string(),
4853 SchemaFieldSlot::new(1),
4854 profile_kind,
4855 vec![PersistedNestedLeafSnapshot::new(
4856 vec!["degree".to_string()],
4857 AcceptedFieldKind::Composite {
4858 type_id: degree_type,
4859 },
4860 false,
4861 )],
4862 false,
4863 SchemaInsertDefault::SlotPayload(profile_default),
4864 FieldStorageDecode::CatalogValue,
4865 LeafCodec::Structural,
4866 ),
4867 PersistedFieldSnapshot::new_initial_with_write_policy(
4868 FieldId::new(3),
4869 "updated_at".to_string(),
4870 SchemaFieldSlot::new(2),
4871 AcceptedFieldKind::Timestamp,
4872 Vec::new(),
4873 false,
4874 SchemaInsertDefault::None,
4875 SchemaFieldWritePolicy::from_model_policies(
4876 None,
4877 Some(FieldWriteManagement::UpdatedAt),
4878 ),
4879 FieldStorageDecode::ByKind,
4880 LeafCodec::Scalar(ScalarCodec::Timestamp),
4881 ),
4882 ];
4883 let mut snapshot = PersistedSchemaSnapshot::new(
4884 SchemaVersion::initial(),
4885 ENTITY_SOURCE.to_string(),
4886 "TargetedMutation".to_string(),
4887 FieldId::new(1),
4888 SchemaRowLayout::initial(
4889 fields
4890 .iter()
4891 .map(|field| (field.id(), field.slot()))
4892 .collect(),
4893 ),
4894 fields,
4895 );
4896 let constraint_catalog = snapshot
4897 .constraint_catalog()
4898 .clone()
4899 .with_added_targeted_rule(
4900 "profile_degree_multiple".to_string(),
4901 ConstraintOrigin::Generated,
4902 AcceptedRuleTarget::new(
4903 FieldId::new(2),
4904 AcceptedNamedTypeIdentity::Composite(degree_type),
4905 ),
4906 AcceptedRuleOperation::MultipleOf {
4907 divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
4908 },
4909 )
4910 .expect("targeted mutation rule should allocate");
4911 let targeted_rule_id = constraint_catalog
4912 .constraints()
4913 .last()
4914 .expect("targeted mutation rule should persist")
4915 .id();
4916 snapshot = snapshot.with_constraint_catalog(constraint_catalog);
4917
4918 let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
4919 let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
4920 let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
4921 let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
4922 let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
4923 let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
4924 let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
4925 let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
4926 let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
4927 BTreeMap::from([(entity_source, entity_tag)]),
4928 BTreeMap::from([
4929 ((entity_tag, id_source), FieldId::new(1)),
4930 ((entity_tag, profile_source), FieldId::new(2)),
4931 ((entity_tag, updated_at_source), FieldId::new(3)),
4932 ]),
4933 BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
4934 BTreeMap::new(),
4935 BTreeMap::new(),
4936 )
4937 .with_initial_named_types_for_tests(
4938 BTreeMap::from([
4939 (
4940 profile_type_source,
4941 AcceptedNamedTypeIdentity::Composite(profile_type),
4942 ),
4943 (
4944 degree_type_source,
4945 AcceptedNamedTypeIdentity::Composite(degree_type),
4946 ),
4947 ]),
4948 BTreeMap::new(),
4949 BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
4950 );
4951 let candidate = accepted_schema_candidate_with_catalogs_for_tests(
4952 STORE_PATH,
4953 AcceptedSchemaRevision::INITIAL,
4954 enum_catalog,
4955 composite_catalog,
4956 source_bindings,
4957 BTreeMap::from([(entity_tag, snapshot)]),
4958 );
4959
4960 let session = DbSession::<TestCanister>::new(
4961 &STORE_REGISTRY,
4962 &crate::db::RequestExecutionRoot::__new_runtime_root(),
4963 );
4964 session
4965 .db
4966 .ensure_recovered_state()
4967 .expect("targeted mutation test database should initialize");
4968 let store = session
4969 .db
4970 .store_handle(STORE_PATH)
4971 .expect("targeted mutation test store should resolve");
4972 crate::db::commit::publish_accepted_schema_candidate(
4973 STORE_PATH,
4974 store,
4975 AcceptedSchemaRevision::NONE,
4976 &candidate,
4977 )
4978 .expect("targeted mutation candidate should publish");
4979
4980 let dynamic_error = session
4981 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4982 entity: "TargetedMutation".to_string(),
4983 patch: structural_patch(1, 12),
4984 })
4985 .expect_err("dynamic write must enforce the targeted rule");
4986 assert_eq!(
4987 targeted_constraint_id(&dynamic_error),
4988 targeted_rule_id.get()
4989 );
4990
4991 let binding = session
4992 .issue_typed_entity_binding(
4993 ENTITY_SOURCE,
4994 &[
4995 DynamicTypedFieldBindingRequest::new(
4996 ID_SOURCE.to_string(),
4997 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4998 false,
4999 ),
5000 DynamicTypedFieldBindingRequest::new(
5001 PROFILE_SOURCE.to_string(),
5002 DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
5003 false,
5004 ),
5005 DynamicTypedFieldBindingRequest::new(
5006 UPDATED_AT_SOURCE.to_string(),
5007 DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
5008 false,
5009 ),
5010 ],
5011 )
5012 .expect("targeted typed binding should issue");
5013 let typed_patch = binding
5014 .bind_write_fields(vec![
5015 (
5016 ID_SOURCE.to_string(),
5017 DynamicWriteCell::Value(InputValue::Nat64(2)),
5018 ),
5019 (
5020 PROFILE_SOURCE.to_string(),
5021 DynamicWriteCell::Value(profile_input(12)),
5022 ),
5023 ])
5024 .expect("targeted typed patch should bind");
5025 let typed_error = session
5026 .execute_trusted_typed_mutation(
5027 &binding,
5028 &DynamicTypedMutation::Insert { patch: typed_patch },
5029 )
5030 .expect_err("typed write must enforce the targeted rule");
5031 assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
5032
5033 #[cfg(feature = "sql")]
5034 {
5035 let sql_error = session
5036 .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
5037 .expect_err("SQL default resolution must enforce the targeted rule");
5038 let crate::db::QueryError::Execute(execute) = sql_error else {
5039 panic!("targeted SQL write should fail at shared execution admission");
5040 };
5041 assert_eq!(
5042 targeted_constraint_id(execute.as_internal()),
5043 targeted_rule_id.get()
5044 );
5045 }
5046
5047 session
5048 .execute_trusted_dynamic_mutation_batch(vec![
5049 DynamicMutation::Insert {
5050 entity: "TargetedMutation".to_string(),
5051 patch: structural_patch(4, 5),
5052 },
5053 DynamicMutation::Insert {
5054 entity: "TargetedMutation".to_string(),
5055 patch: structural_patch(5, 12),
5056 },
5057 ])
5058 .expect_err("one invalid targeted value must reject the whole batch");
5059 assert_eq!(
5060 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
5061 Some(0),
5062 "no frontend or earlier valid batch row may escape targeted admission",
5063 );
5064
5065 let admitted = session
5066 .execute_trusted_dynamic_mutation_batch(vec![
5067 DynamicMutation::Insert {
5068 entity: "TargetedMutation".to_string(),
5069 patch: structural_patch(6, 5),
5070 },
5071 DynamicMutation::Insert {
5072 entity: "TargetedMutation".to_string(),
5073 patch: structural_patch(7, 10),
5074 },
5075 ])
5076 .expect("compliant targeted values should share one accepted batch");
5077 let [first, second] = admitted.rows.as_slice() else {
5078 panic!("the mixed targeted batch should return two rows");
5079 };
5080 let first_timestamp = first
5081 .get(2)
5082 .expect("the first mixed row should contain its managed timestamp");
5083 assert!(matches!(
5084 first_timestamp,
5085 crate::value::OutputValue::Timestamp(_)
5086 ));
5087 assert_eq!(
5088 second.get(2),
5089 Some(first_timestamp),
5090 "one accepted mixed batch must materialize one managed timestamp",
5091 );
5092 assert_eq!(
5093 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
5094 Some(2),
5095 );
5096 }
5097}