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, 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::executor::budget::{
2687 HardExecutionBudget, HardExecutionContext, HardExecutionFailureHeadroom,
2688 with_query_execution_budget_for_tests,
2689 };
2690 use crate::{
2691 db::{
2692 commit::{database_incarnation_id, forget_recovered_domain_for_tests},
2693 data::DataStore,
2694 executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
2695 index::IndexStore,
2696 integrity::{
2697 PhysicalUnitCheckpoint, QuickIntegrityStatus, RowInspectionLimits,
2698 execute_quick_integrity, execute_row_integrity_page,
2699 },
2700 journal::JournalTailStore,
2701 registry::{
2702 StoreAllocationIdentities, StoreAllocationIdentity, StoreRegistry,
2703 StoreRuntimeStorageCapabilities,
2704 },
2705 schema::{
2706 AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
2707 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2708 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2709 PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
2710 SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
2711 accepted_schema_candidate_with_field_bindings_for_tests,
2712 },
2713 write_context::MutationMode,
2714 },
2715 error::{ErrorClass, ErrorOrigin, InternalError},
2716 testing::test_memory,
2717 traits::{CanisterKind, Path},
2718 types::{EntityTag, Timestamp},
2719 value::{InputValue, OutputValue, Value},
2720 };
2721 use icydb_schema::{FieldSourceKey, ScalarType};
2722 use std::{cell::RefCell, collections::BTreeMap, time::Instant};
2723
2724 const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
2725 const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
2726 const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
2727 const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
2728 const ENTITY_NAME: &str = "IdentityRow";
2729 const ENTITY_TAG: EntityTag = EntityTag::new(93);
2730 const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
2731
2732 struct TestCanister;
2733
2734 impl Path for TestCanister {
2735 const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
2736 }
2737
2738 impl CanisterKind for TestCanister {
2739 const COMMIT_MEMORY_ID: u8 = 45;
2740 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
2741 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
2742 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2743 "icydb.identity_pre_key_tests.integrity.progress.v1";
2744 }
2745
2746 thread_local! {
2747 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2748 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2749 static SCHEMA_STORE: RefCell<SchemaStore> =
2750 const { RefCell::new(SchemaStore::init_heap()) };
2751 static STORE_REGISTRY: StoreRegistry = {
2752 let mut registry = StoreRegistry::new();
2753 registry.register_store(
2754 STORE_PATH,
2755 &DATA_STORE,
2756 &INDEX_STORE,
2757 &SCHEMA_STORE,
2758 StoreAllocationIdentities::absent(),
2759 StoreRuntimeStorageCapabilities::heap(),
2760 ).expect("identity pre-key test store should register");
2761 registry
2762 };
2763 static JOURNALED_DATA_STORE: RefCell<DataStore> =
2764 RefCell::new(DataStore::init_journaled(test_memory(186)));
2765 static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
2766 RefCell::new(IndexStore::init_journaled(test_memory(187)));
2767 static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
2768 RefCell::new(SchemaStore::init_journaled(test_memory(188)));
2769 static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
2770 RefCell::new(JournalTailStore::init(test_memory(189)));
2771 static JOURNALED_STORE_REGISTRY: StoreRegistry = {
2772 let mut registry = StoreRegistry::new();
2773 registry.register_journaled_store(
2774 JOURNALED_STORE_PATH,
2775 &JOURNALED_DATA_STORE,
2776 &JOURNALED_INDEX_STORE,
2777 &JOURNALED_SCHEMA_STORE,
2778 &JOURNALED_TAIL_STORE,
2779 StoreAllocationIdentities::new_journaled(
2780 StoreAllocationIdentity::new(186, "icydb.test.identity-range.data.v1"),
2781 StoreAllocationIdentity::new(187, "icydb.test.identity-range.index.v1"),
2782 StoreAllocationIdentity::new(188, "icydb.test.identity-range.schema.v1"),
2783 StoreAllocationIdentity::new(189, "icydb.test.identity-range.journal.v1"),
2784 ),
2785 StoreRuntimeStorageCapabilities::journaled(),
2786 ).expect("identity range journaled store should register");
2787 registry
2788 };
2789 }
2790
2791 struct JournaledTestCanister;
2792
2793 impl Path for JournaledTestCanister {
2794 const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
2795 }
2796
2797 impl CanisterKind for JournaledTestCanister {
2798 const COMMIT_MEMORY_ID: u8 = 190;
2799 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
2800 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
2801 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2802 "icydb.identity_range_tests.integrity.progress.v1";
2803 }
2804
2805 fn source_key(source: &str) -> FieldSourceKey {
2806 FieldSourceKey::try_new(source).expect("identity test field source should admit")
2807 }
2808
2809 fn identity_snapshot(store_path: &str) -> PersistedSchemaSnapshot {
2810 let fields = vec![
2811 PersistedFieldSnapshot::new_initial_with_write_policy(
2812 FieldId::new(1),
2813 "id".to_string(),
2814 SchemaFieldSlot::new(0),
2815 AcceptedFieldKind::Nat64,
2816 Vec::new(),
2817 false,
2818 SchemaInsertDefault::None,
2819 SchemaFieldWritePolicy::from_model_policies(
2820 Some(FieldInsertGeneration::Identity),
2821 None,
2822 ),
2823 FieldStorageDecode::ByKind,
2824 LeafCodec::Scalar(ScalarCodec::Nat64),
2825 ),
2826 PersistedFieldSnapshot::new_initial(
2827 FieldId::new(2),
2828 "payload".to_string(),
2829 SchemaFieldSlot::new(1),
2830 AcceptedFieldKind::Nat64,
2831 Vec::new(),
2832 false,
2833 SchemaInsertDefault::None,
2834 FieldStorageDecode::ByKind,
2835 LeafCodec::Scalar(ScalarCodec::Nat64),
2836 ),
2837 ];
2838 PersistedSchemaSnapshot::new_with_indexes(
2839 SchemaVersion::initial(),
2840 ENTITY_SOURCE.to_string(),
2841 ENTITY_NAME.to_string(),
2842 FieldId::new(1),
2843 SchemaRowLayout::initial(
2844 fields
2845 .iter()
2846 .map(|field| (field.id(), field.slot()))
2847 .collect(),
2848 ),
2849 fields,
2850 vec![PersistedIndexSnapshot::new(
2851 SchemaIndexId::new(1).expect("identity test index ID should admit"),
2852 1,
2853 "by_payload".to_string(),
2854 store_path.to_string(),
2855 false,
2856 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2857 FieldId::new(2),
2858 SchemaFieldSlot::new(1),
2859 vec!["payload".to_string()],
2860 AcceptedFieldKind::Nat64,
2861 false,
2862 )]),
2863 None,
2864 )],
2865 )
2866 }
2867
2868 fn initialize() -> DbSession<TestCanister> {
2869 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2870 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2871 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2872 let session = DbSession::<TestCanister>::new(
2873 &STORE_REGISTRY,
2874 &crate::db::RequestExecutionRoot::__new_runtime_root(),
2875 );
2876 session
2877 .db
2878 .ensure_recovered_state()
2879 .expect("identity pre-key test database should initialize");
2880 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2881 STORE_PATH,
2882 AcceptedSchemaRevision::INITIAL,
2883 BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH))]),
2884 BTreeMap::from([
2885 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2886 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2887 ]),
2888 );
2889 let store = session
2890 .db
2891 .store_handle(STORE_PATH)
2892 .expect("identity pre-key test store should resolve");
2893 crate::db::commit::publish_accepted_schema_candidate(
2894 STORE_PATH,
2895 store,
2896 AcceptedSchemaRevision::NONE,
2897 &candidate,
2898 )
2899 .expect("identity candidate should publish with explicit zero state");
2900 session
2901 }
2902
2903 fn initialize_journaled() -> DbSession<JournaledTestCanister> {
2904 let session = DbSession::<JournaledTestCanister>::new(
2905 &JOURNALED_STORE_REGISTRY,
2906 &crate::db::RequestExecutionRoot::__new_runtime_root(),
2907 );
2908 session
2909 .db
2910 .ensure_recovered_state()
2911 .expect("journaled identity database should initialize");
2912 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2913 JOURNALED_STORE_PATH,
2914 AcceptedSchemaRevision::INITIAL,
2915 BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH))]),
2916 BTreeMap::from([
2917 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2918 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2919 ]),
2920 );
2921 let store = session
2922 .db
2923 .store_handle(JOURNALED_STORE_PATH)
2924 .expect("journaled identity store should resolve");
2925 crate::db::commit::publish_accepted_schema_candidate(
2926 JOURNALED_STORE_PATH,
2927 store,
2928 AcceptedSchemaRevision::NONE,
2929 &candidate,
2930 )
2931 .expect("journaled identity candidate should publish");
2932 session
2933 }
2934
2935 fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
2936 AcceptedMutationIntentPatch::new()
2937 .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
2938 }
2939
2940 fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
2941 DynamicStructuralPatch::new(vec![(
2942 "payload".to_string(),
2943 DynamicWriteCell::Value(InputValue::Nat64(value)),
2944 )])
2945 }
2946
2947 fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
2948 vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
2949 }
2950
2951 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2952 fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
2953 session
2954 .issue_typed_entity_binding(
2955 ENTITY_SOURCE,
2956 &[
2957 DynamicTypedFieldBindingRequest::new(
2958 ID_SOURCE.to_string(),
2959 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2960 false,
2961 ),
2962 DynamicTypedFieldBindingRequest::new(
2963 PAYLOAD_SOURCE.to_string(),
2964 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2965 false,
2966 ),
2967 ],
2968 )
2969 .expect("exact-key test binding should issue")
2970 }
2971
2972 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2973 fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
2974 let output = session
2975 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
2976 entity: ENTITY_NAME.to_string(),
2977 patch: dynamic_payload_patch(payload),
2978 })
2979 .expect("exact-key fixture insert should commit");
2980 match output.rows.as_slice() {
2981 [row] => match row.as_slice() {
2982 [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
2983 if *actual_payload == payload =>
2984 {
2985 *id
2986 }
2987 _ => panic!("exact-key fixture should return its identity and payload"),
2988 },
2989 _ => panic!("exact-key fixture insert should return one row"),
2990 }
2991 }
2992
2993 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2994 fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
2995 let first = insert_exact_key_fixture(session, 41);
2996 let second = insert_exact_key_fixture(session, 42);
2997 let missing = u64::MAX;
2998 let binding = exact_key_binding(session);
2999 let gets_before = DataStore::current_get_call_count();
3000 let result = session
3001 .execute_public_exact_key_batch_for_typed_binding(
3002 &binding,
3003 &[second, missing, first, second],
3004 )
3005 .expect("exact-key batch should execute")
3006 .expect("exact-key binding should remain current");
3007
3008 assert_eq!(result.positions, vec![0, 1, 2, 0]);
3009 assert_eq!(
3010 result.distinct_rows,
3011 vec![
3012 Some(expected_dynamic_row(second, 42)),
3013 None,
3014 Some(expected_dynamic_row(first, 41)),
3015 ],
3016 );
3017 assert_eq!(
3018 DataStore::current_get_call_count().saturating_sub(gets_before),
3019 3,
3020 "four input positions with one duplicate must perform three physical reads",
3021 );
3022 }
3023
3024 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3025 #[test]
3026 fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
3027 assert_exact_key_batch(&initialize());
3028 assert_exact_key_batch(&initialize_journaled());
3029 }
3030
3031 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3032 #[test]
3033 fn exact_key_batch_uses_typed_hard_execution_budget() {
3034 let session = initialize();
3035 let binding = exact_key_binding(&session);
3036 let budget =
3037 HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
3038 let error = session
3039 .execute_exact_key_batch_with_hard_budget_for_tests(&binding, &[u64::MAX], &budget)
3040 .expect_err("zero query budget should reject the exact-key route");
3041
3042 assert!(matches!(
3043 error.diagnostic().detail(),
3044 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3045 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3046 })
3047 ));
3048 let facts = error.diagnostic_facts();
3049 assert_eq!(
3050 &facts[..5],
3051 &[
3052 (
3053 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3054 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions.raw(),
3055 ),
3056 (icydb_diagnostic_code::DiagnosticFactTag::Limit, 0),
3057 (icydb_diagnostic_code::DiagnosticFactTag::Actual, 1),
3058 (
3059 icydb_diagnostic_code::DiagnosticFactTag::ExecutionBudgetScope,
3060 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution.raw(),
3061 ),
3062 (
3063 icydb_diagnostic_code::DiagnosticFactTag::ExecutionLane,
3064 icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead.raw(),
3065 ),
3066 ],
3067 );
3068 assert_eq!(
3069 facts[5].0,
3070 icydb_diagnostic_code::DiagnosticFactTag::QueryShapeFingerprintPrefix,
3071 );
3072 assert_ne!(facts[5].1, 0);
3073 }
3074
3075 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3076 fn assert_planned_query_exhausts(
3077 session: &DbSession<TestCanister>,
3078 query: &crate::db::DynamicQuery,
3079 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3080 ) {
3081 let budget = HardExecutionBudget::uniform_for_tests(
3082 u64::MAX,
3083 HardExecutionFailureHeadroom::new(500, 256),
3084 )
3085 .with_limit_for_tests(resource, 0);
3086 let context = HardExecutionContext::new(
3087 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3088 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3089 0x7068_7973_6963_616c,
3090 );
3091 let error = with_query_execution_budget_for_tests(budget, context, || {
3092 session.execute_trusted_live_page(query, None)
3093 })
3094 .expect_err("the injected zero resource allowance should reject planned execution");
3095
3096 assert!(matches!(
3097 error.diagnostic().detail(),
3098 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3099 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3100 })
3101 ));
3102 assert_eq!(
3103 error.diagnostic_facts()[0],
3104 (
3105 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3106 resource.raw(),
3107 ),
3108 );
3109 }
3110
3111 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3112 fn assert_grouped_query_exhausts(
3113 session: &DbSession<TestCanister>,
3114 query: &crate::db::DynamicQuery,
3115 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3116 ) {
3117 let budget = HardExecutionBudget::uniform_for_tests(
3118 u64::MAX,
3119 HardExecutionFailureHeadroom::new(500, 256),
3120 )
3121 .with_limit_for_tests(resource, 0);
3122 let context = HardExecutionContext::new(
3123 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3124 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3125 0x6772_6f75_7065_642d,
3126 );
3127 let error = with_query_execution_budget_for_tests(budget, context, || {
3128 session.execute_trusted_dynamic_grouped_query(query)
3129 })
3130 .expect_err("the injected zero resource allowance should reject grouped execution");
3131
3132 assert!(matches!(
3133 error.diagnostic().detail(),
3134 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3135 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3136 })
3137 ));
3138 assert_eq!(
3139 error.diagnostic_facts()[0],
3140 (
3141 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3142 resource.raw(),
3143 ),
3144 );
3145 }
3146
3147 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3148 fn assert_sql_query_exhausts(
3149 session: &DbSession<TestCanister>,
3150 sql: &str,
3151 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3152 ) {
3153 let budget = HardExecutionBudget::uniform_for_tests(
3154 u64::MAX,
3155 HardExecutionFailureHeadroom::new(500, 256),
3156 )
3157 .with_limit_for_tests(resource, 0);
3158 let context = HardExecutionContext::new(
3159 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3160 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3161 0x7371_6c2d_736f_7274,
3162 );
3163 let error = with_query_execution_budget_for_tests(budget, context, || {
3164 session.execute_trusted_sql_query(sql)
3165 })
3166 .expect_err("the injected zero resource allowance should reject SQL execution");
3167
3168 assert!(matches!(
3169 error.diagnostic().detail(),
3170 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3171 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3172 })
3173 ));
3174 assert_eq!(
3175 error.diagnostic_facts()[0],
3176 (
3177 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3178 resource.raw(),
3179 ),
3180 );
3181 }
3182
3183 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3184 #[test]
3185 fn planned_read_routes_share_physical_resource_accounting() {
3186 let session = initialize();
3187 let first = insert_exact_key_fixture(&session, 41);
3188 insert_exact_key_fixture(&session, 42);
3189
3190 let fallback = crate::db::DynamicQuery::new(ENTITY_NAME)
3191 .filter(crate::db::FieldRef::new("id").eq(first))
3192 .select(["id", "payload"])
3193 .order_by(crate::db::asc("id"))
3194 .limit(1);
3195 assert_eq!(
3196 session
3197 .execute_trusted_live_page(&fallback, None)
3198 .expect("bounded fallback execution should preserve its result")
3199 .row_count,
3200 1,
3201 );
3202 assert_planned_query_exhausts(
3203 &session,
3204 &fallback,
3205 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::RowsVisited,
3206 );
3207
3208 let covering = crate::db::DynamicQuery::new(ENTITY_NAME)
3209 .filter(crate::db::FieldRef::new("payload").eq(41_u64))
3210 .select(["payload"])
3211 .order_by(crate::db::asc("payload"))
3212 .limit(1);
3213 assert_eq!(
3214 session
3215 .execute_trusted_live_page(&covering, None)
3216 .expect("bounded covering execution should preserve its result")
3217 .row_count,
3218 1,
3219 );
3220 assert_planned_query_exhausts(
3221 &session,
3222 &covering,
3223 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
3224 );
3225
3226 let residual = crate::db::DynamicQuery::new(ENTITY_NAME)
3227 .filter(crate::db::FieldRef::new("payload").eq_field("id"))
3228 .select(["id"])
3229 .order_by(crate::db::asc("id"))
3230 .limit(1);
3231 assert_eq!(
3232 session
3233 .execute_trusted_live_page(&residual, None)
3234 .expect("bounded residual execution should preserve its result")
3235 .row_count,
3236 0,
3237 );
3238 assert_planned_query_exhausts(
3239 &session,
3240 &residual,
3241 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::PredicateExpressionSteps,
3242 );
3243
3244 assert_planned_query_exhausts(
3245 &session,
3246 &fallback,
3247 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::ResultBytes,
3248 );
3249
3250 let grouped = crate::db::DynamicQuery::new(ENTITY_NAME)
3251 .group_by("payload")
3252 .aggregate(crate::db::count())
3253 .order_by(crate::db::asc("payload"))
3254 .grouped_limits(10, 16 * 1_024)
3255 .limit(1);
3256 let grouped_result = session
3257 .execute_trusted_dynamic_grouped_query(&grouped)
3258 .expect("bounded grouped execution should preserve its result");
3259 assert_eq!(grouped_result.row_count, 1);
3260 assert!(grouped_result.next_cursor.is_some());
3261 assert_grouped_query_exhausts(
3262 &session,
3263 &grouped,
3264 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::GroupDistinctEntries,
3265 );
3266 assert_grouped_query_exhausts(
3267 &session,
3268 &grouped,
3269 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::CursorSteps,
3270 );
3271
3272 assert_sql_query_exhausts(
3273 &session,
3274 "SELECT payload, COUNT(*) AS row_count FROM IdentityRow \
3275 GROUP BY payload ORDER BY row_count DESC, payload ASC LIMIT 1",
3276 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::SortEntries,
3277 );
3278 }
3279
3280 fn assert_dynamic_payload(session: &DbSession<TestCanister>, key: u64, expected_payload: u64) {
3281 let unchanged = session
3282 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3283 entity: ENTITY_NAME.to_string(),
3284 key: InputValue::Nat64(key),
3285 patch: dynamic_payload_patch(expected_payload),
3286 })
3287 .expect("the expected row should remain readable through a no-op update");
3288 assert_eq!(unchanged.affected_rows, 0);
3289 assert_eq!(
3290 unchanged.rows,
3291 vec![expected_dynamic_row(key, expected_payload)],
3292 );
3293 }
3294
3295 fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
3296 values
3297 .iter()
3298 .map(|value| {
3299 AcceptedStructuralMutation::save(
3300 MutationMode::Insert,
3301 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
3302 payload_patch(*value),
3303 )
3304 })
3305 .collect()
3306 }
3307
3308 fn assert_identity_boundary(error: &InternalError) {
3309 assert_eq!(error.class(), ErrorClass::Unsupported);
3310 assert_eq!(error.origin(), ErrorOrigin::Identity);
3311 }
3312
3313 #[test]
3314 fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
3315 let generated = insert_key_exists_after_generation(true);
3316 assert_eq!(generated.class(), ErrorClass::Corruption);
3317 assert_eq!(generated.origin(), ErrorOrigin::Identity);
3318
3319 let ordinary = insert_key_exists_after_generation(false);
3320 assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
3321 }
3322
3323 #[cfg(target_pointer_width = "64")]
3324 #[test]
3325 fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
3326 let error = checked_pre_key_candidate_count(
3327 usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
3328 )
3329 .expect_err("candidate counts beyond u32 must reject");
3330 assert_identity_boundary(&error);
3331 }
3332
3333 #[test]
3334 #[expect(
3335 clippy::too_many_lines,
3336 reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
3337 )]
3338 fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
3339 let session = initialize();
3340 let seeded = session
3341 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3342 .expect("seed rows should commit");
3343 assert_eq!(seeded.affected_rows, 1);
3344
3345 let split = session
3346 .execute_trusted_dynamic_mutation_batch(vec![
3347 DynamicMutation::Update {
3348 entity: ENTITY_NAME.to_string(),
3349 key: InputValue::Nat64(1),
3350 patch: dynamic_payload_patch(60),
3351 },
3352 DynamicMutation::Insert {
3353 entity: ENTITY_NAME.to_string(),
3354 patch: dynamic_payload_patch(40),
3355 },
3356 ])
3357 .expect("one holding should split atomically");
3358 assert_eq!(split.affected_rows, 2);
3359 assert_eq!(
3360 split.rows,
3361 vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
3362 "split after-images must retain input order and exact quantity",
3363 );
3364
3365 let rejected_split = session
3366 .execute_trusted_dynamic_mutation_batch(vec![
3367 DynamicMutation::Update {
3368 entity: ENTITY_NAME.to_string(),
3369 key: InputValue::Nat64(1),
3370 patch: dynamic_payload_patch(50),
3371 },
3372 DynamicMutation::Insert {
3373 entity: ENTITY_NAME.to_string(),
3374 patch: DynamicStructuralPatch::new(Vec::new()),
3375 },
3376 ])
3377 .expect_err("an invalid split output must reject the staged source update");
3378 assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
3379 assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
3380 assert_eq!(
3381 rejected_split.diagnostic_facts(),
3382 vec![
3383 (
3384 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3385 ENTITY_TAG.value(),
3386 ),
3387 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
3388 (
3389 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3390 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3391 ),
3392 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3393 ],
3394 );
3395 assert_dynamic_payload(&session, 1, 60);
3396 assert_dynamic_payload(&session, 2, 40);
3397
3398 let transfer = session
3399 .execute_trusted_dynamic_mutation_batch(vec![
3400 DynamicMutation::Update {
3401 entity: ENTITY_NAME.to_string(),
3402 key: InputValue::Nat64(1),
3403 patch: dynamic_payload_patch(70),
3404 },
3405 DynamicMutation::Update {
3406 entity: ENTITY_NAME.to_string(),
3407 key: InputValue::Nat64(2),
3408 patch: dynamic_payload_patch(30),
3409 },
3410 ])
3411 .expect("distinct transfer patches should share one atomic batch");
3412 assert_eq!(
3413 transfer.rows,
3414 vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
3415 "the transfer must preserve the exact total quantity",
3416 );
3417
3418 let merge = session
3419 .execute_trusted_dynamic_mutation_batch(vec![
3420 DynamicMutation::Delete {
3421 entity: ENTITY_NAME.to_string(),
3422 key: InputValue::Nat64(2),
3423 },
3424 DynamicMutation::Update {
3425 entity: ENTITY_NAME.to_string(),
3426 key: InputValue::Nat64(1),
3427 patch: dynamic_payload_patch(100),
3428 },
3429 ])
3430 .expect("two holdings should merge atomically");
3431 assert_eq!(
3432 merge.rows,
3433 vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
3434 "delete before-images and update after-images must retain input order",
3435 );
3436
3437 let resplit = session
3438 .execute_trusted_dynamic_mutation_batch(vec![
3439 DynamicMutation::Update {
3440 entity: ENTITY_NAME.to_string(),
3441 key: InputValue::Nat64(1),
3442 patch: dynamic_payload_patch(60),
3443 },
3444 DynamicMutation::Insert {
3445 entity: ENTITY_NAME.to_string(),
3446 patch: dynamic_payload_patch(40),
3447 },
3448 ])
3449 .expect("the merged holding should split again");
3450 assert_eq!(
3451 resplit.rows,
3452 vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
3453 );
3454
3455 let rejected_merge = session
3456 .execute_trusted_dynamic_mutation_batch(vec![
3457 DynamicMutation::Delete {
3458 entity: ENTITY_NAME.to_string(),
3459 key: InputValue::Nat64(3),
3460 },
3461 DynamicMutation::Update {
3462 entity: ENTITY_NAME.to_string(),
3463 key: InputValue::Nat64(99),
3464 patch: dynamic_payload_patch(100),
3465 },
3466 ])
3467 .expect_err("a late missing merge target must preserve the earlier staged delete");
3468 assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
3469 assert_dynamic_payload(&session, 1, 60);
3470 assert_dynamic_payload(&session, 3, 40);
3471
3472 SCHEMA_STORE.with(|store| {
3473 let cursor = store
3474 .borrow()
3475 .identity_statement_cursor(
3476 database_incarnation_id().expect("database incarnation should remain readable"),
3477 ENTITY_TAG,
3478 FieldId::new(1),
3479 &AcceptedFieldKind::Nat64,
3480 )
3481 .expect("mixed Identity state should remain readable");
3482 assert_eq!(cursor.expected_high_water(), 3);
3483 assert!(!cursor.has_allocations());
3484 });
3485 }
3486
3487 #[test]
3488 fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
3489 let session = initialize();
3490 session
3491 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3492 .expect("the holding fixture should initialize");
3493
3494 let duplicate = session
3495 .execute_trusted_dynamic_mutation_batch(vec![
3496 DynamicMutation::Update {
3497 entity: ENTITY_NAME.to_string(),
3498 key: InputValue::Nat64(1),
3499 patch: dynamic_payload_patch(60),
3500 },
3501 DynamicMutation::Delete {
3502 entity: ENTITY_NAME.to_string(),
3503 key: InputValue::Nat64(1),
3504 },
3505 ])
3506 .expect_err("duplicate targets across operation kinds must reject");
3507 assert!(matches!(
3508 duplicate.diagnostic().detail(),
3509 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3510 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
3511 }),
3512 ));
3513 assert_eq!(
3514 duplicate.diagnostic_facts(),
3515 vec![
3516 (
3517 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3518 ENTITY_TAG.value(),
3519 ),
3520 (
3521 icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
3522 0,
3523 ),
3524 (
3525 icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
3526 1,
3527 ),
3528 ],
3529 );
3530 assert_dynamic_payload(&session, 1, 100);
3531 }
3532
3533 #[test]
3534 fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
3535 let session = initialize();
3536 let empty = session
3537 .execute_trusted_dynamic_mutation_batch(Vec::new())
3538 .expect_err("an empty public batch must reject");
3539 assert!(matches!(
3540 empty.diagnostic().detail(),
3541 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3542 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
3543 }),
3544 ));
3545 assert_eq!(
3546 empty.diagnostic_facts(),
3547 vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
3548 );
3549
3550 let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
3551 .map(|_| DynamicMutation::Delete {
3552 entity: ENTITY_NAME.to_string(),
3553 key: InputValue::Nat64(1),
3554 })
3555 .collect();
3556 let over_bound = session
3557 .execute_trusted_dynamic_mutation_batch(requests)
3558 .expect_err("operation cap plus one must reject before row resolution");
3559 assert!(matches!(
3560 over_bound.diagnostic().detail(),
3561 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3562 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
3563 }),
3564 ));
3565 assert_eq!(
3566 over_bound.diagnostic_facts(),
3567 vec![
3568 (
3569 icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
3570 (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
3571 ),
3572 (
3573 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3574 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
3575 ),
3576 ],
3577 );
3578 }
3579
3580 #[test]
3581 fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
3582 let mut exact = 0;
3583 add_structural_mutation_staged_bytes(
3584 &mut exact,
3585 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
3586 )
3587 .expect("the exact staged-byte boundary should admit");
3588 assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
3589
3590 let error = add_structural_mutation_staged_bytes(&mut exact, [1])
3591 .expect_err("one byte above the staged-byte boundary must reject");
3592 assert!(matches!(
3593 error.diagnostic().detail(),
3594 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3595 boundary:
3596 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
3597 }),
3598 ));
3599 assert_eq!(
3600 error.diagnostic_facts(),
3601 vec![
3602 (
3603 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3604 (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
3605 ),
3606 (
3607 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3608 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
3609 ),
3610 ],
3611 );
3612
3613 validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
3614 .expect("the exact result-byte boundary should admit");
3615 let error = validate_structural_mutation_result_bytes(
3616 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
3617 )
3618 .expect_err("one byte above the result-byte boundary must reject");
3619 assert!(matches!(
3620 error.diagnostic().detail(),
3621 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3622 boundary:
3623 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
3624 }),
3625 ));
3626 assert_eq!(
3627 error.diagnostic_facts(),
3628 vec![
3629 (
3630 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3631 (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
3632 ),
3633 (
3634 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3635 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
3636 ),
3637 ],
3638 );
3639 }
3640
3641 #[expect(
3642 clippy::too_many_lines,
3643 reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
3644 )]
3645 #[test]
3646 fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
3647 let session = initialize();
3648 let catalog = session
3649 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3650 .expect("identity catalog should resolve");
3651 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3652 .expect("identity row layout should build");
3653 let initial_description = session
3654 .try_describe_entity_by_name(ENTITY_NAME)
3655 .expect("accepted Identity description should resolve");
3656 assert_eq!(
3657 initial_description.entity_tag(),
3658 catalog.identity().entity_tag().value()
3659 );
3660 assert_eq!(
3661 initial_description.accepted_schema_fingerprint_method(),
3662 catalog.fingerprint_method_version()
3663 );
3664 assert_eq!(
3665 initial_description.accepted_schema_fingerprint(),
3666 catalog.fingerprint()
3667 );
3668 let initial_identity = initial_description
3669 .identity()
3670 .expect("accepted Identity policy should be described");
3671 assert_eq!(initial_identity.field(), "id");
3672 assert_eq!(initial_identity.generator(), "Identity::next");
3673 assert_eq!(initial_identity.accepted_kind(), "nat64");
3674 assert_eq!(initial_identity.minimum(), 1);
3675 assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
3676 assert_eq!(initial_identity.high_water(), 0);
3677 assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
3678 assert!(!initial_identity.exhausted());
3679
3680 let rejected = session
3681 .execute_accepted_structural_save_batch(
3682 &catalog,
3683 &descriptor,
3684 batch(&[1_000, 2_000]),
3685 Timestamp::from_millis(6),
3686 |_| Err::<(), _>(InternalError::executor_unsupported()),
3687 )
3688 .expect_err("a rejected precommit result must not publish its tentative range");
3689 assert_eq!(rejected.class(), ErrorClass::Unsupported);
3690 assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
3691
3692 let rows = session
3693 .execute_accepted_structural_save_batch(
3694 &catalog,
3695 &descriptor,
3696 batch(&[10, 20, 30]),
3697 Timestamp::from_millis(7),
3698 Ok,
3699 )
3700 .expect("one accepted batch should commit rows and one identity range");
3701 assert_eq!(
3702 rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
3703 vec![
3704 vec![Value::Nat64(1), Value::Nat64(10)],
3705 vec![Value::Nat64(2), Value::Nat64(20)],
3706 vec![Value::Nat64(3), Value::Nat64(30)],
3707 ],
3708 );
3709
3710 let dynamic = session
3711 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3712 entity: ENTITY_NAME.to_string(),
3713 patch: DynamicStructuralPatch::new(vec![(
3714 "payload".to_string(),
3715 DynamicWriteCell::Value(InputValue::Nat64(40)),
3716 )]),
3717 })
3718 .expect("dynamic omission should commit through shared Identity generation");
3719 assert_eq!(dynamic.affected_rows, 1);
3720
3721 for (request, operation) in [
3722 (
3723 DynamicMutation::Insert {
3724 entity: ENTITY_NAME.to_string(),
3725 patch: DynamicStructuralPatch::new(vec![
3726 (
3727 "id".to_string(),
3728 DynamicWriteCell::Value(InputValue::Nat64(41)),
3729 ),
3730 (
3731 "payload".to_string(),
3732 DynamicWriteCell::Value(InputValue::Nat64(42)),
3733 ),
3734 ]),
3735 },
3736 icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
3737 ),
3738 (
3739 DynamicMutation::Update {
3740 entity: ENTITY_NAME.to_string(),
3741 key: InputValue::Nat64(1),
3742 patch: DynamicStructuralPatch::new(vec![(
3743 "id".to_string(),
3744 DynamicWriteCell::Default,
3745 )]),
3746 },
3747 icydb_diagnostic_code::DiagnosticMutationOperation::Update,
3748 ),
3749 ] {
3750 let error = session
3751 .execute_trusted_dynamic_mutation(&request)
3752 .expect_err("structural Identity authorship and regeneration must reject");
3753 assert_eq!(error.class(), ErrorClass::Unsupported);
3754 assert_eq!(error.origin(), ErrorOrigin::Executor);
3755 assert_eq!(
3756 error.diagnostic_facts(),
3757 vec![
3758 (
3759 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3760 ENTITY_TAG.value(),
3761 ),
3762 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3763 (
3764 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3765 operation.raw(),
3766 ),
3767 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3768 ],
3769 );
3770 }
3771
3772 let binding = session
3773 .issue_typed_entity_binding(
3774 ENTITY_SOURCE,
3775 &[
3776 DynamicTypedFieldBindingRequest::new(
3777 ID_SOURCE.to_string(),
3778 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3779 false,
3780 ),
3781 DynamicTypedFieldBindingRequest::new(
3782 PAYLOAD_SOURCE.to_string(),
3783 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3784 false,
3785 ),
3786 ],
3787 )
3788 .expect("typed output should bind the Identity field");
3789 let typed_patch = binding
3790 .bind_write_fields(vec![(
3791 PAYLOAD_SOURCE.to_string(),
3792 DynamicWriteCell::Value(InputValue::Nat64(50)),
3793 )])
3794 .expect("typed payload should lower");
3795 let typed = session
3796 .execute_trusted_typed_mutation(
3797 &binding,
3798 &DynamicTypedMutation::Insert { patch: typed_patch },
3799 )
3800 .expect("typed omission should commit through shared Identity generation");
3801 assert_eq!(
3802 typed
3803 .expect("typed insert should return one mutation result")
3804 .affected_rows,
3805 1,
3806 );
3807 let explicit_typed_patch = binding
3808 .bind_write_fields(vec![
3809 (
3810 ID_SOURCE.to_string(),
3811 DynamicWriteCell::Value(InputValue::Nat64(51)),
3812 ),
3813 (
3814 PAYLOAD_SOURCE.to_string(),
3815 DynamicWriteCell::Value(InputValue::Nat64(52)),
3816 ),
3817 ])
3818 .expect("the low-level binding should retain exact authored intent");
3819 let explicit_typed_error = session
3820 .execute_trusted_typed_mutation(
3821 &binding,
3822 &DynamicTypedMutation::Insert {
3823 patch: explicit_typed_patch,
3824 },
3825 )
3826 .expect_err("typed Identity authorship must reject before allocation");
3827 assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
3828 assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
3829 assert_eq!(
3830 explicit_typed_error.diagnostic_facts(),
3831 vec![
3832 (
3833 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3834 ENTITY_TAG.value(),
3835 ),
3836 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3837 (
3838 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3839 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3840 ),
3841 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3842 ],
3843 );
3844
3845 let replace_error = session
3846 .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
3847 entity: ENTITY_NAME.to_string(),
3848 key: InputValue::Nat64(99),
3849 patch: DynamicStructuralPatch::new(vec![(
3850 "payload".to_string(),
3851 DynamicWriteCell::Value(InputValue::Nat64(60)),
3852 )]),
3853 })
3854 .expect_err("save-as-insert with a chosen Identity must reject");
3855 assert_eq!(replace_error.class(), ErrorClass::Unsupported);
3856 assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
3857
3858 #[cfg(feature = "sql")]
3859 {
3860 for sql in [
3861 "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
3862 "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
3863 ] {
3864 let _result = session
3865 .execute_trusted_sql_mutation(sql)
3866 .expect("SQL omission and DEFAULT should commit Identity generation");
3867 }
3868
3869 let error = session
3870 .execute_trusted_sql_mutation(
3871 "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
3872 )
3873 .expect_err("an explicit SQL Identity value must reject before allocation");
3874 let diagnostic = error.diagnostic();
3875 assert_eq!(
3876 diagnostic.code(),
3877 icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
3878 );
3879 assert!(matches!(
3880 diagnostic.detail(),
3881 Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
3882 boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
3883 }),
3884 ));
3885 }
3886
3887 let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
3888 assert_eq!(
3889 DATA_STORE.with(|store| store.borrow().len()),
3890 expected_committed
3891 );
3892 SCHEMA_STORE.with(|store| {
3893 let cursor = store
3894 .borrow()
3895 .identity_statement_cursor(
3896 database_incarnation_id().expect("database incarnation should remain readable"),
3897 ENTITY_TAG,
3898 FieldId::new(1),
3899 &AcceptedFieldKind::Nat64,
3900 )
3901 .expect("committed writes must leave active state readable");
3902 assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
3903 assert!(!cursor.has_allocations());
3904 });
3905 let committed_description = session
3906 .try_describe_entity_by_name(ENTITY_NAME)
3907 .expect("committed Identity description should resolve");
3908 let committed_identity = committed_description
3909 .identity()
3910 .expect("accepted Identity policy should remain described");
3911 assert_eq!(
3912 committed_identity.high_water(),
3913 u128::from(expected_committed),
3914 );
3915 assert_eq!(
3916 committed_identity.remaining(),
3917 u128::from(u64::MAX - expected_committed),
3918 );
3919 assert!(!committed_identity.exhausted());
3920 }
3921
3922 #[test]
3923 #[expect(
3924 clippy::too_many_lines,
3925 reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
3926 )]
3927 fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
3928 let session = initialize_journaled();
3929 let catalog = session
3930 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3931 .expect("journaled identity catalog should resolve");
3932 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3933 .expect("journaled identity row layout should build");
3934
3935 for (ordinal, interruption) in [
3936 MutationCommitInterruption::MarkerPersisted,
3937 MutationCommitInterruption::JournalPublished,
3938 MutationCommitInterruption::RowsPublished,
3939 MutationCommitInterruption::StateMaterialized,
3940 ]
3941 .into_iter()
3942 .enumerate()
3943 {
3944 interrupt_next_mutation_commit_for_tests(interruption);
3945 let interrupted = session.execute_accepted_structural_save_batch(
3946 &catalog,
3947 &descriptor,
3948 batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
3949 Timestamp::from_millis(8),
3950 Ok,
3951 );
3952 assert!(
3953 interrupted.is_err(),
3954 "the selected durable boundary should interrupt",
3955 );
3956
3957 let committed = session
3958 .execute_accepted_structural_save_batch(
3959 &catalog,
3960 &descriptor,
3961 batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
3962 Timestamp::from_millis(9),
3963 Ok,
3964 )
3965 .expect("the next mutation must recover before allocating");
3966 let expected_high_water =
3967 u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
3968 assert_eq!(
3969 committed
3970 .into_iter()
3971 .map(|row| row.values)
3972 .collect::<Vec<_>>(),
3973 vec![vec![
3974 Value::Nat64(expected_high_water),
3975 Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
3976 ]],
3977 );
3978 assert_eq!(
3979 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
3980 expected_high_water,
3981 );
3982 JOURNALED_SCHEMA_STORE.with(|store| {
3983 let cursor = store
3984 .borrow()
3985 .identity_statement_cursor(
3986 database_incarnation_id()
3987 .expect("database incarnation should remain readable"),
3988 ENTITY_TAG,
3989 FieldId::new(1),
3990 &AcceptedFieldKind::Nat64,
3991 )
3992 .expect("guarded recovery must leave quiescent active state");
3993 assert_eq!(
3994 cursor.expected_high_water(),
3995 u128::from(expected_high_water),
3996 );
3997 assert!(!cursor.has_allocations());
3998 });
3999 }
4000
4001 for (ordinal, (interruption, deleted_key)) in [
4002 (MutationCommitInterruption::MarkerPersisted, 2),
4003 (MutationCommitInterruption::JournalPublished, 4),
4004 (MutationCommitInterruption::RowPrefixPublished, 6),
4005 (MutationCommitInterruption::RowsPublished, 8),
4006 (MutationCommitInterruption::StateMaterialized, 7),
4007 ]
4008 .into_iter()
4009 .enumerate()
4010 {
4011 let expected_payload =
4012 501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
4013 interrupt_next_mutation_commit_for_tests(interruption);
4014 let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
4015 DynamicMutation::Update {
4016 entity: ENTITY_NAME.to_string(),
4017 key: InputValue::Nat64(1),
4018 patch: dynamic_payload_patch(expected_payload),
4019 },
4020 DynamicMutation::Delete {
4021 entity: ENTITY_NAME.to_string(),
4022 key: InputValue::Nat64(deleted_key),
4023 },
4024 ]);
4025 assert!(
4026 interrupted.is_err(),
4027 "the selected caller-key mixed publication boundary should interrupt",
4028 );
4029 let recovered_update = session
4030 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
4031 entity: ENTITY_NAME.to_string(),
4032 key: InputValue::Nat64(1),
4033 patch: dynamic_payload_patch(expected_payload),
4034 })
4035 .expect("guarded reentry should complete the marker-authorized mixed batch");
4036 assert_eq!(
4037 recovered_update.affected_rows, 0,
4038 "the recovered update must already expose its admitted final image",
4039 );
4040 let recovered_delete = session
4041 .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
4042 entity: ENTITY_NAME.to_string(),
4043 key: InputValue::Nat64(deleted_key),
4044 })
4045 .expect_err("the recovered delete must already be materialized");
4046 assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
4047 JOURNALED_SCHEMA_STORE.with(|store| {
4048 let cursor = store
4049 .borrow()
4050 .identity_statement_cursor(
4051 database_incarnation_id()
4052 .expect("database incarnation should remain readable"),
4053 ENTITY_TAG,
4054 FieldId::new(1),
4055 &AcceptedFieldKind::Nat64,
4056 )
4057 .expect("caller-key recovery must preserve active Identity state");
4058 assert_eq!(cursor.expected_high_water(), 8);
4059 assert!(!cursor.has_allocations());
4060 });
4061 }
4062
4063 forget_recovered_domain_for_tests(&session.db)
4064 .expect("the final journal tail should remain recoverable");
4065 session
4066 .db
4067 .ensure_recovered_state()
4068 .expect("derived rebuild must not allocate another identity");
4069
4070 let quick = execute_quick_integrity(&session.db, catalog.inspection_plan())
4071 .expect("quiescent Identity control inventory should be inspectable");
4072 assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
4073 let row_page = execute_row_integrity_page(
4074 &session.db,
4075 catalog.inspection_plan(),
4076 PhysicalUnitCheckpoint::BeforeFirst,
4077 RowInspectionLimits::standard(),
4078 )
4079 .expect("Identity rows should remain within committed high-water");
4080 assert!(row_page.exhausted());
4081 assert!(row_page.findings().is_empty());
4082
4083 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
4084 assert!(
4085 JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
4086 "derived index rebuild should restore witnesses without allocating identities",
4087 );
4088 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
4089 JOURNALED_SCHEMA_STORE.with(|store| {
4090 let cursor = store
4091 .borrow()
4092 .identity_statement_cursor(
4093 database_incarnation_id().expect("database incarnation should remain readable"),
4094 ENTITY_TAG,
4095 FieldId::new(1),
4096 &AcceptedFieldKind::Nat64,
4097 )
4098 .expect("folded identity state should reopen without allocating");
4099 assert_eq!(cursor.expected_high_water(), 8);
4100 assert!(!cursor.has_allocations());
4101 });
4102 }
4103
4104 #[test]
4105 #[ignore = "release-closeout native timing probe for one marker-authorized Identity recovery"]
4106 fn identity_recovery_closeout_reports_guarded_reentry_time() {
4107 let session = initialize_journaled();
4108 let catalog = session
4109 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4110 .expect("journaled identity catalog should resolve");
4111 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4112 .expect("journaled identity row layout should build");
4113
4114 interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
4115 let interrupted = session.execute_accepted_structural_save_batch(
4116 &catalog,
4117 &descriptor,
4118 batch(&[1]),
4119 Timestamp::from_millis(10),
4120 Ok,
4121 );
4122 assert!(
4123 interrupted.is_err(),
4124 "the selected publication boundary should interrupt",
4125 );
4126
4127 let start = Instant::now();
4128 let committed = session
4129 .execute_accepted_structural_save_batch(
4130 &catalog,
4131 &descriptor,
4132 batch(&[2]),
4133 Timestamp::from_millis(11),
4134 Ok,
4135 )
4136 .expect("guarded reentry should recover before allocation");
4137 let elapsed = start.elapsed();
4138 assert_eq!(
4139 committed
4140 .into_iter()
4141 .map(|row| row.values)
4142 .collect::<Vec<_>>(),
4143 vec![vec![Value::Nat64(2), Value::Nat64(2)]],
4144 );
4145
4146 println!(
4147 "identity recovery closeout: guarded_reentry_nanos={}",
4148 elapsed.as_nanos(),
4149 );
4150 }
4151}
4152
4153#[cfg(test)]
4154mod targeted_rule_mutation_tests {
4155 use super::{
4156 DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
4157 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
4158 };
4159 use crate::{
4160 db::{
4161 data::{DataStore, encode_input_value_for_candidate_field_contract},
4162 index::IndexStore,
4163 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
4164 schema::{
4165 AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
4166 AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
4167 AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
4168 ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
4169 PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
4170 ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
4171 SchemaRowLayout, SchemaStore, SchemaVersion,
4172 accepted_schema_candidate_with_catalogs_for_tests,
4173 build_record_newtype_composite_catalog_for_tests,
4174 empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
4175 },
4176 },
4177 error::InternalError,
4178 traits::{CanisterKind, Path},
4179 types::EntityTag,
4180 value::InputValue,
4181 };
4182 use icydb_schema::{
4183 ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
4184 };
4185 use std::{cell::RefCell, collections::BTreeMap};
4186
4187 const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
4188 const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
4189 const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
4190 const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
4191 const UPDATED_AT_SOURCE: &str =
4192 "session::write::targeted_rule_mutation_tests::Entity::updated_at";
4193 const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
4194 const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
4195 const DEGREE_MEMBER_SOURCE: &str =
4196 "session::write::targeted_rule_mutation_tests::Profile::degree";
4197 const DEGREE_RULE_SOURCE: &str =
4198 "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
4199
4200 struct TestCanister;
4201
4202 impl Path for TestCanister {
4203 const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
4204 }
4205
4206 impl CanisterKind for TestCanister {
4207 const COMMIT_MEMORY_ID: u8 = 43;
4208 const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
4209 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
4210 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
4211 "icydb.targeted_mutation_tests.integrity.progress.v1";
4212 }
4213
4214 thread_local! {
4215 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
4216 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
4217 static SCHEMA_STORE: RefCell<SchemaStore> =
4218 const { RefCell::new(SchemaStore::init_heap()) };
4219 static STORE_REGISTRY: StoreRegistry = {
4220 let mut registry = StoreRegistry::new();
4221 registry.register_store(
4222 STORE_PATH,
4223 &DATA_STORE,
4224 &INDEX_STORE,
4225 &SCHEMA_STORE,
4226 StoreAllocationIdentities::absent(),
4227 StoreRuntimeStorageCapabilities::heap(),
4228 ).expect("targeted mutation test store should register");
4229 registry
4230 };
4231 }
4232
4233 fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
4234 parse(raw.to_string()).expect("test source identity should admit")
4235 }
4236
4237 fn profile_input(degree: u64) -> InputValue {
4238 InputValue::Map(vec![(
4239 InputValue::Text("degree".to_string()),
4240 InputValue::Nat64(degree),
4241 )])
4242 }
4243
4244 fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
4245 DynamicStructuralPatch::new(vec![
4246 (
4247 "id".to_string(),
4248 DynamicWriteCell::Value(InputValue::Nat64(id)),
4249 ),
4250 (
4251 "profile".to_string(),
4252 DynamicWriteCell::Value(profile_input(degree)),
4253 ),
4254 ])
4255 }
4256
4257 fn encoded_value(
4258 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4259 composite_catalog: &AcceptedCompositeCatalog,
4260 name: &str,
4261 kind: &AcceptedFieldKind,
4262 storage_decode: FieldStorageDecode,
4263 leaf_codec: LeafCodec,
4264 value: InputValue,
4265 ) -> Vec<u8> {
4266 let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
4267 encode_input_value_for_candidate_field_contract(
4268 enum_catalog,
4269 composite_catalog,
4270 field,
4271 value,
4272 &mut ValueAdmissionBudget::standard(),
4273 )
4274 .expect("test accepted value should encode")
4275 }
4276
4277 fn nat64_literal(
4278 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4279 composite_catalog: &AcceptedCompositeCatalog,
4280 value: u64,
4281 ) -> AcceptedCheckLiteralV1 {
4282 let kind = AcceptedFieldKind::Nat64;
4283 AcceptedCheckLiteralV1::from_accepted_parts(
4284 kind.clone(),
4285 FieldStorageDecode::ByKind,
4286 LeafCodec::Scalar(ScalarCodec::Nat64),
4287 encoded_value(
4288 enum_catalog,
4289 composite_catalog,
4290 "degree_bound",
4291 &kind,
4292 FieldStorageDecode::ByKind,
4293 LeafCodec::Scalar(ScalarCodec::Nat64),
4294 InputValue::Nat64(value),
4295 ),
4296 )
4297 }
4298
4299 fn targeted_constraint_id(error: &InternalError) -> u32 {
4300 let facts = error.diagnostic_facts();
4301 assert!(facts.contains(&(
4302 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4303 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4304 )));
4305 assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
4306 assert!(facts.contains(&(
4307 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
4308 icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
4309 )));
4310 assert_eq!(
4311 facts
4312 .iter()
4313 .filter(|(tag, _)| matches!(
4314 tag,
4315 icydb_diagnostic_code::DiagnosticFactTag::RootField
4316 | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
4317 ))
4318 .copied()
4319 .collect::<Vec<_>>(),
4320 vec![
4321 (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
4322 (
4323 icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
4324 icydb_diagnostic_code::pack_u32_pair(1, 1),
4325 ),
4326 ]
4327 );
4328 let value = facts
4329 .iter()
4330 .find_map(|(tag, value)| {
4331 (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
4332 })
4333 .expect("targeted mutation should retain its accepted constraint ID");
4334 u32::try_from(value).expect("accepted constraint ID fits u32")
4335 }
4336
4337 #[expect(
4338 clippy::too_many_lines,
4339 reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
4340 )]
4341 #[test]
4342 fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
4343 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
4344 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
4345 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
4346
4347 let entity_tag = EntityTag::new(93);
4348 let enum_catalog = empty_accepted_enum_catalog_for_tests();
4349 let (composite_catalog, profile_type, degree_type, degree_member) =
4350 build_record_newtype_composite_catalog_for_tests(
4351 "tests::TargetedProfile".to_string(),
4352 "degree".to_string(),
4353 "tests::TargetedDegree".to_string(),
4354 AcceptedFieldKind::Nat64,
4355 &enum_catalog,
4356 )
4357 .expect("targeted mutation composites should close");
4358 let profile_kind = AcceptedFieldKind::Composite {
4359 type_id: profile_type,
4360 };
4361 let profile_default = encoded_value(
4362 &enum_catalog,
4363 &composite_catalog,
4364 "profile",
4365 &profile_kind,
4366 FieldStorageDecode::CatalogValue,
4367 LeafCodec::Structural,
4368 profile_input(12),
4369 );
4370 let fields = vec![
4371 PersistedFieldSnapshot::new_initial(
4372 FieldId::new(1),
4373 "id".to_string(),
4374 SchemaFieldSlot::new(0),
4375 AcceptedFieldKind::Nat64,
4376 Vec::new(),
4377 false,
4378 SchemaInsertDefault::None,
4379 FieldStorageDecode::ByKind,
4380 LeafCodec::Scalar(ScalarCodec::Nat64),
4381 ),
4382 PersistedFieldSnapshot::new_initial(
4383 FieldId::new(2),
4384 "profile".to_string(),
4385 SchemaFieldSlot::new(1),
4386 profile_kind,
4387 vec![PersistedNestedLeafSnapshot::new(
4388 vec!["degree".to_string()],
4389 AcceptedFieldKind::Composite {
4390 type_id: degree_type,
4391 },
4392 false,
4393 )],
4394 false,
4395 SchemaInsertDefault::SlotPayload(profile_default),
4396 FieldStorageDecode::CatalogValue,
4397 LeafCodec::Structural,
4398 ),
4399 PersistedFieldSnapshot::new_initial_with_write_policy(
4400 FieldId::new(3),
4401 "updated_at".to_string(),
4402 SchemaFieldSlot::new(2),
4403 AcceptedFieldKind::Timestamp,
4404 Vec::new(),
4405 false,
4406 SchemaInsertDefault::None,
4407 SchemaFieldWritePolicy::from_model_policies(
4408 None,
4409 Some(FieldWriteManagement::UpdatedAt),
4410 ),
4411 FieldStorageDecode::ByKind,
4412 LeafCodec::Scalar(ScalarCodec::Timestamp),
4413 ),
4414 ];
4415 let mut snapshot = PersistedSchemaSnapshot::new(
4416 SchemaVersion::initial(),
4417 ENTITY_SOURCE.to_string(),
4418 "TargetedMutation".to_string(),
4419 FieldId::new(1),
4420 SchemaRowLayout::initial(
4421 fields
4422 .iter()
4423 .map(|field| (field.id(), field.slot()))
4424 .collect(),
4425 ),
4426 fields,
4427 );
4428 let constraint_catalog = snapshot
4429 .constraint_catalog()
4430 .clone()
4431 .with_added_targeted_rule(
4432 "profile_degree_multiple".to_string(),
4433 ConstraintOrigin::Generated,
4434 AcceptedRuleTarget::new(
4435 FieldId::new(2),
4436 AcceptedNamedTypeIdentity::Composite(degree_type),
4437 ),
4438 AcceptedRuleOperation::MultipleOf {
4439 divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
4440 },
4441 )
4442 .expect("targeted mutation rule should allocate");
4443 let targeted_rule_id = constraint_catalog
4444 .constraints()
4445 .last()
4446 .expect("targeted mutation rule should persist")
4447 .id();
4448 snapshot = snapshot.with_constraint_catalog(constraint_catalog);
4449
4450 let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
4451 let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
4452 let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
4453 let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
4454 let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
4455 let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
4456 let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
4457 let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
4458 let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
4459 BTreeMap::from([(entity_source, entity_tag)]),
4460 BTreeMap::from([
4461 ((entity_tag, id_source), FieldId::new(1)),
4462 ((entity_tag, profile_source), FieldId::new(2)),
4463 ((entity_tag, updated_at_source), FieldId::new(3)),
4464 ]),
4465 BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
4466 BTreeMap::new(),
4467 BTreeMap::new(),
4468 )
4469 .with_initial_named_types_for_tests(
4470 BTreeMap::from([
4471 (
4472 profile_type_source,
4473 AcceptedNamedTypeIdentity::Composite(profile_type),
4474 ),
4475 (
4476 degree_type_source,
4477 AcceptedNamedTypeIdentity::Composite(degree_type),
4478 ),
4479 ]),
4480 BTreeMap::new(),
4481 BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
4482 );
4483 let candidate = accepted_schema_candidate_with_catalogs_for_tests(
4484 STORE_PATH,
4485 AcceptedSchemaRevision::INITIAL,
4486 enum_catalog,
4487 composite_catalog,
4488 source_bindings,
4489 BTreeMap::from([(entity_tag, snapshot)]),
4490 );
4491
4492 let session = DbSession::<TestCanister>::new(
4493 &STORE_REGISTRY,
4494 &crate::db::RequestExecutionRoot::__new_runtime_root(),
4495 );
4496 session
4497 .db
4498 .ensure_recovered_state()
4499 .expect("targeted mutation test database should initialize");
4500 let store = session
4501 .db
4502 .store_handle(STORE_PATH)
4503 .expect("targeted mutation test store should resolve");
4504 crate::db::commit::publish_accepted_schema_candidate(
4505 STORE_PATH,
4506 store,
4507 AcceptedSchemaRevision::NONE,
4508 &candidate,
4509 )
4510 .expect("targeted mutation candidate should publish");
4511
4512 let dynamic_error = session
4513 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4514 entity: "TargetedMutation".to_string(),
4515 patch: structural_patch(1, 12),
4516 })
4517 .expect_err("dynamic write must enforce the targeted rule");
4518 assert_eq!(
4519 targeted_constraint_id(&dynamic_error),
4520 targeted_rule_id.get()
4521 );
4522
4523 let binding = session
4524 .issue_typed_entity_binding(
4525 ENTITY_SOURCE,
4526 &[
4527 DynamicTypedFieldBindingRequest::new(
4528 ID_SOURCE.to_string(),
4529 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4530 false,
4531 ),
4532 DynamicTypedFieldBindingRequest::new(
4533 PROFILE_SOURCE.to_string(),
4534 DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
4535 false,
4536 ),
4537 DynamicTypedFieldBindingRequest::new(
4538 UPDATED_AT_SOURCE.to_string(),
4539 DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
4540 false,
4541 ),
4542 ],
4543 )
4544 .expect("targeted typed binding should issue");
4545 let typed_patch = binding
4546 .bind_write_fields(vec![
4547 (
4548 ID_SOURCE.to_string(),
4549 DynamicWriteCell::Value(InputValue::Nat64(2)),
4550 ),
4551 (
4552 PROFILE_SOURCE.to_string(),
4553 DynamicWriteCell::Value(profile_input(12)),
4554 ),
4555 ])
4556 .expect("targeted typed patch should bind");
4557 let typed_error = session
4558 .execute_trusted_typed_mutation(
4559 &binding,
4560 &DynamicTypedMutation::Insert { patch: typed_patch },
4561 )
4562 .expect_err("typed write must enforce the targeted rule");
4563 assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
4564
4565 #[cfg(feature = "sql")]
4566 {
4567 let sql_error = session
4568 .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
4569 .expect_err("SQL default resolution must enforce the targeted rule");
4570 let crate::db::QueryError::Execute(execute) = sql_error else {
4571 panic!("targeted SQL write should fail at shared execution admission");
4572 };
4573 assert_eq!(
4574 targeted_constraint_id(execute.as_internal()),
4575 targeted_rule_id.get()
4576 );
4577 }
4578
4579 session
4580 .execute_trusted_dynamic_mutation_batch(vec![
4581 DynamicMutation::Insert {
4582 entity: "TargetedMutation".to_string(),
4583 patch: structural_patch(4, 5),
4584 },
4585 DynamicMutation::Insert {
4586 entity: "TargetedMutation".to_string(),
4587 patch: structural_patch(5, 12),
4588 },
4589 ])
4590 .expect_err("one invalid targeted value must reject the whole batch");
4591 assert_eq!(
4592 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4593 Some(0),
4594 "no frontend or earlier valid batch row may escape targeted admission",
4595 );
4596
4597 let admitted = session
4598 .execute_trusted_dynamic_mutation_batch(vec![
4599 DynamicMutation::Insert {
4600 entity: "TargetedMutation".to_string(),
4601 patch: structural_patch(6, 5),
4602 },
4603 DynamicMutation::Insert {
4604 entity: "TargetedMutation".to_string(),
4605 patch: structural_patch(7, 10),
4606 },
4607 ])
4608 .expect("compliant targeted values should share one accepted batch");
4609 let [first, second] = admitted.rows.as_slice() else {
4610 panic!("the mixed targeted batch should return two rows");
4611 };
4612 let first_timestamp = first
4613 .get(2)
4614 .expect("the first mixed row should contain its managed timestamp");
4615 assert!(matches!(
4616 first_timestamp,
4617 crate::value::OutputValue::Timestamp(_)
4618 ));
4619 assert_eq!(
4620 second.get(2),
4621 Some(first_timestamp),
4622 "one accepted mixed batch must materialize one managed timestamp",
4623 );
4624 assert_eq!(
4625 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4626 Some(2),
4627 );
4628 }
4629}