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(&STORE_REGISTRY);
1600 session
1601 .db
1602 .ensure_recovered_state()
1603 .expect("typed adapter test database should initialize");
1604 publish(
1605 &session,
1606 AcceptedSchemaRevision::NONE,
1607 AcceptedSchemaRevision::INITIAL,
1608 BTreeMap::from([(
1609 entity_tag,
1610 snapshot(
1611 ENTITY_SOURCE,
1612 "Entity",
1613 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
1614 ),
1615 )]),
1616 BTreeMap::from([
1617 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1618 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1619 ]),
1620 );
1621
1622 let initial_catalog = session
1623 .find_accepted_schema_catalog_context_for_entity_source_key(ENTITY_SOURCE)
1624 .expect("initial source catalog lookup should inspect")
1625 .expect("initial source catalog should exist");
1626 assert_eq!(initial_catalog.identity().entity_tag(), entity_tag);
1627 let initial = session
1628 .issue_typed_entity_binding(
1629 entity_source(ENTITY_SOURCE).as_str(),
1630 &[request(ID_SOURCE), request(VALUE_SOURCE)],
1631 )
1632 .expect("initial typed binding should issue");
1633 assert_eq!(initial.field_slot(ID_SOURCE), Some(0));
1634 assert_eq!(initial.field_slot(VALUE_SOURCE), Some(1));
1635 assert_eq!(initial.output_field_slot("value"), Some(1));
1636 let initial_patch = initial
1637 .bind_write_fields(vec![(
1638 VALUE_SOURCE.to_string(),
1639 DynamicWriteCell::Value(InputValue::Nat64(7)),
1640 )])
1641 .expect("source-bound patch should lower");
1642 assert_eq!(
1643 initial_patch.fields(),
1644 &[(2, 1, DynamicWriteCell::Value(InputValue::Nat64(7)))]
1645 );
1646
1647 publish(
1648 &session,
1649 AcceptedSchemaRevision::INITIAL,
1650 AcceptedSchemaRevision::new(2),
1651 BTreeMap::from([
1652 (
1653 entity_tag,
1654 snapshot(
1655 ENTITY_SOURCE,
1656 "RenamedEntity",
1657 vec![
1658 nat64_field(1, "id", 0),
1659 nat64_field(2, "renamed_value", 1),
1660 nat64_field(3, "value", 2),
1661 ],
1662 ),
1663 ),
1664 (
1665 other_entity_tag,
1666 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
1667 ),
1668 ]),
1669 BTreeMap::from([
1670 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1671 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1672 (
1673 (entity_tag, field_source(REPLACEMENT_SOURCE)),
1674 FieldId::new(3),
1675 ),
1676 (
1677 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
1678 FieldId::new(1),
1679 ),
1680 ]),
1681 );
1682
1683 let stale_authority = session
1684 .ensure_accepted_schema_authority_is_current_for_store_path(
1685 STORE_PATH,
1686 initial_catalog.value_catalog_handle().authority(),
1687 )
1688 .expect_err("the initial accepted authority must be stale after revision two");
1689 assert_eq!(
1690 stale_authority.diagnostic_facts(),
1691 vec![
1692 (
1693 icydb_diagnostic_code::DiagnosticFactTag::ExpectedRevision,
1694 AcceptedSchemaRevision::INITIAL.get(),
1695 ),
1696 (
1697 icydb_diagnostic_code::DiagnosticFactTag::CurrentRevision,
1698 AcceptedSchemaRevision::new(2).get(),
1699 ),
1700 ],
1701 );
1702
1703 assert!(
1704 !session
1705 .typed_entity_binding_is_current(&initial)
1706 .expect("renamed binding currentness should inspect")
1707 );
1708 let renamed = session
1709 .issue_typed_entity_binding(ENTITY_SOURCE, &[request(ID_SOURCE), request(VALUE_SOURCE)])
1710 .expect("renamed source-bound adapter should rebind");
1711 assert_eq!(renamed.entity(), "RenamedEntity");
1712 assert_eq!(renamed.field_slot(VALUE_SOURCE), Some(1));
1713 assert_eq!(renamed.output_field_slot("renamed_value"), Some(1));
1714 assert_eq!(renamed.output_field_slot("value"), None);
1715
1716 publish(
1717 &session,
1718 AcceptedSchemaRevision::new(2),
1719 AcceptedSchemaRevision::new(3),
1720 BTreeMap::from([
1721 (
1722 entity_tag,
1723 snapshot(
1724 ENTITY_SOURCE,
1725 "RenamedEntity",
1726 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
1727 ),
1728 ),
1729 (
1730 other_entity_tag,
1731 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
1732 ),
1733 ]),
1734 BTreeMap::from([
1735 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1736 (
1737 (entity_tag, field_source(REPLACEMENT_SOURCE)),
1738 FieldId::new(2),
1739 ),
1740 (
1741 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
1742 FieldId::new(1),
1743 ),
1744 ]),
1745 );
1746
1747 assert!(matches!(
1748 session.issue_typed_entity_binding(
1749 ENTITY_SOURCE,
1750 &[request(ID_SOURCE), request(VALUE_SOURCE)],
1751 ),
1752 Err(DynamicTypedBindingError::FieldUnavailable),
1753 ));
1754 assert!(
1755 !session
1756 .typed_entity_binding_is_current(&renamed)
1757 .expect("removed source binding should become stale")
1758 );
1759
1760 let replacement = session
1761 .issue_typed_entity_binding(
1762 ENTITY_SOURCE,
1763 &[request(ID_SOURCE), request(REPLACEMENT_SOURCE)],
1764 )
1765 .expect("explicit replacement source should bind");
1766 assert!(
1767 session
1768 .execute_trusted_typed_mutation(
1769 &replacement,
1770 &DynamicTypedMutation::Insert {
1771 patch: initial_patch
1772 },
1773 )
1774 .expect("cross-binding patch should fail closed")
1775 .is_none()
1776 );
1777 let patch = replacement
1778 .bind_write_fields(vec![
1779 (
1780 ID_SOURCE.to_string(),
1781 DynamicWriteCell::Value(InputValue::Nat64(1)),
1782 ),
1783 (
1784 REPLACEMENT_SOURCE.to_string(),
1785 DynamicWriteCell::Value(InputValue::Nat64(9)),
1786 ),
1787 ])
1788 .expect("replacement source write should bind by accepted IDs and slots");
1789 let result = session
1790 .execute_trusted_typed_mutation(&replacement, &DynamicTypedMutation::Insert { patch })
1791 .expect("typed insert should use the accepted mutation pipeline")
1792 .expect("replacement binding should remain current");
1793 assert_eq!(result.entity, "RenamedEntity");
1794 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
1795 assert_eq!(
1796 result.rows,
1797 vec![vec![
1798 crate::value::OutputValue::Nat64(1),
1799 crate::value::OutputValue::Nat64(9)
1800 ]]
1801 );
1802 assert_eq!(result.affected_rows, 1);
1803
1804 let second_patch = replacement
1805 .bind_write_fields(vec![
1806 (
1807 ID_SOURCE.to_string(),
1808 DynamicWriteCell::Value(InputValue::Nat64(2)),
1809 ),
1810 (
1811 REPLACEMENT_SOURCE.to_string(),
1812 DynamicWriteCell::Value(InputValue::Nat64(10)),
1813 ),
1814 ])
1815 .expect("second source-bound patch should lower");
1816 session
1817 .execute_trusted_typed_mutation(
1818 &replacement,
1819 &DynamicTypedMutation::Insert {
1820 patch: second_patch,
1821 },
1822 )
1823 .expect("second typed insert should use the accepted mutation pipeline")
1824 .expect("replacement binding should remain current");
1825
1826 {
1827 let query = crate::db::DynamicQuery::new("RenamedEntity")
1828 .select(["id", "value"])
1829 .order_by(crate::db::asc("id"))
1830 .limit(1);
1831 let result = session
1832 .execute_trusted_dynamic_query(&query)
1833 .expect("SQL-free dynamic execution should use accepted authority");
1834 assert_eq!(result.entity, "RenamedEntity");
1835 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
1836 assert_eq!(
1837 result.rows,
1838 vec![vec![
1839 crate::value::OutputValue::Nat64(1),
1840 crate::value::OutputValue::Nat64(9)
1841 ]]
1842 );
1843 assert_eq!(result.row_count, 1);
1844 assert_query_diagnostic(
1845 session
1846 .execute_trusted_dynamic_query(&query.cursor("00"))
1847 .expect_err("scalar execution must reject grouped cursor state"),
1848 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1849 icydb_diagnostic_code::ErrorOrigin::Query,
1850 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1851 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1852 },
1853 );
1854 assert_query_diagnostic(
1855 session
1856 .execute_public_dynamic_grouped_query(
1857 &crate::db::DynamicQuery::new("RenamedEntity").grouped_limits(1, 1024),
1858 )
1859 .expect_err("grouped execution must reject scalar query state"),
1860 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1861 icydb_diagnostic_code::ErrorOrigin::Query,
1862 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1863 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1864 },
1865 );
1866
1867 let grouped_query = crate::db::DynamicQuery::new("RenamedEntity")
1868 .filter(crate::db::FieldRef::new("id").eq(1_u64))
1869 .group_by("value")
1870 .aggregate(crate::db::count())
1871 .grouped_limits(1, 1024)
1872 .limit(1);
1873 let grouped = session
1874 .execute_public_dynamic_grouped_query(&grouped_query)
1875 .expect("SQL-free grouped execution should use accepted authority");
1876 let typed_grouped = session
1877 .execute_public_dynamic_grouped_query_for_typed_binding(
1878 &replacement,
1879 &grouped_query,
1880 )
1881 .expect("typed grouped execution should inspect accepted authority")
1882 .expect("replacement binding should remain current");
1883 assert_eq!(typed_grouped, grouped);
1884 assert!(
1885 session
1886 .execute_public_dynamic_grouped_query_for_typed_binding(
1887 &renamed,
1888 &grouped_query,
1889 )
1890 .expect("stale grouped binding should inspect accepted authority")
1891 .is_none(),
1892 "stale typed grouped bindings must fail closed before execution"
1893 );
1894 assert_eq!(grouped.entity, "RenamedEntity");
1895 assert_eq!(grouped.row_count, 1);
1896 assert_eq!(grouped.rows.len(), 1);
1897 assert_eq!(
1898 grouped.rows[0].group_key(),
1899 &[crate::value::OutputValue::Nat64(9)]
1900 );
1901 assert_eq!(
1902 grouped.rows[0].aggregate_values(),
1903 &[crate::value::OutputValue::Nat64(1)]
1904 );
1905 assert_eq!(grouped.next_cursor, None);
1906
1907 assert_query_diagnostic(
1908 session
1909 .execute_public_dynamic_grouped_query(&grouped_query.clone().select(["value"]))
1910 .expect_err("grouped output must reject scalar selection"),
1911 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
1912 icydb_diagnostic_code::ErrorOrigin::Query,
1913 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1914 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
1915 },
1916 );
1917 assert_query_diagnostic(
1918 session
1919 .execute_public_dynamic_grouped_query(
1920 &crate::db::DynamicQuery::new("RenamedEntity")
1921 .group_by("value")
1922 .aggregate(crate::db::count()),
1923 )
1924 .expect_err("public grouped execution must require explicit limits"),
1925 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1926 icydb_diagnostic_code::ErrorOrigin::Query,
1927 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1928 reason:
1929 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
1930 },
1931 );
1932 assert_query_diagnostic(
1933 session
1934 .execute_trusted_dynamic_grouped_query(
1935 &crate::db::DynamicQuery::new("RenamedEntity")
1936 .group_by("value")
1937 .aggregate(crate::db::count())
1938 .grouped_limits(0, 1024),
1939 )
1940 .expect_err("trusted grouped execution must reject zero limits"),
1941 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1942 icydb_diagnostic_code::ErrorOrigin::Query,
1943 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1944 reason:
1945 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
1946 },
1947 );
1948 assert_query_diagnostic(
1949 session
1950 .execute_public_dynamic_grouped_query(&grouped_query.grouped_limits(101, 1024))
1951 .expect_err("public grouped execution must enforce its group budget"),
1952 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1953 icydb_diagnostic_code::ErrorOrigin::Query,
1954 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1955 reason:
1956 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryExceedsBudget,
1957 },
1958 );
1959
1960 let paged_query = crate::db::DynamicQuery::new("RenamedEntity")
1961 .group_by("value")
1962 .aggregate(crate::db::count())
1963 .grouped_limits(2, 1024)
1964 .limit(1);
1965 assert_query_diagnostic(
1966 session
1967 .execute_public_dynamic_grouped_query(&paged_query)
1968 .expect_err("public grouped execution must reject an unbounded full scan"),
1969 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
1970 icydb_diagnostic_code::ErrorOrigin::Query,
1971 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
1972 reason:
1973 icydb_diagnostic_code::QueryReadAdmissionCode::UnboundedFullScanRejected,
1974 },
1975 );
1976 let first_page = session
1977 .execute_trusted_dynamic_grouped_query(&paged_query)
1978 .expect("SQL-free grouped first page should execute");
1979 assert_eq!(first_page.row_count, 1);
1980 assert_eq!(
1981 first_page.rows[0].group_key(),
1982 &[crate::value::OutputValue::Nat64(9)]
1983 );
1984 let cursor = first_page
1985 .next_cursor
1986 .expect("first grouped page should return a continuation cursor");
1987 assert_query_diagnostic(
1988 session
1989 .execute_trusted_dynamic_grouped_query(
1990 &paged_query.clone().cursor(format!("{cursor}0")),
1991 )
1992 .expect_err("tampered grouped cursor must fail closed"),
1993 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
1994 icydb_diagnostic_code::ErrorOrigin::Cursor,
1995 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
1996 kind: icydb_diagnostic_code::QueryErrorKind::InvalidContinuationCursor,
1997 },
1998 );
1999 let second_page = session
2000 .execute_trusted_dynamic_grouped_query(&paged_query.cursor(cursor))
2001 .expect("SQL-free grouped continuation should execute");
2002 assert_eq!(second_page.row_count, 1);
2003 assert_eq!(
2004 second_page.rows[0].group_key(),
2005 &[crate::value::OutputValue::Nat64(10)]
2006 );
2007 assert_eq!(second_page.next_cursor, None);
2008 }
2009 }
2010}
2011
2012#[cfg(test)]
2013mod mixed_relation_batch_tests {
2014 use super::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell};
2015 use crate::{
2016 db::{
2017 data::DataStore,
2018 index::IndexStore,
2019 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
2020 schema::{
2021 AcceptedConstraintCatalog, AcceptedFieldKind, AcceptedSchemaRevision, FieldId,
2022 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2023 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2024 PersistedRelationEdgeSnapshot, PersistedSchemaSnapshot, RelationId, ScalarCodec,
2025 SchemaFieldSlot, SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore,
2026 SchemaVersion, accepted_schema_candidate_with_field_bindings_for_tests,
2027 },
2028 },
2029 error::ErrorClass,
2030 traits::{CanisterKind, Path},
2031 types::EntityTag,
2032 value::{InputValue, OutputValue},
2033 };
2034 use icydb_schema::FieldSourceKey;
2035 use std::{cell::RefCell, collections::BTreeMap};
2036
2037 const STORE_PATH: &str = "session::write::mixed_relation_batch_tests::Store";
2038 const ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node";
2039 const ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::id";
2040 const PARENT_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::parent_id";
2041 const CODE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::code";
2042 const ENTITY_NAME: &str = "MixedRelationNode";
2043 const ENTITY_TAG: EntityTag = EntityTag::new(94);
2044 const OTHER_ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other";
2045 const OTHER_ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::id";
2046 const OTHER_VALUE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::value";
2047 const OTHER_ENTITY_NAME: &str = "MixedRelationOther";
2048 const OTHER_ENTITY_TAG: EntityTag = EntityTag::new(95);
2049
2050 struct TestCanister;
2051
2052 impl Path for TestCanister {
2053 const PATH: &'static str = "session::write::mixed_relation_batch_tests::Canister";
2054 }
2055
2056 impl CanisterKind for TestCanister {
2057 const COMMIT_MEMORY_ID: u8 = 47;
2058 const COMMIT_STABLE_KEY: &'static str = "icydb.mixed_relation_batch_tests.commit.v1";
2059 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 48;
2060 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2061 "icydb.mixed_relation_batch_tests.integrity.progress.v1";
2062 }
2063
2064 thread_local! {
2065 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2066 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2067 static SCHEMA_STORE: RefCell<SchemaStore> =
2068 const { RefCell::new(SchemaStore::init_heap()) };
2069 static STORE_REGISTRY: StoreRegistry = {
2070 let mut registry = StoreRegistry::new();
2071 registry.register_store(
2072 STORE_PATH,
2073 &DATA_STORE,
2074 &INDEX_STORE,
2075 &SCHEMA_STORE,
2076 StoreAllocationIdentities::absent(),
2077 StoreRuntimeStorageCapabilities::heap(),
2078 ).expect("mixed relation test store should register");
2079 registry
2080 };
2081 }
2082
2083 fn source_key(source: &str) -> FieldSourceKey {
2084 FieldSourceKey::try_new(source).expect("mixed relation field source should admit")
2085 }
2086
2087 fn relation_snapshot() -> PersistedSchemaSnapshot {
2088 let fields = vec![
2089 PersistedFieldSnapshot::new_initial(
2090 FieldId::new(1),
2091 "id".to_string(),
2092 SchemaFieldSlot::new(0),
2093 AcceptedFieldKind::Nat64,
2094 Vec::new(),
2095 false,
2096 SchemaInsertDefault::None,
2097 FieldStorageDecode::ByKind,
2098 LeafCodec::Scalar(ScalarCodec::Nat64),
2099 ),
2100 PersistedFieldSnapshot::new_initial(
2101 FieldId::new(2),
2102 "parent_id".to_string(),
2103 SchemaFieldSlot::new(1),
2104 AcceptedFieldKind::Relation {
2105 target_path: ENTITY_SOURCE.to_string(),
2106 target_entity_name: ENTITY_NAME.to_string(),
2107 target_entity_tag: ENTITY_TAG,
2108 target_store_path: STORE_PATH.to_string(),
2109 key_kind: Box::new(AcceptedFieldKind::Nat64),
2110 },
2111 Vec::new(),
2112 true,
2113 SchemaInsertDefault::None,
2114 FieldStorageDecode::ByKind,
2115 LeafCodec::Scalar(ScalarCodec::Nat64),
2116 ),
2117 PersistedFieldSnapshot::new_initial(
2118 FieldId::new(3),
2119 "code".to_string(),
2120 SchemaFieldSlot::new(2),
2121 AcceptedFieldKind::Nat64,
2122 Vec::new(),
2123 false,
2124 SchemaInsertDefault::None,
2125 FieldStorageDecode::ByKind,
2126 LeafCodec::Scalar(ScalarCodec::Nat64),
2127 ),
2128 ];
2129 let relation = PersistedRelationEdgeSnapshot::new(
2130 RelationId::new(1).expect("mixed relation identity should be non-zero"),
2131 "parent".to_string(),
2132 ENTITY_SOURCE.to_string(),
2133 vec![FieldId::new(2)],
2134 );
2135 let snapshot = PersistedSchemaSnapshot::new_with_indexes(
2136 SchemaVersion::initial(),
2137 ENTITY_SOURCE.to_string(),
2138 ENTITY_NAME.to_string(),
2139 FieldId::new(1),
2140 SchemaRowLayout::initial(
2141 fields
2142 .iter()
2143 .map(|field| (field.id(), field.slot()))
2144 .collect(),
2145 ),
2146 fields,
2147 vec![PersistedIndexSnapshot::new(
2148 SchemaIndexId::new(1).expect("mixed unique index identity should be non-zero"),
2149 1,
2150 "by_code".to_string(),
2151 STORE_PATH.to_string(),
2152 true,
2153 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2154 FieldId::new(3),
2155 SchemaFieldSlot::new(2),
2156 vec!["code".to_string()],
2157 AcceptedFieldKind::Nat64,
2158 false,
2159 )]),
2160 None,
2161 )],
2162 )
2163 .with_relations(vec![relation]);
2164 let constraints = AcceptedConstraintCatalog::initial(
2165 snapshot.fields(),
2166 snapshot.indexes(),
2167 snapshot.relations(),
2168 )
2169 .expect("mixed relation constraints should close");
2170 snapshot.with_constraint_catalog(constraints)
2171 }
2172
2173 fn other_snapshot() -> PersistedSchemaSnapshot {
2174 let fields = vec![
2175 PersistedFieldSnapshot::new_initial(
2176 FieldId::new(1),
2177 "id".to_string(),
2178 SchemaFieldSlot::new(0),
2179 AcceptedFieldKind::Nat64,
2180 Vec::new(),
2181 false,
2182 SchemaInsertDefault::None,
2183 FieldStorageDecode::ByKind,
2184 LeafCodec::Scalar(ScalarCodec::Nat64),
2185 ),
2186 PersistedFieldSnapshot::new_initial(
2187 FieldId::new(2),
2188 "value".to_string(),
2189 SchemaFieldSlot::new(1),
2190 AcceptedFieldKind::Nat64,
2191 Vec::new(),
2192 false,
2193 SchemaInsertDefault::None,
2194 FieldStorageDecode::ByKind,
2195 LeafCodec::Scalar(ScalarCodec::Nat64),
2196 ),
2197 ];
2198 PersistedSchemaSnapshot::new(
2199 SchemaVersion::initial(),
2200 OTHER_ENTITY_SOURCE.to_string(),
2201 OTHER_ENTITY_NAME.to_string(),
2202 FieldId::new(1),
2203 SchemaRowLayout::initial(
2204 fields
2205 .iter()
2206 .map(|field| (field.id(), field.slot()))
2207 .collect(),
2208 ),
2209 fields,
2210 )
2211 }
2212
2213 fn initialize() -> DbSession<TestCanister> {
2214 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2215 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2216 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2217 let session = DbSession::<TestCanister>::new(&STORE_REGISTRY);
2218 session
2219 .db
2220 .ensure_recovered_state()
2221 .expect("mixed relation database should initialize");
2222 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2223 STORE_PATH,
2224 AcceptedSchemaRevision::INITIAL,
2225 BTreeMap::from([
2226 (ENTITY_TAG, relation_snapshot()),
2227 (OTHER_ENTITY_TAG, other_snapshot()),
2228 ]),
2229 BTreeMap::from([
2230 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2231 ((ENTITY_TAG, source_key(PARENT_SOURCE)), FieldId::new(2)),
2232 ((ENTITY_TAG, source_key(CODE_SOURCE)), FieldId::new(3)),
2233 (
2234 (OTHER_ENTITY_TAG, source_key(OTHER_ID_SOURCE)),
2235 FieldId::new(1),
2236 ),
2237 (
2238 (OTHER_ENTITY_TAG, source_key(OTHER_VALUE_SOURCE)),
2239 FieldId::new(2),
2240 ),
2241 ]),
2242 );
2243 let store = session
2244 .db
2245 .store_handle(STORE_PATH)
2246 .expect("mixed relation store should resolve");
2247 crate::db::commit::publish_accepted_schema_candidate(
2248 STORE_PATH,
2249 store,
2250 AcceptedSchemaRevision::NONE,
2251 &candidate,
2252 )
2253 .expect("mixed relation candidate should publish");
2254 session
2255 }
2256
2257 fn patch(id: Option<u64>, parent: Option<u64>, code: Option<u64>) -> DynamicStructuralPatch {
2258 let mut fields = Vec::new();
2259 if let Some(id) = id {
2260 fields.push((
2261 "id".to_string(),
2262 DynamicWriteCell::Value(InputValue::Nat64(id)),
2263 ));
2264 }
2265 fields.push((
2266 "parent_id".to_string(),
2267 parent.map_or(DynamicWriteCell::Null, |parent| {
2268 DynamicWriteCell::Value(InputValue::Nat64(parent))
2269 }),
2270 ));
2271 if let Some(code) = code {
2272 fields.push((
2273 "code".to_string(),
2274 DynamicWriteCell::Value(InputValue::Nat64(code)),
2275 ));
2276 }
2277 DynamicStructuralPatch::new(fields)
2278 }
2279
2280 fn insert(id: u64, parent: Option<u64>) -> DynamicMutation {
2281 insert_with_code(id, parent, id)
2282 }
2283
2284 fn insert_with_code(id: u64, parent: Option<u64>, code: u64) -> DynamicMutation {
2285 DynamicMutation::Insert {
2286 entity: ENTITY_NAME.to_string(),
2287 patch: patch(Some(id), parent, Some(code)),
2288 }
2289 }
2290
2291 fn update_parent(id: u64, parent: Option<u64>) -> DynamicMutation {
2292 DynamicMutation::Update {
2293 entity: ENTITY_NAME.to_string(),
2294 key: InputValue::Nat64(id),
2295 patch: patch(None, parent, None),
2296 }
2297 }
2298
2299 fn update_code(id: u64, code: u64) -> DynamicMutation {
2300 DynamicMutation::Update {
2301 entity: ENTITY_NAME.to_string(),
2302 key: InputValue::Nat64(id),
2303 patch: DynamicStructuralPatch::new(vec![(
2304 "code".to_string(),
2305 DynamicWriteCell::Value(InputValue::Nat64(code)),
2306 )]),
2307 }
2308 }
2309
2310 fn delete(id: u64) -> DynamicMutation {
2311 DynamicMutation::Delete {
2312 entity: ENTITY_NAME.to_string(),
2313 key: InputValue::Nat64(id),
2314 }
2315 }
2316
2317 fn expected_row(id: u64, parent: Option<u64>) -> Vec<OutputValue> {
2318 expected_row_with_code(id, parent, id)
2319 }
2320
2321 fn expected_row_with_code(id: u64, parent: Option<u64>, code: u64) -> Vec<OutputValue> {
2322 vec![
2323 OutputValue::Nat64(id),
2324 parent.map_or(OutputValue::Null, OutputValue::Nat64),
2325 OutputValue::Nat64(code),
2326 ]
2327 }
2328
2329 fn other_patch(id: Option<u64>, value: u64) -> DynamicStructuralPatch {
2330 let mut fields = Vec::new();
2331 if let Some(id) = id {
2332 fields.push((
2333 "id".to_string(),
2334 DynamicWriteCell::Value(InputValue::Nat64(id)),
2335 ));
2336 }
2337 fields.push((
2338 "value".to_string(),
2339 DynamicWriteCell::Value(InputValue::Nat64(value)),
2340 ));
2341 DynamicStructuralPatch::new(fields)
2342 }
2343
2344 fn assert_relation_violation(error: &crate::error::InternalError) {
2345 assert!(error.diagnostic_facts().contains(&(
2346 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
2347 icydb_diagnostic_code::DiagnosticConstraintKind::Relation.raw(),
2348 )));
2349 }
2350
2351 #[test]
2352 fn mixed_relation_validation_uses_the_complete_final_row_overlay() {
2353 let session = initialize();
2354 session
2355 .execute_trusted_dynamic_mutation_batch(vec![insert(1, None), insert(2, Some(1))])
2356 .expect("the initial relation should commit");
2357
2358 let blocked = session
2359 .execute_trusted_dynamic_mutation(&delete(1))
2360 .expect_err("an unaffected committed source must block target deletion");
2361 assert_relation_violation(&blocked);
2362
2363 let deleted = session
2364 .execute_trusted_dynamic_mutation_batch(vec![delete(2), delete(1)])
2365 .expect("a source and its target should delete atomically");
2366 assert_eq!(
2367 deleted.rows,
2368 vec![expected_row(2, Some(1)), expected_row(1, None)],
2369 );
2370
2371 session
2372 .execute_trusted_dynamic_mutation_batch(vec![insert(3, None), insert(4, Some(3))])
2373 .expect("the update-away fixture should commit");
2374 let updated_away = session
2375 .execute_trusted_dynamic_mutation_batch(vec![update_parent(4, None), delete(3)])
2376 .expect("an updated final source may release a deleted target");
2377 assert_eq!(
2378 updated_away.rows,
2379 vec![expected_row(4, None), expected_row(3, None)],
2380 );
2381
2382 session
2383 .execute_trusted_dynamic_mutation_batch(vec![insert(5, None), insert(6, Some(5))])
2384 .expect("the retained-reference fixture should commit");
2385 let retained = session
2386 .execute_trusted_dynamic_mutation_batch(vec![update_parent(6, Some(5)), delete(5)])
2387 .expect_err("a final updated source must still block target deletion");
2388 assert_relation_violation(&retained);
2389
2390 session
2391 .execute_trusted_dynamic_mutation(&insert(7, None))
2392 .expect("the inserted-reference fixture target should commit");
2393 let inserted_reference = session
2394 .execute_trusted_dynamic_mutation_batch(vec![insert(8, Some(7)), delete(7)])
2395 .expect_err("a final inserted source must not reference a deleted target");
2396 assert_relation_violation(&inserted_reference);
2397
2398 let inserted_target = session
2399 .execute_trusted_dynamic_mutation_batch(vec![insert(10, Some(9)), insert(9, None)])
2400 .expect("an inserted relation should see its batch-final target");
2401 assert_eq!(
2402 inserted_target.rows,
2403 vec![expected_row(10, Some(9)), expected_row(9, None)],
2404 );
2405
2406 session
2407 .execute_trusted_dynamic_mutation(&insert(11, None))
2408 .expect("the updated-reference fixture source should commit");
2409 let updated_target = session
2410 .execute_trusted_dynamic_mutation_batch(vec![
2411 update_parent(11, Some(12)),
2412 insert(12, None),
2413 ])
2414 .expect("an updated relation should see its batch-final target");
2415 assert_eq!(
2416 updated_target.rows,
2417 vec![expected_row(11, Some(12)), expected_row(12, None)],
2418 );
2419 }
2420
2421 #[test]
2422 fn mixed_batch_rejects_cross_entity_missing_and_collision_then_honors_replace() {
2423 let session = initialize();
2424 session
2425 .execute_trusted_dynamic_mutation(&insert(1, None))
2426 .expect("the primary mixed fixture row should commit");
2427 session
2428 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
2429 entity: OTHER_ENTITY_NAME.to_string(),
2430 patch: other_patch(Some(1), 10),
2431 })
2432 .expect("the secondary mixed fixture row should commit");
2433
2434 let mixed_entity = session
2435 .execute_trusted_dynamic_mutation_batch(vec![
2436 update_code(1, 11),
2437 DynamicMutation::Update {
2438 entity: OTHER_ENTITY_NAME.to_string(),
2439 key: InputValue::Nat64(1),
2440 patch: other_patch(None, 11),
2441 },
2442 ])
2443 .expect_err("one atomic batch must not cross accepted entities");
2444 assert_eq!(mixed_entity.class(), ErrorClass::Conflict);
2445 assert_eq!(
2446 mixed_entity.diagnostic_facts(),
2447 vec![
2448 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
2449 (
2450 icydb_diagnostic_code::DiagnosticFactTag::ExpectedEntityTag,
2451 ENTITY_TAG.value(),
2452 ),
2453 (
2454 icydb_diagnostic_code::DiagnosticFactTag::ActualEntityTag,
2455 OTHER_ENTITY_TAG.value(),
2456 ),
2457 ],
2458 );
2459
2460 let missing = session
2461 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 12), delete(99)])
2462 .expect_err("a late missing delete must reject the earlier staged update");
2463 assert_eq!(missing.class(), ErrorClass::NotFound);
2464
2465 session
2466 .execute_trusted_dynamic_mutation(&insert(2, None))
2467 .expect("the collision fixture should commit");
2468 let collision = session
2469 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 13), insert(2, None)])
2470 .expect_err("an insert collision must reject the earlier staged update");
2471 assert_eq!(collision.class(), ErrorClass::Conflict);
2472 let failures_unchanged = session
2473 .execute_trusted_dynamic_mutation(&update_code(1, 1))
2474 .expect("failed batches must preserve the original unique value");
2475 assert_eq!(failures_unchanged.affected_rows, 0);
2476
2477 let replaced = session
2478 .execute_trusted_dynamic_mutation_batch(vec![
2479 update_code(1, 14),
2480 DynamicMutation::Replace {
2481 entity: ENTITY_NAME.to_string(),
2482 key: InputValue::Nat64(99),
2483 patch: patch(None, None, Some(99)),
2484 },
2485 ])
2486 .expect("ordinary caller-key replace should insert its absent final row");
2487 assert_eq!(
2488 replaced.rows,
2489 vec![
2490 expected_row_with_code(1, None, 14),
2491 expected_row_with_code(99, None, 99),
2492 ],
2493 );
2494
2495 let unchanged = session
2496 .execute_trusted_dynamic_mutation(&update_code(1, 14))
2497 .expect("the successful mixed replace must publish its preceding update");
2498 assert_eq!(unchanged.affected_rows, 0);
2499 let other_unchanged = session
2500 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
2501 entity: OTHER_ENTITY_NAME.to_string(),
2502 key: InputValue::Nat64(1),
2503 patch: other_patch(None, 10),
2504 })
2505 .expect("cross-entity rejection must preserve the secondary row");
2506 assert_eq!(other_unchanged.affected_rows, 0);
2507 }
2508
2509 #[test]
2510 fn mixed_batch_unique_swap_and_delete_release_use_the_final_overlay() {
2511 let session = initialize();
2512 session
2513 .execute_trusted_dynamic_mutation_batch(vec![
2514 insert_with_code(1, None, 10),
2515 insert_with_code(2, None, 20),
2516 ])
2517 .expect("the unique-overlay fixture should commit");
2518
2519 let swapped = session
2520 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 20), update_code(2, 10)])
2521 .expect("two final rows should atomically swap unique memberships");
2522 assert_eq!(
2523 swapped.rows,
2524 vec![
2525 expected_row_with_code(1, None, 20),
2526 expected_row_with_code(2, None, 10),
2527 ],
2528 );
2529
2530 let released = session
2531 .execute_trusted_dynamic_mutation_batch(vec![delete(1), insert_with_code(3, None, 20)])
2532 .expect("a delete should release unique membership to a final inserted row");
2533 assert_eq!(
2534 released.rows,
2535 vec![
2536 expected_row_with_code(1, None, 20),
2537 expected_row_with_code(3, None, 20),
2538 ],
2539 );
2540 }
2541}
2542
2543#[cfg(test)]
2544mod identity_pre_key_tests {
2545 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2546 use super::DynamicTypedEntityBinding;
2547 use super::{
2548 AcceptedMutationIntentPatch, AcceptedRowLayoutRuntimeContract, AcceptedStructuralMutation,
2549 AcceptedStructuralMutationTarget, DbSession, DynamicMutation, DynamicStructuralPatch,
2550 DynamicTypedFieldBindingRequest, DynamicTypedFieldType, DynamicTypedMutation,
2551 DynamicWriteCell, FieldSlot, MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
2552 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
2553 add_structural_mutation_staged_bytes, checked_pre_key_candidate_count,
2554 insert_key_exists_after_generation, validate_structural_mutation_result_bytes,
2555 };
2556 use crate::{
2557 db::{
2558 commit::{database_incarnation_id, forget_recovered_domain_for_tests},
2559 data::DataStore,
2560 executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
2561 index::IndexStore,
2562 integrity::{
2563 PhysicalUnitCheckpoint, QuickIntegrityStatus, RowInspectionLimits,
2564 execute_quick_integrity, execute_row_integrity_page,
2565 },
2566 journal::JournalTailStore,
2567 registry::{
2568 StoreAllocationIdentities, StoreAllocationIdentity, StoreRegistry,
2569 StoreRuntimeStorageCapabilities,
2570 },
2571 schema::{
2572 AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
2573 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2574 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2575 PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
2576 SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
2577 accepted_schema_candidate_with_field_bindings_for_tests,
2578 },
2579 write_context::MutationMode,
2580 },
2581 error::{ErrorClass, ErrorOrigin, InternalError},
2582 testing::test_memory,
2583 traits::{CanisterKind, Path},
2584 types::{EntityTag, Timestamp},
2585 value::{InputValue, OutputValue, Value},
2586 };
2587 use icydb_schema::{FieldSourceKey, ScalarType};
2588 use std::{cell::RefCell, collections::BTreeMap, time::Instant};
2589
2590 const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
2591 const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
2592 const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
2593 const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
2594 const ENTITY_NAME: &str = "IdentityRow";
2595 const ENTITY_TAG: EntityTag = EntityTag::new(93);
2596 const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
2597
2598 struct TestCanister;
2599
2600 impl Path for TestCanister {
2601 const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
2602 }
2603
2604 impl CanisterKind for TestCanister {
2605 const COMMIT_MEMORY_ID: u8 = 45;
2606 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
2607 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
2608 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2609 "icydb.identity_pre_key_tests.integrity.progress.v1";
2610 }
2611
2612 thread_local! {
2613 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2614 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2615 static SCHEMA_STORE: RefCell<SchemaStore> =
2616 const { RefCell::new(SchemaStore::init_heap()) };
2617 static STORE_REGISTRY: StoreRegistry = {
2618 let mut registry = StoreRegistry::new();
2619 registry.register_store(
2620 STORE_PATH,
2621 &DATA_STORE,
2622 &INDEX_STORE,
2623 &SCHEMA_STORE,
2624 StoreAllocationIdentities::absent(),
2625 StoreRuntimeStorageCapabilities::heap(),
2626 ).expect("identity pre-key test store should register");
2627 registry
2628 };
2629 static JOURNALED_DATA_STORE: RefCell<DataStore> =
2630 RefCell::new(DataStore::init_journaled(test_memory(186)));
2631 static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
2632 RefCell::new(IndexStore::init_journaled(test_memory(187)));
2633 static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
2634 RefCell::new(SchemaStore::init_journaled(test_memory(188)));
2635 static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
2636 RefCell::new(JournalTailStore::init(test_memory(189)));
2637 static JOURNALED_STORE_REGISTRY: StoreRegistry = {
2638 let mut registry = StoreRegistry::new();
2639 registry.register_journaled_store(
2640 JOURNALED_STORE_PATH,
2641 &JOURNALED_DATA_STORE,
2642 &JOURNALED_INDEX_STORE,
2643 &JOURNALED_SCHEMA_STORE,
2644 &JOURNALED_TAIL_STORE,
2645 StoreAllocationIdentities::new_journaled(
2646 StoreAllocationIdentity::new(186, "icydb.test.identity-range.data.v1"),
2647 StoreAllocationIdentity::new(187, "icydb.test.identity-range.index.v1"),
2648 StoreAllocationIdentity::new(188, "icydb.test.identity-range.schema.v1"),
2649 StoreAllocationIdentity::new(189, "icydb.test.identity-range.journal.v1"),
2650 ),
2651 StoreRuntimeStorageCapabilities::journaled(),
2652 ).expect("identity range journaled store should register");
2653 registry
2654 };
2655 }
2656
2657 struct JournaledTestCanister;
2658
2659 impl Path for JournaledTestCanister {
2660 const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
2661 }
2662
2663 impl CanisterKind for JournaledTestCanister {
2664 const COMMIT_MEMORY_ID: u8 = 190;
2665 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
2666 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
2667 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2668 "icydb.identity_range_tests.integrity.progress.v1";
2669 }
2670
2671 fn source_key(source: &str) -> FieldSourceKey {
2672 FieldSourceKey::try_new(source).expect("identity test field source should admit")
2673 }
2674
2675 fn identity_snapshot(store_path: &str) -> PersistedSchemaSnapshot {
2676 let fields = vec![
2677 PersistedFieldSnapshot::new_initial_with_write_policy(
2678 FieldId::new(1),
2679 "id".to_string(),
2680 SchemaFieldSlot::new(0),
2681 AcceptedFieldKind::Nat64,
2682 Vec::new(),
2683 false,
2684 SchemaInsertDefault::None,
2685 SchemaFieldWritePolicy::from_model_policies(
2686 Some(FieldInsertGeneration::Identity),
2687 None,
2688 ),
2689 FieldStorageDecode::ByKind,
2690 LeafCodec::Scalar(ScalarCodec::Nat64),
2691 ),
2692 PersistedFieldSnapshot::new_initial(
2693 FieldId::new(2),
2694 "payload".to_string(),
2695 SchemaFieldSlot::new(1),
2696 AcceptedFieldKind::Nat64,
2697 Vec::new(),
2698 false,
2699 SchemaInsertDefault::None,
2700 FieldStorageDecode::ByKind,
2701 LeafCodec::Scalar(ScalarCodec::Nat64),
2702 ),
2703 ];
2704 PersistedSchemaSnapshot::new_with_indexes(
2705 SchemaVersion::initial(),
2706 ENTITY_SOURCE.to_string(),
2707 ENTITY_NAME.to_string(),
2708 FieldId::new(1),
2709 SchemaRowLayout::initial(
2710 fields
2711 .iter()
2712 .map(|field| (field.id(), field.slot()))
2713 .collect(),
2714 ),
2715 fields,
2716 vec![PersistedIndexSnapshot::new(
2717 SchemaIndexId::new(1).expect("identity test index ID should admit"),
2718 1,
2719 "by_payload".to_string(),
2720 store_path.to_string(),
2721 false,
2722 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2723 FieldId::new(2),
2724 SchemaFieldSlot::new(1),
2725 vec!["payload".to_string()],
2726 AcceptedFieldKind::Nat64,
2727 false,
2728 )]),
2729 None,
2730 )],
2731 )
2732 }
2733
2734 fn initialize() -> DbSession<TestCanister> {
2735 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2736 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2737 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2738 let session = DbSession::<TestCanister>::new(&STORE_REGISTRY);
2739 session
2740 .db
2741 .ensure_recovered_state()
2742 .expect("identity pre-key test database should initialize");
2743 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2744 STORE_PATH,
2745 AcceptedSchemaRevision::INITIAL,
2746 BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH))]),
2747 BTreeMap::from([
2748 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2749 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2750 ]),
2751 );
2752 let store = session
2753 .db
2754 .store_handle(STORE_PATH)
2755 .expect("identity pre-key test store should resolve");
2756 crate::db::commit::publish_accepted_schema_candidate(
2757 STORE_PATH,
2758 store,
2759 AcceptedSchemaRevision::NONE,
2760 &candidate,
2761 )
2762 .expect("identity candidate should publish with explicit zero state");
2763 session
2764 }
2765
2766 fn initialize_journaled() -> DbSession<JournaledTestCanister> {
2767 let session = DbSession::<JournaledTestCanister>::new(&JOURNALED_STORE_REGISTRY);
2768 session
2769 .db
2770 .ensure_recovered_state()
2771 .expect("journaled identity database should initialize");
2772 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2773 JOURNALED_STORE_PATH,
2774 AcceptedSchemaRevision::INITIAL,
2775 BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH))]),
2776 BTreeMap::from([
2777 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2778 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2779 ]),
2780 );
2781 let store = session
2782 .db
2783 .store_handle(JOURNALED_STORE_PATH)
2784 .expect("journaled identity store should resolve");
2785 crate::db::commit::publish_accepted_schema_candidate(
2786 JOURNALED_STORE_PATH,
2787 store,
2788 AcceptedSchemaRevision::NONE,
2789 &candidate,
2790 )
2791 .expect("journaled identity candidate should publish");
2792 session
2793 }
2794
2795 fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
2796 AcceptedMutationIntentPatch::new()
2797 .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
2798 }
2799
2800 fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
2801 DynamicStructuralPatch::new(vec![(
2802 "payload".to_string(),
2803 DynamicWriteCell::Value(InputValue::Nat64(value)),
2804 )])
2805 }
2806
2807 fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
2808 vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
2809 }
2810
2811 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2812 fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
2813 session
2814 .issue_typed_entity_binding(
2815 ENTITY_SOURCE,
2816 &[
2817 DynamicTypedFieldBindingRequest::new(
2818 ID_SOURCE.to_string(),
2819 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2820 false,
2821 ),
2822 DynamicTypedFieldBindingRequest::new(
2823 PAYLOAD_SOURCE.to_string(),
2824 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2825 false,
2826 ),
2827 ],
2828 )
2829 .expect("exact-key test binding should issue")
2830 }
2831
2832 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2833 fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
2834 let output = session
2835 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
2836 entity: ENTITY_NAME.to_string(),
2837 patch: dynamic_payload_patch(payload),
2838 })
2839 .expect("exact-key fixture insert should commit");
2840 match output.rows.as_slice() {
2841 [row] => match row.as_slice() {
2842 [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
2843 if *actual_payload == payload =>
2844 {
2845 *id
2846 }
2847 _ => panic!("exact-key fixture should return its identity and payload"),
2848 },
2849 _ => panic!("exact-key fixture insert should return one row"),
2850 }
2851 }
2852
2853 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2854 fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
2855 let first = insert_exact_key_fixture(session, 41);
2856 let second = insert_exact_key_fixture(session, 42);
2857 let missing = u64::MAX;
2858 let binding = exact_key_binding(session);
2859 let gets_before = DataStore::current_get_call_count();
2860 let result = session
2861 .execute_public_exact_key_batch_for_typed_binding(
2862 &binding,
2863 &[second, missing, first, second],
2864 )
2865 .expect("exact-key batch should execute")
2866 .expect("exact-key binding should remain current");
2867
2868 assert_eq!(result.positions, vec![0, 1, 2, 0]);
2869 assert_eq!(
2870 result.distinct_rows,
2871 vec![
2872 Some(expected_dynamic_row(second, 42)),
2873 None,
2874 Some(expected_dynamic_row(first, 41)),
2875 ],
2876 );
2877 assert_eq!(
2878 DataStore::current_get_call_count().saturating_sub(gets_before),
2879 3,
2880 "four input positions with one duplicate must perform three physical reads",
2881 );
2882 }
2883
2884 #[cfg(all(feature = "sql", feature = "diagnostics"))]
2885 #[test]
2886 fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
2887 assert_exact_key_batch(&initialize());
2888 assert_exact_key_batch(&initialize_journaled());
2889 }
2890
2891 fn assert_dynamic_payload(session: &DbSession<TestCanister>, key: u64, expected_payload: u64) {
2892 let unchanged = session
2893 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
2894 entity: ENTITY_NAME.to_string(),
2895 key: InputValue::Nat64(key),
2896 patch: dynamic_payload_patch(expected_payload),
2897 })
2898 .expect("the expected row should remain readable through a no-op update");
2899 assert_eq!(unchanged.affected_rows, 0);
2900 assert_eq!(
2901 unchanged.rows,
2902 vec![expected_dynamic_row(key, expected_payload)],
2903 );
2904 }
2905
2906 fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
2907 values
2908 .iter()
2909 .map(|value| {
2910 AcceptedStructuralMutation::save(
2911 MutationMode::Insert,
2912 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
2913 payload_patch(*value),
2914 )
2915 })
2916 .collect()
2917 }
2918
2919 fn assert_identity_boundary(error: &InternalError) {
2920 assert_eq!(error.class(), ErrorClass::Unsupported);
2921 assert_eq!(error.origin(), ErrorOrigin::Identity);
2922 }
2923
2924 #[test]
2925 fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
2926 let generated = insert_key_exists_after_generation(true);
2927 assert_eq!(generated.class(), ErrorClass::Corruption);
2928 assert_eq!(generated.origin(), ErrorOrigin::Identity);
2929
2930 let ordinary = insert_key_exists_after_generation(false);
2931 assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
2932 }
2933
2934 #[cfg(target_pointer_width = "64")]
2935 #[test]
2936 fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
2937 let error = checked_pre_key_candidate_count(
2938 usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
2939 )
2940 .expect_err("candidate counts beyond u32 must reject");
2941 assert_identity_boundary(&error);
2942 }
2943
2944 #[test]
2945 #[expect(
2946 clippy::too_many_lines,
2947 reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
2948 )]
2949 fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
2950 let session = initialize();
2951 let seeded = session
2952 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
2953 .expect("seed rows should commit");
2954 assert_eq!(seeded.affected_rows, 1);
2955
2956 let split = session
2957 .execute_trusted_dynamic_mutation_batch(vec![
2958 DynamicMutation::Update {
2959 entity: ENTITY_NAME.to_string(),
2960 key: InputValue::Nat64(1),
2961 patch: dynamic_payload_patch(60),
2962 },
2963 DynamicMutation::Insert {
2964 entity: ENTITY_NAME.to_string(),
2965 patch: dynamic_payload_patch(40),
2966 },
2967 ])
2968 .expect("one holding should split atomically");
2969 assert_eq!(split.affected_rows, 2);
2970 assert_eq!(
2971 split.rows,
2972 vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
2973 "split after-images must retain input order and exact quantity",
2974 );
2975
2976 let rejected_split = session
2977 .execute_trusted_dynamic_mutation_batch(vec![
2978 DynamicMutation::Update {
2979 entity: ENTITY_NAME.to_string(),
2980 key: InputValue::Nat64(1),
2981 patch: dynamic_payload_patch(50),
2982 },
2983 DynamicMutation::Insert {
2984 entity: ENTITY_NAME.to_string(),
2985 patch: DynamicStructuralPatch::new(Vec::new()),
2986 },
2987 ])
2988 .expect_err("an invalid split output must reject the staged source update");
2989 assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
2990 assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
2991 assert_eq!(
2992 rejected_split.diagnostic_facts(),
2993 vec![
2994 (
2995 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
2996 ENTITY_TAG.value(),
2997 ),
2998 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
2999 (
3000 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3001 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3002 ),
3003 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3004 ],
3005 );
3006 assert_dynamic_payload(&session, 1, 60);
3007 assert_dynamic_payload(&session, 2, 40);
3008
3009 let transfer = session
3010 .execute_trusted_dynamic_mutation_batch(vec![
3011 DynamicMutation::Update {
3012 entity: ENTITY_NAME.to_string(),
3013 key: InputValue::Nat64(1),
3014 patch: dynamic_payload_patch(70),
3015 },
3016 DynamicMutation::Update {
3017 entity: ENTITY_NAME.to_string(),
3018 key: InputValue::Nat64(2),
3019 patch: dynamic_payload_patch(30),
3020 },
3021 ])
3022 .expect("distinct transfer patches should share one atomic batch");
3023 assert_eq!(
3024 transfer.rows,
3025 vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
3026 "the transfer must preserve the exact total quantity",
3027 );
3028
3029 let merge = session
3030 .execute_trusted_dynamic_mutation_batch(vec![
3031 DynamicMutation::Delete {
3032 entity: ENTITY_NAME.to_string(),
3033 key: InputValue::Nat64(2),
3034 },
3035 DynamicMutation::Update {
3036 entity: ENTITY_NAME.to_string(),
3037 key: InputValue::Nat64(1),
3038 patch: dynamic_payload_patch(100),
3039 },
3040 ])
3041 .expect("two holdings should merge atomically");
3042 assert_eq!(
3043 merge.rows,
3044 vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
3045 "delete before-images and update after-images must retain input order",
3046 );
3047
3048 let resplit = session
3049 .execute_trusted_dynamic_mutation_batch(vec![
3050 DynamicMutation::Update {
3051 entity: ENTITY_NAME.to_string(),
3052 key: InputValue::Nat64(1),
3053 patch: dynamic_payload_patch(60),
3054 },
3055 DynamicMutation::Insert {
3056 entity: ENTITY_NAME.to_string(),
3057 patch: dynamic_payload_patch(40),
3058 },
3059 ])
3060 .expect("the merged holding should split again");
3061 assert_eq!(
3062 resplit.rows,
3063 vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
3064 );
3065
3066 let rejected_merge = session
3067 .execute_trusted_dynamic_mutation_batch(vec![
3068 DynamicMutation::Delete {
3069 entity: ENTITY_NAME.to_string(),
3070 key: InputValue::Nat64(3),
3071 },
3072 DynamicMutation::Update {
3073 entity: ENTITY_NAME.to_string(),
3074 key: InputValue::Nat64(99),
3075 patch: dynamic_payload_patch(100),
3076 },
3077 ])
3078 .expect_err("a late missing merge target must preserve the earlier staged delete");
3079 assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
3080 assert_dynamic_payload(&session, 1, 60);
3081 assert_dynamic_payload(&session, 3, 40);
3082
3083 SCHEMA_STORE.with(|store| {
3084 let cursor = store
3085 .borrow()
3086 .identity_statement_cursor(
3087 database_incarnation_id().expect("database incarnation should remain readable"),
3088 ENTITY_TAG,
3089 FieldId::new(1),
3090 &AcceptedFieldKind::Nat64,
3091 )
3092 .expect("mixed Identity state should remain readable");
3093 assert_eq!(cursor.expected_high_water(), 3);
3094 assert!(!cursor.has_allocations());
3095 });
3096 }
3097
3098 #[test]
3099 fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
3100 let session = initialize();
3101 session
3102 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3103 .expect("the holding fixture should initialize");
3104
3105 let duplicate = session
3106 .execute_trusted_dynamic_mutation_batch(vec![
3107 DynamicMutation::Update {
3108 entity: ENTITY_NAME.to_string(),
3109 key: InputValue::Nat64(1),
3110 patch: dynamic_payload_patch(60),
3111 },
3112 DynamicMutation::Delete {
3113 entity: ENTITY_NAME.to_string(),
3114 key: InputValue::Nat64(1),
3115 },
3116 ])
3117 .expect_err("duplicate targets across operation kinds must reject");
3118 assert!(matches!(
3119 duplicate.diagnostic().detail(),
3120 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3121 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
3122 }),
3123 ));
3124 assert_eq!(
3125 duplicate.diagnostic_facts(),
3126 vec![
3127 (
3128 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3129 ENTITY_TAG.value(),
3130 ),
3131 (
3132 icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
3133 0,
3134 ),
3135 (
3136 icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
3137 1,
3138 ),
3139 ],
3140 );
3141 assert_dynamic_payload(&session, 1, 100);
3142 }
3143
3144 #[test]
3145 fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
3146 let session = initialize();
3147 let empty = session
3148 .execute_trusted_dynamic_mutation_batch(Vec::new())
3149 .expect_err("an empty public batch must reject");
3150 assert!(matches!(
3151 empty.diagnostic().detail(),
3152 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3153 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
3154 }),
3155 ));
3156 assert_eq!(
3157 empty.diagnostic_facts(),
3158 vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
3159 );
3160
3161 let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
3162 .map(|_| DynamicMutation::Delete {
3163 entity: ENTITY_NAME.to_string(),
3164 key: InputValue::Nat64(1),
3165 })
3166 .collect();
3167 let over_bound = session
3168 .execute_trusted_dynamic_mutation_batch(requests)
3169 .expect_err("operation cap plus one must reject before row resolution");
3170 assert!(matches!(
3171 over_bound.diagnostic().detail(),
3172 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3173 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
3174 }),
3175 ));
3176 assert_eq!(
3177 over_bound.diagnostic_facts(),
3178 vec![
3179 (
3180 icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
3181 (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
3182 ),
3183 (
3184 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3185 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
3186 ),
3187 ],
3188 );
3189 }
3190
3191 #[test]
3192 fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
3193 let mut exact = 0;
3194 add_structural_mutation_staged_bytes(
3195 &mut exact,
3196 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
3197 )
3198 .expect("the exact staged-byte boundary should admit");
3199 assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
3200
3201 let error = add_structural_mutation_staged_bytes(&mut exact, [1])
3202 .expect_err("one byte above the staged-byte boundary must reject");
3203 assert!(matches!(
3204 error.diagnostic().detail(),
3205 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3206 boundary:
3207 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
3208 }),
3209 ));
3210 assert_eq!(
3211 error.diagnostic_facts(),
3212 vec![
3213 (
3214 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3215 (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
3216 ),
3217 (
3218 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3219 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
3220 ),
3221 ],
3222 );
3223
3224 validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
3225 .expect("the exact result-byte boundary should admit");
3226 let error = validate_structural_mutation_result_bytes(
3227 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
3228 )
3229 .expect_err("one byte above the result-byte boundary must reject");
3230 assert!(matches!(
3231 error.diagnostic().detail(),
3232 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3233 boundary:
3234 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
3235 }),
3236 ));
3237 assert_eq!(
3238 error.diagnostic_facts(),
3239 vec![
3240 (
3241 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3242 (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
3243 ),
3244 (
3245 icydb_diagnostic_code::DiagnosticFactTag::Limit,
3246 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
3247 ),
3248 ],
3249 );
3250 }
3251
3252 #[expect(
3253 clippy::too_many_lines,
3254 reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
3255 )]
3256 #[test]
3257 fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
3258 let session = initialize();
3259 let catalog = session
3260 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3261 .expect("identity catalog should resolve");
3262 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3263 .expect("identity row layout should build");
3264 let initial_description = session
3265 .try_describe_entity_by_name(ENTITY_NAME)
3266 .expect("accepted Identity description should resolve");
3267 assert_eq!(
3268 initial_description.entity_tag(),
3269 catalog.identity().entity_tag().value()
3270 );
3271 assert_eq!(
3272 initial_description.accepted_schema_fingerprint_method(),
3273 catalog.fingerprint_method_version()
3274 );
3275 assert_eq!(
3276 initial_description.accepted_schema_fingerprint(),
3277 catalog.fingerprint()
3278 );
3279 let initial_identity = initial_description
3280 .identity()
3281 .expect("accepted Identity policy should be described");
3282 assert_eq!(initial_identity.field(), "id");
3283 assert_eq!(initial_identity.generator(), "Identity::next");
3284 assert_eq!(initial_identity.accepted_kind(), "nat64");
3285 assert_eq!(initial_identity.minimum(), 1);
3286 assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
3287 assert_eq!(initial_identity.high_water(), 0);
3288 assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
3289 assert!(!initial_identity.exhausted());
3290
3291 let rejected = session
3292 .execute_accepted_structural_save_batch(
3293 &catalog,
3294 &descriptor,
3295 batch(&[1_000, 2_000]),
3296 Timestamp::from_millis(6),
3297 |_| Err::<(), _>(InternalError::executor_unsupported()),
3298 )
3299 .expect_err("a rejected precommit result must not publish its tentative range");
3300 assert_eq!(rejected.class(), ErrorClass::Unsupported);
3301 assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
3302
3303 let rows = session
3304 .execute_accepted_structural_save_batch(
3305 &catalog,
3306 &descriptor,
3307 batch(&[10, 20, 30]),
3308 Timestamp::from_millis(7),
3309 Ok,
3310 )
3311 .expect("one accepted batch should commit rows and one identity range");
3312 assert_eq!(
3313 rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
3314 vec![
3315 vec![Value::Nat64(1), Value::Nat64(10)],
3316 vec![Value::Nat64(2), Value::Nat64(20)],
3317 vec![Value::Nat64(3), Value::Nat64(30)],
3318 ],
3319 );
3320
3321 let dynamic = session
3322 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3323 entity: ENTITY_NAME.to_string(),
3324 patch: DynamicStructuralPatch::new(vec![(
3325 "payload".to_string(),
3326 DynamicWriteCell::Value(InputValue::Nat64(40)),
3327 )]),
3328 })
3329 .expect("dynamic omission should commit through shared Identity generation");
3330 assert_eq!(dynamic.affected_rows, 1);
3331
3332 for (request, operation) in [
3333 (
3334 DynamicMutation::Insert {
3335 entity: ENTITY_NAME.to_string(),
3336 patch: DynamicStructuralPatch::new(vec![
3337 (
3338 "id".to_string(),
3339 DynamicWriteCell::Value(InputValue::Nat64(41)),
3340 ),
3341 (
3342 "payload".to_string(),
3343 DynamicWriteCell::Value(InputValue::Nat64(42)),
3344 ),
3345 ]),
3346 },
3347 icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
3348 ),
3349 (
3350 DynamicMutation::Update {
3351 entity: ENTITY_NAME.to_string(),
3352 key: InputValue::Nat64(1),
3353 patch: DynamicStructuralPatch::new(vec![(
3354 "id".to_string(),
3355 DynamicWriteCell::Default,
3356 )]),
3357 },
3358 icydb_diagnostic_code::DiagnosticMutationOperation::Update,
3359 ),
3360 ] {
3361 let error = session
3362 .execute_trusted_dynamic_mutation(&request)
3363 .expect_err("structural Identity authorship and regeneration must reject");
3364 assert_eq!(error.class(), ErrorClass::Unsupported);
3365 assert_eq!(error.origin(), ErrorOrigin::Executor);
3366 assert_eq!(
3367 error.diagnostic_facts(),
3368 vec![
3369 (
3370 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3371 ENTITY_TAG.value(),
3372 ),
3373 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3374 (
3375 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3376 operation.raw(),
3377 ),
3378 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3379 ],
3380 );
3381 }
3382
3383 let binding = session
3384 .issue_typed_entity_binding(
3385 ENTITY_SOURCE,
3386 &[
3387 DynamicTypedFieldBindingRequest::new(
3388 ID_SOURCE.to_string(),
3389 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3390 false,
3391 ),
3392 DynamicTypedFieldBindingRequest::new(
3393 PAYLOAD_SOURCE.to_string(),
3394 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3395 false,
3396 ),
3397 ],
3398 )
3399 .expect("typed output should bind the Identity field");
3400 let typed_patch = binding
3401 .bind_write_fields(vec![(
3402 PAYLOAD_SOURCE.to_string(),
3403 DynamicWriteCell::Value(InputValue::Nat64(50)),
3404 )])
3405 .expect("typed payload should lower");
3406 let typed = session
3407 .execute_trusted_typed_mutation(
3408 &binding,
3409 &DynamicTypedMutation::Insert { patch: typed_patch },
3410 )
3411 .expect("typed omission should commit through shared Identity generation");
3412 assert_eq!(
3413 typed
3414 .expect("typed insert should return one mutation result")
3415 .affected_rows,
3416 1,
3417 );
3418 let explicit_typed_patch = binding
3419 .bind_write_fields(vec![
3420 (
3421 ID_SOURCE.to_string(),
3422 DynamicWriteCell::Value(InputValue::Nat64(51)),
3423 ),
3424 (
3425 PAYLOAD_SOURCE.to_string(),
3426 DynamicWriteCell::Value(InputValue::Nat64(52)),
3427 ),
3428 ])
3429 .expect("the low-level binding should retain exact authored intent");
3430 let explicit_typed_error = session
3431 .execute_trusted_typed_mutation(
3432 &binding,
3433 &DynamicTypedMutation::Insert {
3434 patch: explicit_typed_patch,
3435 },
3436 )
3437 .expect_err("typed Identity authorship must reject before allocation");
3438 assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
3439 assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
3440 assert_eq!(
3441 explicit_typed_error.diagnostic_facts(),
3442 vec![
3443 (
3444 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3445 ENTITY_TAG.value(),
3446 ),
3447 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3448 (
3449 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3450 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3451 ),
3452 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3453 ],
3454 );
3455
3456 let replace_error = session
3457 .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
3458 entity: ENTITY_NAME.to_string(),
3459 key: InputValue::Nat64(99),
3460 patch: DynamicStructuralPatch::new(vec![(
3461 "payload".to_string(),
3462 DynamicWriteCell::Value(InputValue::Nat64(60)),
3463 )]),
3464 })
3465 .expect_err("save-as-insert with a chosen Identity must reject");
3466 assert_eq!(replace_error.class(), ErrorClass::Unsupported);
3467 assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
3468
3469 #[cfg(feature = "sql")]
3470 {
3471 for sql in [
3472 "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
3473 "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
3474 ] {
3475 let _result = session
3476 .execute_trusted_sql_mutation(sql)
3477 .expect("SQL omission and DEFAULT should commit Identity generation");
3478 }
3479
3480 let error = session
3481 .execute_trusted_sql_mutation(
3482 "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
3483 )
3484 .expect_err("an explicit SQL Identity value must reject before allocation");
3485 let diagnostic = error.diagnostic();
3486 assert_eq!(
3487 diagnostic.code(),
3488 icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
3489 );
3490 assert!(matches!(
3491 diagnostic.detail(),
3492 Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
3493 boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
3494 }),
3495 ));
3496 }
3497
3498 let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
3499 assert_eq!(
3500 DATA_STORE.with(|store| store.borrow().len()),
3501 expected_committed
3502 );
3503 SCHEMA_STORE.with(|store| {
3504 let cursor = store
3505 .borrow()
3506 .identity_statement_cursor(
3507 database_incarnation_id().expect("database incarnation should remain readable"),
3508 ENTITY_TAG,
3509 FieldId::new(1),
3510 &AcceptedFieldKind::Nat64,
3511 )
3512 .expect("committed writes must leave active state readable");
3513 assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
3514 assert!(!cursor.has_allocations());
3515 });
3516 let committed_description = session
3517 .try_describe_entity_by_name(ENTITY_NAME)
3518 .expect("committed Identity description should resolve");
3519 let committed_identity = committed_description
3520 .identity()
3521 .expect("accepted Identity policy should remain described");
3522 assert_eq!(
3523 committed_identity.high_water(),
3524 u128::from(expected_committed),
3525 );
3526 assert_eq!(
3527 committed_identity.remaining(),
3528 u128::from(u64::MAX - expected_committed),
3529 );
3530 assert!(!committed_identity.exhausted());
3531 }
3532
3533 #[test]
3534 #[expect(
3535 clippy::too_many_lines,
3536 reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
3537 )]
3538 fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
3539 let session = initialize_journaled();
3540 let catalog = session
3541 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3542 .expect("journaled identity catalog should resolve");
3543 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3544 .expect("journaled identity row layout should build");
3545
3546 for (ordinal, interruption) in [
3547 MutationCommitInterruption::MarkerPersisted,
3548 MutationCommitInterruption::JournalPublished,
3549 MutationCommitInterruption::RowsPublished,
3550 MutationCommitInterruption::StateMaterialized,
3551 ]
3552 .into_iter()
3553 .enumerate()
3554 {
3555 interrupt_next_mutation_commit_for_tests(interruption);
3556 let interrupted = session.execute_accepted_structural_save_batch(
3557 &catalog,
3558 &descriptor,
3559 batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
3560 Timestamp::from_millis(8),
3561 Ok,
3562 );
3563 assert!(
3564 interrupted.is_err(),
3565 "the selected durable boundary should interrupt",
3566 );
3567
3568 let committed = session
3569 .execute_accepted_structural_save_batch(
3570 &catalog,
3571 &descriptor,
3572 batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
3573 Timestamp::from_millis(9),
3574 Ok,
3575 )
3576 .expect("the next mutation must recover before allocating");
3577 let expected_high_water =
3578 u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
3579 assert_eq!(
3580 committed
3581 .into_iter()
3582 .map(|row| row.values)
3583 .collect::<Vec<_>>(),
3584 vec![vec![
3585 Value::Nat64(expected_high_water),
3586 Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
3587 ]],
3588 );
3589 assert_eq!(
3590 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
3591 expected_high_water,
3592 );
3593 JOURNALED_SCHEMA_STORE.with(|store| {
3594 let cursor = store
3595 .borrow()
3596 .identity_statement_cursor(
3597 database_incarnation_id()
3598 .expect("database incarnation should remain readable"),
3599 ENTITY_TAG,
3600 FieldId::new(1),
3601 &AcceptedFieldKind::Nat64,
3602 )
3603 .expect("guarded recovery must leave quiescent active state");
3604 assert_eq!(
3605 cursor.expected_high_water(),
3606 u128::from(expected_high_water),
3607 );
3608 assert!(!cursor.has_allocations());
3609 });
3610 }
3611
3612 for (ordinal, (interruption, deleted_key)) in [
3613 (MutationCommitInterruption::MarkerPersisted, 2),
3614 (MutationCommitInterruption::JournalPublished, 4),
3615 (MutationCommitInterruption::RowPrefixPublished, 6),
3616 (MutationCommitInterruption::RowsPublished, 8),
3617 (MutationCommitInterruption::StateMaterialized, 7),
3618 ]
3619 .into_iter()
3620 .enumerate()
3621 {
3622 let expected_payload =
3623 501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
3624 interrupt_next_mutation_commit_for_tests(interruption);
3625 let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
3626 DynamicMutation::Update {
3627 entity: ENTITY_NAME.to_string(),
3628 key: InputValue::Nat64(1),
3629 patch: dynamic_payload_patch(expected_payload),
3630 },
3631 DynamicMutation::Delete {
3632 entity: ENTITY_NAME.to_string(),
3633 key: InputValue::Nat64(deleted_key),
3634 },
3635 ]);
3636 assert!(
3637 interrupted.is_err(),
3638 "the selected caller-key mixed publication boundary should interrupt",
3639 );
3640 let recovered_update = session
3641 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3642 entity: ENTITY_NAME.to_string(),
3643 key: InputValue::Nat64(1),
3644 patch: dynamic_payload_patch(expected_payload),
3645 })
3646 .expect("guarded reentry should complete the marker-authorized mixed batch");
3647 assert_eq!(
3648 recovered_update.affected_rows, 0,
3649 "the recovered update must already expose its admitted final image",
3650 );
3651 let recovered_delete = session
3652 .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
3653 entity: ENTITY_NAME.to_string(),
3654 key: InputValue::Nat64(deleted_key),
3655 })
3656 .expect_err("the recovered delete must already be materialized");
3657 assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
3658 JOURNALED_SCHEMA_STORE.with(|store| {
3659 let cursor = store
3660 .borrow()
3661 .identity_statement_cursor(
3662 database_incarnation_id()
3663 .expect("database incarnation should remain readable"),
3664 ENTITY_TAG,
3665 FieldId::new(1),
3666 &AcceptedFieldKind::Nat64,
3667 )
3668 .expect("caller-key recovery must preserve active Identity state");
3669 assert_eq!(cursor.expected_high_water(), 8);
3670 assert!(!cursor.has_allocations());
3671 });
3672 }
3673
3674 forget_recovered_domain_for_tests(&session.db)
3675 .expect("the final journal tail should remain recoverable");
3676 session
3677 .db
3678 .ensure_recovered_state()
3679 .expect("derived rebuild must not allocate another identity");
3680
3681 let quick = execute_quick_integrity(&session.db, catalog.inspection_plan())
3682 .expect("quiescent Identity control inventory should be inspectable");
3683 assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
3684 let row_page = execute_row_integrity_page(
3685 &session.db,
3686 catalog.inspection_plan(),
3687 PhysicalUnitCheckpoint::BeforeFirst,
3688 RowInspectionLimits::standard(),
3689 )
3690 .expect("Identity rows should remain within committed high-water");
3691 assert!(row_page.exhausted());
3692 assert!(row_page.findings().is_empty());
3693
3694 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
3695 assert!(
3696 JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
3697 "derived index rebuild should restore witnesses without allocating identities",
3698 );
3699 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
3700 JOURNALED_SCHEMA_STORE.with(|store| {
3701 let cursor = store
3702 .borrow()
3703 .identity_statement_cursor(
3704 database_incarnation_id().expect("database incarnation should remain readable"),
3705 ENTITY_TAG,
3706 FieldId::new(1),
3707 &AcceptedFieldKind::Nat64,
3708 )
3709 .expect("folded identity state should reopen without allocating");
3710 assert_eq!(cursor.expected_high_water(), 8);
3711 assert!(!cursor.has_allocations());
3712 });
3713 }
3714
3715 #[test]
3716 #[ignore = "release-closeout native timing probe for one marker-authorized Identity recovery"]
3717 fn identity_recovery_closeout_reports_guarded_reentry_time() {
3718 let session = initialize_journaled();
3719 let catalog = session
3720 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3721 .expect("journaled identity catalog should resolve");
3722 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3723 .expect("journaled identity row layout should build");
3724
3725 interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
3726 let interrupted = session.execute_accepted_structural_save_batch(
3727 &catalog,
3728 &descriptor,
3729 batch(&[1]),
3730 Timestamp::from_millis(10),
3731 Ok,
3732 );
3733 assert!(
3734 interrupted.is_err(),
3735 "the selected publication boundary should interrupt",
3736 );
3737
3738 let start = Instant::now();
3739 let committed = session
3740 .execute_accepted_structural_save_batch(
3741 &catalog,
3742 &descriptor,
3743 batch(&[2]),
3744 Timestamp::from_millis(11),
3745 Ok,
3746 )
3747 .expect("guarded reentry should recover before allocation");
3748 let elapsed = start.elapsed();
3749 assert_eq!(
3750 committed
3751 .into_iter()
3752 .map(|row| row.values)
3753 .collect::<Vec<_>>(),
3754 vec![vec![Value::Nat64(2), Value::Nat64(2)]],
3755 );
3756
3757 println!(
3758 "identity recovery closeout: guarded_reentry_nanos={}",
3759 elapsed.as_nanos(),
3760 );
3761 }
3762}
3763
3764#[cfg(test)]
3765mod targeted_rule_mutation_tests {
3766 use super::{
3767 DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
3768 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
3769 };
3770 use crate::{
3771 db::{
3772 data::{DataStore, encode_input_value_for_candidate_field_contract},
3773 index::IndexStore,
3774 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
3775 schema::{
3776 AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
3777 AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
3778 AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
3779 ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
3780 PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
3781 ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
3782 SchemaRowLayout, SchemaStore, SchemaVersion,
3783 accepted_schema_candidate_with_catalogs_for_tests,
3784 build_record_newtype_composite_catalog_for_tests,
3785 empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
3786 },
3787 },
3788 error::InternalError,
3789 traits::{CanisterKind, Path},
3790 types::EntityTag,
3791 value::InputValue,
3792 };
3793 use icydb_schema::{
3794 ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
3795 };
3796 use std::{cell::RefCell, collections::BTreeMap};
3797
3798 const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
3799 const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
3800 const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
3801 const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
3802 const UPDATED_AT_SOURCE: &str =
3803 "session::write::targeted_rule_mutation_tests::Entity::updated_at";
3804 const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
3805 const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
3806 const DEGREE_MEMBER_SOURCE: &str =
3807 "session::write::targeted_rule_mutation_tests::Profile::degree";
3808 const DEGREE_RULE_SOURCE: &str =
3809 "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
3810
3811 struct TestCanister;
3812
3813 impl Path for TestCanister {
3814 const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
3815 }
3816
3817 impl CanisterKind for TestCanister {
3818 const COMMIT_MEMORY_ID: u8 = 43;
3819 const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
3820 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
3821 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3822 "icydb.targeted_mutation_tests.integrity.progress.v1";
3823 }
3824
3825 thread_local! {
3826 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
3827 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
3828 static SCHEMA_STORE: RefCell<SchemaStore> =
3829 const { RefCell::new(SchemaStore::init_heap()) };
3830 static STORE_REGISTRY: StoreRegistry = {
3831 let mut registry = StoreRegistry::new();
3832 registry.register_store(
3833 STORE_PATH,
3834 &DATA_STORE,
3835 &INDEX_STORE,
3836 &SCHEMA_STORE,
3837 StoreAllocationIdentities::absent(),
3838 StoreRuntimeStorageCapabilities::heap(),
3839 ).expect("targeted mutation test store should register");
3840 registry
3841 };
3842 }
3843
3844 fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
3845 parse(raw.to_string()).expect("test source identity should admit")
3846 }
3847
3848 fn profile_input(degree: u64) -> InputValue {
3849 InputValue::Map(vec![(
3850 InputValue::Text("degree".to_string()),
3851 InputValue::Nat64(degree),
3852 )])
3853 }
3854
3855 fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
3856 DynamicStructuralPatch::new(vec![
3857 (
3858 "id".to_string(),
3859 DynamicWriteCell::Value(InputValue::Nat64(id)),
3860 ),
3861 (
3862 "profile".to_string(),
3863 DynamicWriteCell::Value(profile_input(degree)),
3864 ),
3865 ])
3866 }
3867
3868 fn encoded_value(
3869 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
3870 composite_catalog: &AcceptedCompositeCatalog,
3871 name: &str,
3872 kind: &AcceptedFieldKind,
3873 storage_decode: FieldStorageDecode,
3874 leaf_codec: LeafCodec,
3875 value: InputValue,
3876 ) -> Vec<u8> {
3877 let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
3878 encode_input_value_for_candidate_field_contract(
3879 enum_catalog,
3880 composite_catalog,
3881 field,
3882 value,
3883 &mut ValueAdmissionBudget::standard(),
3884 )
3885 .expect("test accepted value should encode")
3886 }
3887
3888 fn nat64_literal(
3889 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
3890 composite_catalog: &AcceptedCompositeCatalog,
3891 value: u64,
3892 ) -> AcceptedCheckLiteralV1 {
3893 let kind = AcceptedFieldKind::Nat64;
3894 AcceptedCheckLiteralV1::from_accepted_parts(
3895 kind.clone(),
3896 FieldStorageDecode::ByKind,
3897 LeafCodec::Scalar(ScalarCodec::Nat64),
3898 encoded_value(
3899 enum_catalog,
3900 composite_catalog,
3901 "degree_bound",
3902 &kind,
3903 FieldStorageDecode::ByKind,
3904 LeafCodec::Scalar(ScalarCodec::Nat64),
3905 InputValue::Nat64(value),
3906 ),
3907 )
3908 }
3909
3910 fn targeted_constraint_id(error: &InternalError) -> u32 {
3911 let facts = error.diagnostic_facts();
3912 assert!(facts.contains(&(
3913 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3914 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3915 )));
3916 assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
3917 assert!(facts.contains(&(
3918 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
3919 icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
3920 )));
3921 assert_eq!(
3922 facts
3923 .iter()
3924 .filter(|(tag, _)| matches!(
3925 tag,
3926 icydb_diagnostic_code::DiagnosticFactTag::RootField
3927 | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
3928 ))
3929 .copied()
3930 .collect::<Vec<_>>(),
3931 vec![
3932 (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
3933 (
3934 icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
3935 icydb_diagnostic_code::pack_u32_pair(1, 1),
3936 ),
3937 ]
3938 );
3939 let value = facts
3940 .iter()
3941 .find_map(|(tag, value)| {
3942 (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
3943 })
3944 .expect("targeted mutation should retain its accepted constraint ID");
3945 u32::try_from(value).expect("accepted constraint ID fits u32")
3946 }
3947
3948 #[expect(
3949 clippy::too_many_lines,
3950 reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
3951 )]
3952 #[test]
3953 fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
3954 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
3955 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
3956 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
3957
3958 let entity_tag = EntityTag::new(93);
3959 let enum_catalog = empty_accepted_enum_catalog_for_tests();
3960 let (composite_catalog, profile_type, degree_type, degree_member) =
3961 build_record_newtype_composite_catalog_for_tests(
3962 "tests::TargetedProfile".to_string(),
3963 "degree".to_string(),
3964 "tests::TargetedDegree".to_string(),
3965 AcceptedFieldKind::Nat64,
3966 &enum_catalog,
3967 )
3968 .expect("targeted mutation composites should close");
3969 let profile_kind = AcceptedFieldKind::Composite {
3970 type_id: profile_type,
3971 };
3972 let profile_default = encoded_value(
3973 &enum_catalog,
3974 &composite_catalog,
3975 "profile",
3976 &profile_kind,
3977 FieldStorageDecode::CatalogValue,
3978 LeafCodec::Structural,
3979 profile_input(12),
3980 );
3981 let fields = vec![
3982 PersistedFieldSnapshot::new_initial(
3983 FieldId::new(1),
3984 "id".to_string(),
3985 SchemaFieldSlot::new(0),
3986 AcceptedFieldKind::Nat64,
3987 Vec::new(),
3988 false,
3989 SchemaInsertDefault::None,
3990 FieldStorageDecode::ByKind,
3991 LeafCodec::Scalar(ScalarCodec::Nat64),
3992 ),
3993 PersistedFieldSnapshot::new_initial(
3994 FieldId::new(2),
3995 "profile".to_string(),
3996 SchemaFieldSlot::new(1),
3997 profile_kind,
3998 vec![PersistedNestedLeafSnapshot::new(
3999 vec!["degree".to_string()],
4000 AcceptedFieldKind::Composite {
4001 type_id: degree_type,
4002 },
4003 false,
4004 )],
4005 false,
4006 SchemaInsertDefault::SlotPayload(profile_default),
4007 FieldStorageDecode::CatalogValue,
4008 LeafCodec::Structural,
4009 ),
4010 PersistedFieldSnapshot::new_initial_with_write_policy(
4011 FieldId::new(3),
4012 "updated_at".to_string(),
4013 SchemaFieldSlot::new(2),
4014 AcceptedFieldKind::Timestamp,
4015 Vec::new(),
4016 false,
4017 SchemaInsertDefault::None,
4018 SchemaFieldWritePolicy::from_model_policies(
4019 None,
4020 Some(FieldWriteManagement::UpdatedAt),
4021 ),
4022 FieldStorageDecode::ByKind,
4023 LeafCodec::Scalar(ScalarCodec::Timestamp),
4024 ),
4025 ];
4026 let mut snapshot = PersistedSchemaSnapshot::new(
4027 SchemaVersion::initial(),
4028 ENTITY_SOURCE.to_string(),
4029 "TargetedMutation".to_string(),
4030 FieldId::new(1),
4031 SchemaRowLayout::initial(
4032 fields
4033 .iter()
4034 .map(|field| (field.id(), field.slot()))
4035 .collect(),
4036 ),
4037 fields,
4038 );
4039 let constraint_catalog = snapshot
4040 .constraint_catalog()
4041 .clone()
4042 .with_added_targeted_rule(
4043 "profile_degree_multiple".to_string(),
4044 ConstraintOrigin::Generated,
4045 AcceptedRuleTarget::new(
4046 FieldId::new(2),
4047 AcceptedNamedTypeIdentity::Composite(degree_type),
4048 ),
4049 AcceptedRuleOperation::MultipleOf {
4050 divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
4051 },
4052 )
4053 .expect("targeted mutation rule should allocate");
4054 let targeted_rule_id = constraint_catalog
4055 .constraints()
4056 .last()
4057 .expect("targeted mutation rule should persist")
4058 .id();
4059 snapshot = snapshot.with_constraint_catalog(constraint_catalog);
4060
4061 let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
4062 let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
4063 let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
4064 let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
4065 let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
4066 let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
4067 let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
4068 let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
4069 let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
4070 BTreeMap::from([(entity_source, entity_tag)]),
4071 BTreeMap::from([
4072 ((entity_tag, id_source), FieldId::new(1)),
4073 ((entity_tag, profile_source), FieldId::new(2)),
4074 ((entity_tag, updated_at_source), FieldId::new(3)),
4075 ]),
4076 BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
4077 BTreeMap::new(),
4078 BTreeMap::new(),
4079 )
4080 .with_initial_named_types_for_tests(
4081 BTreeMap::from([
4082 (
4083 profile_type_source,
4084 AcceptedNamedTypeIdentity::Composite(profile_type),
4085 ),
4086 (
4087 degree_type_source,
4088 AcceptedNamedTypeIdentity::Composite(degree_type),
4089 ),
4090 ]),
4091 BTreeMap::new(),
4092 BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
4093 );
4094 let candidate = accepted_schema_candidate_with_catalogs_for_tests(
4095 STORE_PATH,
4096 AcceptedSchemaRevision::INITIAL,
4097 enum_catalog,
4098 composite_catalog,
4099 source_bindings,
4100 BTreeMap::from([(entity_tag, snapshot)]),
4101 );
4102
4103 let session = DbSession::<TestCanister>::new(&STORE_REGISTRY);
4104 session
4105 .db
4106 .ensure_recovered_state()
4107 .expect("targeted mutation test database should initialize");
4108 let store = session
4109 .db
4110 .store_handle(STORE_PATH)
4111 .expect("targeted mutation test store should resolve");
4112 crate::db::commit::publish_accepted_schema_candidate(
4113 STORE_PATH,
4114 store,
4115 AcceptedSchemaRevision::NONE,
4116 &candidate,
4117 )
4118 .expect("targeted mutation candidate should publish");
4119
4120 let dynamic_error = session
4121 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4122 entity: "TargetedMutation".to_string(),
4123 patch: structural_patch(1, 12),
4124 })
4125 .expect_err("dynamic write must enforce the targeted rule");
4126 assert_eq!(
4127 targeted_constraint_id(&dynamic_error),
4128 targeted_rule_id.get()
4129 );
4130
4131 let binding = session
4132 .issue_typed_entity_binding(
4133 ENTITY_SOURCE,
4134 &[
4135 DynamicTypedFieldBindingRequest::new(
4136 ID_SOURCE.to_string(),
4137 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4138 false,
4139 ),
4140 DynamicTypedFieldBindingRequest::new(
4141 PROFILE_SOURCE.to_string(),
4142 DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
4143 false,
4144 ),
4145 DynamicTypedFieldBindingRequest::new(
4146 UPDATED_AT_SOURCE.to_string(),
4147 DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
4148 false,
4149 ),
4150 ],
4151 )
4152 .expect("targeted typed binding should issue");
4153 let typed_patch = binding
4154 .bind_write_fields(vec![
4155 (
4156 ID_SOURCE.to_string(),
4157 DynamicWriteCell::Value(InputValue::Nat64(2)),
4158 ),
4159 (
4160 PROFILE_SOURCE.to_string(),
4161 DynamicWriteCell::Value(profile_input(12)),
4162 ),
4163 ])
4164 .expect("targeted typed patch should bind");
4165 let typed_error = session
4166 .execute_trusted_typed_mutation(
4167 &binding,
4168 &DynamicTypedMutation::Insert { patch: typed_patch },
4169 )
4170 .expect_err("typed write must enforce the targeted rule");
4171 assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
4172
4173 #[cfg(feature = "sql")]
4174 {
4175 let sql_error = session
4176 .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
4177 .expect_err("SQL default resolution must enforce the targeted rule");
4178 let crate::db::QueryError::Execute(execute) = sql_error else {
4179 panic!("targeted SQL write should fail at shared execution admission");
4180 };
4181 assert_eq!(
4182 targeted_constraint_id(execute.as_internal()),
4183 targeted_rule_id.get()
4184 );
4185 }
4186
4187 session
4188 .execute_trusted_dynamic_mutation_batch(vec![
4189 DynamicMutation::Insert {
4190 entity: "TargetedMutation".to_string(),
4191 patch: structural_patch(4, 5),
4192 },
4193 DynamicMutation::Insert {
4194 entity: "TargetedMutation".to_string(),
4195 patch: structural_patch(5, 12),
4196 },
4197 ])
4198 .expect_err("one invalid targeted value must reject the whole batch");
4199 assert_eq!(
4200 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4201 Some(0),
4202 "no frontend or earlier valid batch row may escape targeted admission",
4203 );
4204
4205 let admitted = session
4206 .execute_trusted_dynamic_mutation_batch(vec![
4207 DynamicMutation::Insert {
4208 entity: "TargetedMutation".to_string(),
4209 patch: structural_patch(6, 5),
4210 },
4211 DynamicMutation::Insert {
4212 entity: "TargetedMutation".to_string(),
4213 patch: structural_patch(7, 10),
4214 },
4215 ])
4216 .expect("compliant targeted values should share one accepted batch");
4217 let [first, second] = admitted.rows.as_slice() else {
4218 panic!("the mixed targeted batch should return two rows");
4219 };
4220 let first_timestamp = first
4221 .get(2)
4222 .expect("the first mixed row should contain its managed timestamp");
4223 assert!(matches!(
4224 first_timestamp,
4225 crate::value::OutputValue::Timestamp(_)
4226 ));
4227 assert_eq!(
4228 second.get(2),
4229 Some(first_timestamp),
4230 "one accepted mixed batch must materialize one managed timestamp",
4231 );
4232 assert_eq!(
4233 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4234 Some(2),
4235 );
4236 }
4237}