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