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,
24 commit_structural_row_ops_with_mutation_progress_for_path,
25 commit_structural_row_ops_with_window_for_path, mutation_key_exists_error,
26 },
27 integrity::MutationProgressRecordOp,
28 schema::{
29 AcceptedFieldKind, AcceptedIdentityAllocation, AcceptedRowLayoutRuntimeContract,
30 FieldId, FieldInsertGeneration, IdentityStatementCursor, lower_field_type,
31 output_value_from_runtime,
32 },
33 write_context::{AcceptedWriteContext, MutationMode},
34 },
35 error::{InternalError, MutationDiagnosticContext},
36 metrics::sink::{MetricsEvent, SaveMutationKind, record},
37 traits::CanisterKind,
38 types::{CurrentTimestamp, Timestamp},
39 value::{InputValue, Value},
40};
41use icydb_schema::{EntitySourceKey, FieldSourceKey, FieldType, TypeSourceKey};
42
43#[derive(Clone, Debug, Eq, PartialEq)]
44struct AcceptedIdentityInsertField {
45 field_id: FieldId,
46 field_slot: usize,
47 accepted_kind: AcceptedFieldKind,
48}
49
50struct AcceptedStructuralMutationCommitOptions {
51 capture_output_values: bool,
52 packing: AcceptedStructuralMutationPacking,
53}
54
55impl AcceptedStructuralMutationCommitOptions {
56 const fn standard() -> Self {
57 Self {
58 capture_output_values: true,
59 packing: AcceptedStructuralMutationPacking::Complete,
60 }
61 }
62
63 #[cfg(test)]
64 const fn with_mutation_progress() -> Self {
65 Self {
66 capture_output_values: false,
67 packing: AcceptedStructuralMutationPacking::Complete,
68 }
69 }
70
71 const fn bounded_prefix() -> Self {
72 Self {
73 capture_output_values: false,
74 packing: AcceptedStructuralMutationPacking::BoundedPrefix,
75 }
76 }
77}
78
79#[derive(Clone, Copy)]
80enum AcceptedStructuralMutationPacking {
81 Complete,
82 BoundedPrefix,
83}
84
85pub(in crate::db::session) enum AcceptedStructuralMutationCommitDirective {
86 Standard,
87 WithMutationProgress(MutationProgressRecordOp),
88 Skip,
89}
90
91pub(in crate::db::session) enum AcceptedStructuralMutationTarget {
94 ResolveFromAfterImage,
95 Expected(Box<DecodedDataStoreKey>),
96 ExpectedLoaded(AcceptedLoadedStructuralRow),
97}
98
99pub(in crate::db::session) struct AcceptedLoadedStructuralRow {
102 key: Box<DecodedDataStoreKey>,
103 row: RawRow,
104}
105
106impl AcceptedLoadedStructuralRow {
107 pub(in crate::db::session) fn from_validated_parts(
108 key: DecodedDataStoreKey,
109 row: RawRow,
110 ) -> Self {
111 Self {
112 key: Box::new(key),
113 row,
114 }
115 }
116
117 fn into_parts(self) -> (DecodedDataStoreKey, RawRow) {
118 (*self.key, self.row)
119 }
120}
121
122impl AcceptedStructuralMutationTarget {
123 pub(in crate::db::session) fn expected(key: DecodedDataStoreKey) -> Self {
124 Self::Expected(Box::new(key))
125 }
126
127 pub(in crate::db::session) const fn expected_loaded(row: AcceptedLoadedStructuralRow) -> Self {
130 Self::ExpectedLoaded(row)
131 }
132}
133
134pub(in crate::db::session) enum AcceptedStructuralMutation {
137 Save {
138 mode: MutationMode,
139 target: AcceptedStructuralMutationTarget,
140 patch: AcceptedMutationIntentPatch,
141 },
142 Delete {
143 key: Box<DecodedDataStoreKey>,
144 },
145}
146
147impl AcceptedStructuralMutation {
148 pub(in crate::db::session) const fn save(
149 mode: MutationMode,
150 target: AcceptedStructuralMutationTarget,
151 patch: AcceptedMutationIntentPatch,
152 ) -> Self {
153 Self::Save {
154 mode,
155 target,
156 patch,
157 }
158 }
159
160 pub(in crate::db::session) fn delete(key: DecodedDataStoreKey) -> Self {
161 Self::Delete { key: Box::new(key) }
162 }
163}
164
165const MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS: usize = 4_096;
166pub(in crate::db::session) const STRUCTURAL_MUTATION_BATCH_STAGED_BYTES_POLICY: u32 =
167 16 * 1024 * 1024;
168pub(in crate::db::session) const MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES: usize =
169 STRUCTURAL_MUTATION_BATCH_STAGED_BYTES_POLICY as usize;
170const MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES: usize = 1024 * 1024;
171
172#[derive(Clone, Copy, Debug, Eq, PartialEq)]
173pub(in crate::db::session) struct AcceptedStructuralMutationPackingReport {
174 admitted_mutations: usize,
175 stopped_before_candidate: bool,
176 candidate_exceeds_batch_policy: bool,
177}
178
179impl AcceptedStructuralMutationPackingReport {
180 #[must_use]
181 pub(in crate::db::session) const fn admitted_mutations(self) -> usize {
182 self.admitted_mutations
183 }
184
185 #[must_use]
186 pub(in crate::db::session) const fn stopped_before_candidate(self) -> bool {
187 self.stopped_before_candidate
188 }
189
190 #[must_use]
191 pub(in crate::db::session) const fn candidate_exceeds_batch_policy(self) -> bool {
192 self.candidate_exceeds_batch_policy
193 }
194}
195
196fn structural_mutation_staged_charge(
197 lengths: impl IntoIterator<Item = usize>,
198) -> Result<usize, InternalError> {
199 lengths.into_iter().try_fold(0_usize, |total, length| {
200 total.checked_add(length).ok_or_else(|| {
201 InternalError::mutation_batch_staged_bytes_exceeded(
202 None,
203 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
204 )
205 })
206 })
207}
208
209fn add_structural_mutation_staged_bytes(
210 total: &mut usize,
211 lengths: impl IntoIterator<Item = usize>,
212) -> Result<(), InternalError> {
213 let charge = structural_mutation_staged_charge(lengths)?;
214 *total = total.checked_add(charge).ok_or_else(|| {
215 InternalError::mutation_batch_staged_bytes_exceeded(
216 None,
217 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
218 )
219 })?;
220 if *total > MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES {
221 return Err(InternalError::mutation_batch_staged_bytes_exceeded(
222 Some(*total),
223 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
224 ));
225 }
226 Ok(())
227}
228
229fn admit_structural_mutation_staged_charge(
230 total: &mut usize,
231 lengths: impl IntoIterator<Item = usize>,
232 packing: AcceptedStructuralMutationPacking,
233) -> Result<AcceptedStructuralMutationStagedAdmission, InternalError> {
234 if matches!(packing, AcceptedStructuralMutationPacking::Complete) {
235 add_structural_mutation_staged_bytes(total, lengths)?;
236 return Ok(AcceptedStructuralMutationStagedAdmission::Admitted);
237 }
238
239 let charge = structural_mutation_staged_charge(lengths)?;
240 if charge > MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES {
241 return Ok(AcceptedStructuralMutationStagedAdmission::CandidateExceedsPolicy);
242 }
243 let Some(next_total) = total.checked_add(charge) else {
244 return Ok(AcceptedStructuralMutationStagedAdmission::PageFull);
245 };
246 if next_total > MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES {
247 return Ok(AcceptedStructuralMutationStagedAdmission::PageFull);
248 }
249 *total = next_total;
250 Ok(AcceptedStructuralMutationStagedAdmission::Admitted)
251}
252
253#[derive(Clone, Copy, Debug, Eq, PartialEq)]
254enum AcceptedStructuralMutationStagedAdmission {
255 Admitted,
256 PageFull,
257 CandidateExceedsPolicy,
258}
259
260fn validate_structural_mutation_result_bytes(encoded_bytes: usize) -> Result<(), InternalError> {
261 if encoded_bytes > MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES {
262 return Err(InternalError::mutation_batch_result_bytes_exceeded(
263 encoded_bytes,
264 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES,
265 ));
266 }
267 Ok(())
268}
269
270pub(in crate::db::session) struct AcceptedStructuralMutationRow {
272 values: Vec<Value>,
273 logical_changed: bool,
274}
275
276impl AcceptedStructuralMutationRow {
277 #[cfg(any(feature = "sql", test))]
278 pub(in crate::db::session) fn into_values(self) -> Vec<Value> {
279 self.values
280 }
281
282 pub(in crate::db::session) const fn logical_changed(&self) -> bool {
283 self.logical_changed
284 }
285}
286
287const fn dynamic_mutation_mode(request: &DynamicMutation) -> Option<MutationMode> {
288 match request {
289 DynamicMutation::Insert { .. } => Some(MutationMode::Insert),
290 DynamicMutation::Update { .. } => Some(MutationMode::Update),
291 DynamicMutation::Replace { .. } => Some(MutationMode::Replace),
292 DynamicMutation::Delete { .. } => None,
293 }
294}
295
296const fn dynamic_typed_mutation_mode(request: &DynamicTypedMutation) -> MutationMode {
297 match request {
298 DynamicTypedMutation::Insert { .. } => MutationMode::Insert,
299 DynamicTypedMutation::Update { .. } => MutationMode::Update,
300 DynamicTypedMutation::Replace { .. } => MutationMode::Replace,
301 }
302}
303
304const fn diagnostic_mutation_operation(
305 mode: MutationMode,
306) -> icydb_diagnostic_code::DiagnosticMutationOperation {
307 match mode {
308 MutationMode::Insert => icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
309 MutationMode::Replace => icydb_diagnostic_code::DiagnosticMutationOperation::Replace,
310 MutationMode::Update => icydb_diagnostic_code::DiagnosticMutationOperation::Update,
311 }
312}
313
314const fn mutation_diagnostic_context(
315 entity_tag: crate::types::EntityTag,
316 mode: MutationMode,
317 batch_position: u32,
318) -> MutationDiagnosticContext {
319 MutationDiagnosticContext::new(
320 entity_tag.value(),
321 diagnostic_mutation_operation(mode),
322 batch_position,
323 )
324}
325
326const fn dynamic_write_context(operation_timestamp: Timestamp) -> AcceptedWriteContext {
327 AcceptedWriteContext::new(operation_timestamp)
328}
329
330fn insert_key_exists_after_generation(identity_generated: bool) -> InternalError {
331 if identity_generated {
332 InternalError::identity_state_corruption()
333 } else {
334 mutation_key_exists_error()
335 }
336}
337
338fn dynamic_key(
339 entity_tag: crate::types::EntityTag,
340 key: &InputValue,
341) -> Result<DecodedDataStoreKey, InternalError> {
342 let value = key
343 .clone()
344 .try_into_runtime_non_enum()
345 .ok_or_else(InternalError::executor_unsupported)?;
346 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
347}
348
349fn lower_dynamic_patch(
350 entity_path: &str,
351 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
352 patch: &DynamicStructuralPatch,
353 mode: MutationMode,
354 mutation_context: MutationDiagnosticContext,
355) -> Result<AcceptedMutationIntentPatch, InternalError> {
356 let mut lowered = AcceptedMutationIntentPatch::new();
357 for (field_name, cell) in patch.fields() {
358 let slot = descriptor
359 .field_slot_index_by_name(field_name)
360 .ok_or_else(|| {
361 InternalError::mutation_structural_field_unknown(entity_path, field_name)
362 })?;
363 let field = descriptor
364 .field_for_slot_index(slot)
365 .ok_or_else(InternalError::executor_invariant)?;
366 if !matches!(cell, DynamicWriteCell::Omitted)
367 && (field.write_policy().insert_generation().is_some()
368 || field.write_policy().write_management().is_some())
369 {
370 return Err(InternalError::mutation_database_owned_field_explicit(
371 mutation_context,
372 field.field_id().get(),
373 ));
374 }
375 let slot = FieldSlot::from_validated_index(slot);
376 lowered = match cell {
377 DynamicWriteCell::Omitted => lowered,
378 DynamicWriteCell::Default => match mode {
379 MutationMode::Insert | MutationMode::Replace => {
380 lowered.set_explicit_insert_default(slot)
381 }
382 MutationMode::Update => lowered.set_explicit_update_default(slot),
383 },
384 DynamicWriteCell::Null => lowered.set_authored(slot, InputValue::Null),
385 DynamicWriteCell::Value(value) => lowered.set_authored(slot, value.clone()),
386 };
387 }
388 Ok(lowered)
389}
390
391fn lower_dynamic_mutation_intent(
392 entity_tag: crate::types::EntityTag,
393 entity_path: &str,
394 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
395 request: &DynamicMutation,
396 batch_position: u32,
397) -> Result<(AcceptedStructuralMutation, Option<SaveMutationKind>), InternalError> {
398 match request {
399 DynamicMutation::Insert { patch, .. } => {
400 let mode = MutationMode::Insert;
401 Ok((
402 AcceptedStructuralMutation::save(
403 mode,
404 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
405 lower_dynamic_patch(
406 entity_path,
407 descriptor,
408 patch,
409 mode,
410 mutation_diagnostic_context(entity_tag, mode, batch_position),
411 )?,
412 ),
413 Some(SaveMutationKind::Insert),
414 ))
415 }
416 DynamicMutation::Update { key, patch, .. }
417 | DynamicMutation::Replace { key, patch, .. } => {
418 let mode =
419 dynamic_mutation_mode(request).ok_or_else(InternalError::executor_invariant)?;
420 let kind = match mode {
421 MutationMode::Insert => SaveMutationKind::Insert,
422 MutationMode::Replace => SaveMutationKind::Replace,
423 MutationMode::Update => SaveMutationKind::Update,
424 };
425 Ok((
426 AcceptedStructuralMutation::save(
427 mode,
428 AcceptedStructuralMutationTarget::expected(dynamic_key(entity_tag, key)?),
429 lower_dynamic_patch(
430 entity_path,
431 descriptor,
432 patch,
433 mode,
434 mutation_diagnostic_context(entity_tag, mode, batch_position),
435 )?,
436 ),
437 Some(kind),
438 ))
439 }
440 DynamicMutation::Delete { key, .. } => Ok((
441 AcceptedStructuralMutation::delete(dynamic_key(entity_tag, key)?),
442 None,
443 )),
444 }
445}
446
447fn lower_typed_patch(
448 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
449 patch: &DynamicTypedStructuralPatch,
450 mode: MutationMode,
451 mutation_context: MutationDiagnosticContext,
452) -> Result<AcceptedMutationIntentPatch, InternalError> {
453 let mut lowered = AcceptedMutationIntentPatch::new();
454 for (field_id, slot, cell) in patch.fields() {
455 let slot_index = usize::from(*slot);
456 let field = descriptor
457 .field_for_slot_index(slot_index)
458 .ok_or_else(InternalError::store_invariant)?;
459 if field.field_id().get() != *field_id {
460 return Err(InternalError::store_invariant());
461 }
462 if !matches!(cell, DynamicWriteCell::Omitted)
463 && (field.write_policy().insert_generation().is_some()
464 || field.write_policy().write_management().is_some())
465 {
466 return Err(InternalError::mutation_database_owned_field_explicit(
467 mutation_context,
468 field.field_id().get(),
469 ));
470 }
471 let slot = FieldSlot::from_validated_index(slot_index);
472 lowered = match cell {
473 DynamicWriteCell::Omitted => lowered,
474 DynamicWriteCell::Default => match mode {
475 MutationMode::Insert | MutationMode::Replace => {
476 lowered.set_explicit_insert_default(slot)
477 }
478 MutationMode::Update => lowered.set_explicit_update_default(slot),
479 },
480 DynamicWriteCell::Null => lowered.set_authored(slot, InputValue::Null),
481 DynamicWriteCell::Value(value) => lowered.set_authored(slot, value.clone()),
482 };
483 }
484 Ok(lowered)
485}
486
487fn preserve_dynamic_replacement_identity(
488 key: &DecodedDataStoreKey,
489 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
490 mut patch: AcceptedMutationIntentPatch,
491) -> Result<AcceptedMutationIntentPatch, InternalError> {
492 let primary_key_slots = descriptor.primary_key_slot_indices();
493 let runtime_key = key.primary_key_runtime_value();
494 let components = match runtime_key {
495 Value::List(values) if primary_key_slots.len() > 1 => values,
496 value if primary_key_slots.len() == 1 => vec![value],
497 _ => return Err(InternalError::executor_invariant()),
498 };
499 if components.len() != primary_key_slots.len() {
500 return Err(InternalError::executor_invariant());
501 }
502
503 for (slot, value) in primary_key_slots.iter().copied().zip(components) {
504 let _ = descriptor
505 .field_for_slot_index(slot)
506 .ok_or_else(InternalError::executor_invariant)?;
507 let has_explicit_intent = patch
508 .entries()
509 .iter()
510 .any(|entry| entry.slot().index() == slot);
511 if has_explicit_intent {
512 continue;
513 }
514 let value = InputValue::try_from_runtime_non_enum(&value)
515 .ok_or_else(InternalError::executor_invariant)?;
516 patch =
517 patch.set_preserved_replacement_identity(FieldSlot::from_validated_index(slot), value);
518 }
519
520 Ok(patch)
521}
522
523fn accepted_identity_insert_field(
527 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
528) -> Result<Option<AcceptedIdentityInsertField>, InternalError> {
529 let mut identity = None;
530 for field in descriptor.fields() {
531 if field.write_policy().insert_generation() != Some(FieldInsertGeneration::Identity) {
532 continue;
533 }
534 let field_slot = usize::from(field.slot().get());
535 if identity
536 .replace(AcceptedIdentityInsertField {
537 field_id: field.field_id(),
538 field_slot,
539 accepted_kind: field.kind().clone(),
540 })
541 .is_some()
542 || descriptor.primary_key_slot_indices() != [field_slot]
543 {
544 return Err(InternalError::identity_corruption());
545 }
546 }
547 Ok(identity)
548}
549
550fn checked_pre_key_candidate_count(count: usize) -> Result<u32, InternalError> {
551 u32::try_from(count).map_err(|_| InternalError::identity_candidate_count_exhausted())
552}
553
554fn validate_identity_materialization(
555 entity_tag: crate::types::EntityTag,
556 identity_field: &AcceptedIdentityInsertField,
557 candidate: &AcceptedPreKeyInsert,
558 allocation: &AcceptedIdentityAllocation,
559 data_key: &DecodedDataStoreKey,
560 reader: &StructuralSlotReader<'_>,
561) -> Result<(), InternalError> {
562 let owner = allocation.owner();
563 let slot_value = reader.required_cached_value(identity_field.field_slot)?;
564 if candidate.entity_tag() != entity_tag
565 || candidate.input_ordinal() != allocation.input_ordinal()
566 || owner.entity_tag() != entity_tag
567 || owner.field_id() != identity_field.field_id
568 || allocation.field_slot() != identity_field.field_slot
569 || slot_value != allocation.value()
570 || data_key.primary_key_runtime_value() != *allocation.value()
571 {
572 return Err(InternalError::identity_corruption());
573 }
574 Ok(())
575}
576
577fn data_key_from_row(
578 entity_tag: crate::types::EntityTag,
579 contract: &StructuralRowContract,
580 row: &RawRow,
581) -> Result<DecodedDataStoreKey, InternalError> {
582 let reader =
583 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(row, contract)?;
584 let values = contract
585 .primary_key_slot_indices()
586 .iter()
587 .map(|slot| reader.required_cached_value(*slot).cloned())
588 .collect::<Result<Vec<_>, _>>()?;
589 let value = match values.as_slice() {
590 [value] => value.clone(),
591 _ => Value::List(values),
592 };
593 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
594}
595
596#[cfg(feature = "sql")]
597pub(in crate::db::session) fn structural_data_key_from_runtime_values(
598 entity_tag: crate::types::EntityTag,
599 values: Vec<Value>,
600) -> Result<DecodedDataStoreKey, InternalError> {
601 let value = match values.as_slice() {
602 [value] => value.clone(),
603 _ => Value::List(values),
604 };
605 DecodedDataStoreKey::try_from_structural_key(entity_tag, &value)
606}
607
608fn validated_existing_row(
609 store: crate::db::registry::StoreHandle,
610 data_key: &DecodedDataStoreKey,
611 contract: &StructuralRowContract,
612) -> Result<Option<RawRow>, InternalError> {
613 let raw_key = data_key.to_raw()?;
614 let row = store.with_data(|data| data.get(&raw_key));
615 if let Some(row) = row.as_ref() {
616 let reader =
617 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(row, contract)?;
618 reader.validate_primary_key(data_key)?;
619 }
620 Ok(row)
621}
622
623fn prepare_dynamic_mutation_result(
624 catalog: &AcceptedSchemaCatalogContext,
625 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
626 rows: Vec<AcceptedStructuralMutationRow>,
627 enforce_mixed_batch_result_bound: bool,
628) -> Result<DynamicMutationResult, InternalError> {
629 let affected_rows = rows.iter().try_fold(0_u32, |total, row| {
630 total
631 .checked_add(u32::from(row.logical_changed()))
632 .ok_or_else(InternalError::executor_invariant)
633 })?;
634 let columns = descriptor
635 .fields()
636 .iter()
637 .map(|field| field.name().to_string())
638 .collect();
639 let rows = rows
640 .into_iter()
641 .map(|row| {
642 row.values
643 .iter()
644 .map(|value| {
645 output_value_from_runtime(catalog.enum_catalog(), value)
646 .map_err(|_| InternalError::store_invariant())
647 })
648 .collect::<Result<Vec<_>, _>>()
649 })
650 .collect::<Result<Vec<_>, _>>()?;
651 let result = DynamicMutationResult {
652 entity: catalog.snapshot().entity_name().to_string(),
653 columns,
654 rows,
655 affected_rows,
656 };
657 if enforce_mixed_batch_result_bound {
658 let encoded =
659 candid::encode_one(&result).map_err(|_| InternalError::executor_invariant())?;
660 validate_structural_mutation_result_bytes(encoded.len())?;
661 }
662 Ok(result)
663}
664
665fn dynamic_typed_field_type(
666 field_type: DynamicTypedFieldType,
667) -> Result<FieldType, DynamicTypedBindingError> {
668 match field_type {
669 DynamicTypedFieldType::Scalar(scalar) => Ok(FieldType::Scalar(scalar)),
670 DynamicTypedFieldType::List(item) => {
671 Ok(FieldType::List(Box::new(dynamic_typed_field_type(*item)?)))
672 }
673 DynamicTypedFieldType::Named(source_key) => TypeSourceKey::try_new(source_key)
674 .map(FieldType::Named)
675 .map_err(|_| DynamicTypedBindingError::FieldUnavailable),
676 }
677}
678
679fn typed_adapter_field_kind_matches(
680 accepted: &AcceptedFieldKind,
681 expected: &AcceptedFieldKind,
682) -> bool {
683 if accepted == expected {
684 return true;
685 }
686 match (accepted, expected) {
687 (AcceptedFieldKind::Relation { key_kind, .. }, expected) => {
688 typed_adapter_field_kind_matches(key_kind, expected)
689 }
690 (AcceptedFieldKind::List(accepted), AcceptedFieldKind::List(expected)) => {
691 typed_adapter_field_kind_matches(accepted, expected)
692 }
693 _ => false,
694 }
695}
696
697impl<C: CanisterKind> DbSession<C> {
698 pub fn issue_typed_entity_binding(
700 &self,
701 entity_source_key: &str,
702 field_requests: &[DynamicTypedFieldBindingRequest],
703 ) -> Result<DynamicTypedEntityBinding, DynamicTypedBindingError> {
704 let entity_source = EntitySourceKey::try_new(entity_source_key)
705 .map_err(|_| DynamicTypedBindingError::FieldUnavailable)?;
706 let field_requests = field_requests
707 .iter()
708 .map(|request| {
709 Ok((
710 FieldSourceKey::try_new(request.source_key.clone())
711 .map_err(|_| DynamicTypedBindingError::FieldUnavailable)?,
712 dynamic_typed_field_type(request.field_type.clone())?,
713 request.nullable,
714 ))
715 })
716 .collect::<Result<Vec<_>, DynamicTypedBindingError>>()?;
717 let catalog = self
718 .find_accepted_schema_catalog_context_for_entity_source_key(entity_source.as_str())?
719 .ok_or(DynamicTypedBindingError::FieldUnavailable)?;
720 let identity = catalog.identity();
721 if identity.entity_path() != entity_source.as_str() {
722 return Err(InternalError::store_invariant().into());
723 }
724 let store = self.db.recovered_store(identity.store_path())?;
725 let bundle = store
726 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
727 .ok_or_else(InternalError::store_invariant)?;
728 let entity_tag = identity.entity_tag();
729 if bundle.source_bindings().entity(&entity_source) != Some(entity_tag)
730 || bundle.revision() != catalog.revision()
731 {
732 return Err(InternalError::store_invariant().into());
733 }
734 let snapshot = bundle
735 .entity_snapshots()
736 .get(&entity_tag)
737 .ok_or_else(InternalError::store_invariant)?;
738 let row_contract = catalog.inspection_plan().row_contract();
739 let mut fields = Vec::with_capacity(field_requests.len());
740 for (source, field_type, nullable) in &field_requests {
741 let field_id = bundle
742 .source_bindings()
743 .field(entity_tag, source)
744 .ok_or(DynamicTypedBindingError::FieldUnavailable)?;
745 let field = snapshot
746 .fields()
747 .iter()
748 .find(|field| field.id() == field_id)
749 .ok_or_else(InternalError::store_invariant)?;
750 let runtime_field =
751 row_contract.required_accepted_field_contract(usize::from(field.slot().get()))?;
752 if runtime_field.field_id() != field_id {
753 return Err(InternalError::store_invariant().into());
754 }
755 let expected_kind = lower_field_type(field_type, bundle.source_bindings())
756 .map_err(|_| DynamicTypedBindingError::IncompatibleField)?;
757 if field.nullable() != *nullable
758 || !typed_adapter_field_kind_matches(field.kind(), &expected_kind)
759 {
760 return Err(DynamicTypedBindingError::IncompatibleField);
761 }
762 fields.push((
763 source.as_str().to_string(),
764 field_id.get(),
765 field.slot().get(),
766 field.name().to_string(),
767 ));
768 }
769 let adapter_names = bundle.typed_adapter_names()?;
770
771 DynamicTypedEntityBinding::new(
772 database_incarnation_id()?.to_bytes(),
773 entity_source.as_str().to_string(),
774 snapshot.entity_name().to_string(),
775 entity_tag.value(),
776 catalog.revision().get(),
777 catalog.fingerprint(),
778 row_contract.current_layout_version().get(),
779 fields,
780 adapter_names.named_types,
781 adapter_names.enum_variants,
782 adapter_names.composite_fields,
783 )
784 .map_err(Into::into)
785 }
786
787 pub(in crate::db::session) fn current_typed_entity_binding_catalog(
788 &self,
789 binding: &DynamicTypedEntityBinding,
790 ) -> Result<Option<AcceptedSchemaCatalogContext>, InternalError> {
791 if database_incarnation_id()?.to_bytes() != binding.database_incarnation {
792 return Ok(None);
793 }
794 let Some(catalog) = self.find_accepted_schema_catalog_context_for_entity_source_key(
795 binding.entity_source.as_str(),
796 )?
797 else {
798 return Ok(None);
799 };
800 let row_contract = catalog.inspection_plan().row_contract();
801 let identity = catalog.identity();
802 if identity.entity_path() != binding.entity_source.as_str()
803 || identity.entity_tag().value() != binding.entity_tag
804 || catalog.revision().get() != binding.accepted_revision
805 || catalog.fingerprint() != binding.accepted_fingerprint
806 || row_contract.current_layout_version().get() != binding.entity_generation
807 {
808 return Ok(None);
809 }
810 let entity_source = EntitySourceKey::try_new(binding.entity_source.clone())
811 .map_err(|_| InternalError::store_invariant())?;
812 let store = self.db.recovered_store(identity.store_path())?;
813 let bundle = store
814 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
815 .ok_or_else(InternalError::store_invariant)?;
816 if bundle.revision() != catalog.revision()
817 || bundle.source_bindings().entity(&entity_source) != Some(identity.entity_tag())
818 {
819 return Ok(None);
820 }
821 let snapshot = bundle
822 .entity_snapshots()
823 .get(&identity.entity_tag())
824 .ok_or_else(InternalError::store_invariant)?;
825 for (source_key, expected_field_id, expected_slot) in binding.field_identity_bindings() {
826 let source = FieldSourceKey::try_new(source_key)
827 .map_err(|_| InternalError::store_invariant())?;
828 let Some(field_id) = bundle
829 .source_bindings()
830 .field(identity.entity_tag(), &source)
831 else {
832 return Ok(None);
833 };
834 let Some(field) = snapshot
835 .fields()
836 .iter()
837 .find(|field| field.id() == field_id)
838 else {
839 return Err(InternalError::store_invariant());
840 };
841 if field_id.get() != expected_field_id || field.slot().get() != expected_slot {
842 return Ok(None);
843 }
844 }
845 Ok(Some(catalog))
846 }
847
848 pub fn typed_entity_binding_is_current(
850 &self,
851 binding: &DynamicTypedEntityBinding,
852 ) -> Result<bool, InternalError> {
853 self.current_typed_entity_binding_catalog(binding)
854 .map(|catalog| catalog.is_some())
855 }
856
857 #[cfg(feature = "sql")]
860 pub(in crate::db::session) fn execute_accepted_structural_delete_batch(
861 &self,
862 catalog: &AcceptedSchemaCatalogContext,
863 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
864 keys: Vec<DecodedDataStoreKey>,
865 precommit_validation: impl FnOnce(&[Vec<Value>]) -> Result<(), InternalError>,
866 ) -> Result<Vec<Vec<Value>>, InternalError> {
867 let mutations = keys
868 .into_iter()
869 .map(AcceptedStructuralMutation::delete)
870 .collect::<Vec<_>>();
871 let mutation_capacity = mutations.len();
872 let mut mutations = mutations.into_iter();
873 self.execute_accepted_structural_mutation_batch_inner(
874 catalog,
875 descriptor,
876 mutation_capacity,
877 0,
878 || Ok(mutations.next()),
879 Timestamp::now(),
880 AcceptedStructuralMutationCommitOptions::standard(),
881 |rows, _report| {
882 let rows = rows
883 .into_iter()
884 .map(AcceptedStructuralMutationRow::into_values)
885 .collect::<Vec<_>>();
886 precommit_validation(rows.as_slice())?;
887 Ok((rows, AcceptedStructuralMutationCommitDirective::Standard))
888 },
889 )
890 }
891
892 pub(in crate::db::session) fn execute_accepted_structural_save_batch<T>(
900 &self,
901 catalog: &AcceptedSchemaCatalogContext,
902 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
903 mutations: Vec<AcceptedStructuralMutation>,
904 operation_timestamp: Timestamp,
905 precommit_preparation: impl FnOnce(
906 Vec<AcceptedStructuralMutationRow>,
907 ) -> Result<T, InternalError>,
908 ) -> Result<T, InternalError> {
909 let mutation_capacity = mutations.len();
910 let identity_candidate_count = mutations
911 .iter()
912 .filter(|mutation| {
913 matches!(
914 mutation,
915 AcceptedStructuralMutation::Save {
916 mode: MutationMode::Insert,
917 target: AcceptedStructuralMutationTarget::ResolveFromAfterImage,
918 ..
919 }
920 )
921 })
922 .count();
923 let mut mutations = mutations.into_iter();
924 self.execute_accepted_structural_mutation_batch_inner(
925 catalog,
926 descriptor,
927 mutation_capacity,
928 identity_candidate_count,
929 || Ok(mutations.next()),
930 operation_timestamp,
931 AcceptedStructuralMutationCommitOptions::standard(),
932 |rows, _report| {
933 precommit_preparation(rows).map(|prepared| {
934 (
935 prepared,
936 AcceptedStructuralMutationCommitDirective::Standard,
937 )
938 })
939 },
940 )
941 }
942
943 #[cfg(test)]
945 pub(in crate::db::session) fn execute_accepted_structural_update_with_mutation_progress(
946 &self,
947 catalog: &AcceptedSchemaCatalogContext,
948 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
949 mutations: Vec<AcceptedStructuralMutation>,
950 operation_timestamp: Timestamp,
951 mutation_progress: MutationProgressRecordOp,
952 ) -> Result<usize, InternalError> {
953 let mutation_capacity = mutations.len();
954 let mut mutations = mutations.into_iter();
955 self.execute_accepted_structural_mutation_batch_inner(
956 catalog,
957 descriptor,
958 mutation_capacity,
959 0,
960 || Ok(mutations.next()),
961 operation_timestamp,
962 AcceptedStructuralMutationCommitOptions::with_mutation_progress(),
963 |rows, _report| {
964 Ok((
965 rows.len(),
966 AcceptedStructuralMutationCommitDirective::WithMutationProgress(
967 mutation_progress,
968 ),
969 ))
970 },
971 )
972 }
973
974 #[cfg(any(feature = "sql", test))]
977 pub(in crate::db::session) fn execute_accepted_structural_update_bounded_prefix<T>(
978 &self,
979 catalog: &AcceptedSchemaCatalogContext,
980 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
981 mutation_capacity: usize,
982 mut next_mutation: impl FnMut() -> Result<Option<AcceptedStructuralMutation>, InternalError>,
983 operation_timestamp: Timestamp,
984 precommit_preparation: impl FnOnce(
985 AcceptedStructuralMutationPackingReport,
986 ) -> Result<
987 (T, AcceptedStructuralMutationCommitDirective),
988 InternalError,
989 >,
990 ) -> Result<T, InternalError> {
991 self.execute_accepted_structural_mutation_batch_inner(
992 catalog,
993 descriptor,
994 mutation_capacity,
995 0,
996 &mut next_mutation,
997 operation_timestamp,
998 AcceptedStructuralMutationCommitOptions::bounded_prefix(),
999 |rows, report| {
1000 if rows.len() != report.admitted_mutations() {
1001 return Err(InternalError::executor_invariant());
1002 }
1003 precommit_preparation(report)
1004 },
1005 )
1006 }
1007
1008 #[expect(
1009 clippy::too_many_arguments,
1010 clippy::too_many_lines,
1011 reason = "one phased owner keeps accepted authority, mutation context, precommit preparation, output capture, and commit staging inseparable"
1012 )]
1013 fn execute_accepted_structural_mutation_batch_inner<T>(
1014 &self,
1015 catalog: &AcceptedSchemaCatalogContext,
1016 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
1017 mutation_capacity: usize,
1018 identity_candidate_count: usize,
1019 mut next_mutation: impl FnMut() -> Result<Option<AcceptedStructuralMutation>, InternalError>,
1020 operation_timestamp: Timestamp,
1021 options: AcceptedStructuralMutationCommitOptions,
1022 precommit_preparation: impl FnOnce(
1023 Vec<AcceptedStructuralMutationRow>,
1024 AcceptedStructuralMutationPackingReport,
1025 ) -> Result<
1026 (T, AcceptedStructuralMutationCommitDirective),
1027 InternalError,
1028 >,
1029 ) -> Result<T, InternalError> {
1030 let identity = catalog.identity();
1031 let AcceptedStructuralMutationCommitOptions {
1032 capture_output_values,
1033 packing,
1034 } = options;
1035 let entity_path = identity.entity_path();
1036 let store_path = identity.store_path();
1037 let row_decode_contract =
1038 descriptor.row_decode_contract(catalog.value_catalog_handle().clone());
1039 let row_contract = StructuralRowContract::from_accepted_decode_contract(
1040 entity_path,
1041 row_decode_contract.clone(),
1042 );
1043 let store = self.db.recovered_store(store_path)?;
1044 let write_context = dynamic_write_context(operation_timestamp);
1045 let identity_field = accepted_identity_insert_field(descriptor)?;
1046 let identity_incarnation = identity_field
1047 .as_ref()
1048 .map(|_| database_incarnation_id())
1049 .transpose()?;
1050 if mutation_capacity > MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS {
1051 return Err(InternalError::mutation_batch_too_many_items(
1052 mutation_capacity,
1053 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1054 ));
1055 }
1056 let _ = checked_pre_key_candidate_count(identity_candidate_count)?;
1057 let mut identity_cursor: Option<IdentityStatementCursor> = None;
1058 let mut identity_insert_ordinal = 0_u32;
1059 let mut scheduler = AcceptedMutationConstraintScheduler::new(
1060 entity_path,
1061 identity.entity_tag(),
1062 row_decode_contract.clone(),
1063 catalog.fingerprint(),
1064 catalog.fingerprint_method_version(),
1065 catalog.accepted_row_constraints(),
1066 mutation_capacity,
1067 );
1068 let mut output = Vec::with_capacity(mutation_capacity);
1069 let mut staged_bytes = 0_usize;
1070 let mut stopped_before_candidate = false;
1071 let mut candidate_exceeds_batch_policy = false;
1072 let mut input_index = 0_usize;
1073
1074 while let Some(mutation) = next_mutation()? {
1075 if input_index >= mutation_capacity {
1076 return Err(InternalError::mutation_batch_too_many_items(
1077 input_index.saturating_add(1),
1078 mutation_capacity,
1079 ));
1080 }
1081 let batch_input_ordinal = u32::try_from(input_index).map_err(|_| {
1082 InternalError::mutation_batch_too_many_items(
1083 mutation_capacity,
1084 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1085 )
1086 })?;
1087 input_index = input_index.saturating_add(1);
1088 let AcceptedStructuralMutation::Save {
1089 mode,
1090 target,
1091 patch: authored_patch,
1092 } = mutation
1093 else {
1094 let AcceptedStructuralMutation::Delete { key } = mutation else {
1095 return Err(InternalError::executor_invariant());
1096 };
1097 let before = validated_existing_row(store, &key, &row_contract)?
1098 .ok_or_else(|| InternalError::store_not_found(&key))?;
1099 let raw_key = key.to_raw()?;
1100 let canonical_before = canonical_row_from_raw_row_with_accepted_decode_contract(
1101 entity_path,
1102 row_decode_contract.clone(),
1103 &before,
1104 )?;
1105 let admission = admit_structural_mutation_staged_charge(
1106 &mut staged_bytes,
1107 [
1108 raw_key.as_bytes().len(),
1109 canonical_before.as_raw_row().as_bytes().len(),
1110 ],
1111 packing,
1112 )?;
1113 match admission {
1114 AcceptedStructuralMutationStagedAdmission::Admitted => {}
1115 AcceptedStructuralMutationStagedAdmission::PageFull => {
1116 stopped_before_candidate = true;
1117 break;
1118 }
1119 AcceptedStructuralMutationStagedAdmission::CandidateExceedsPolicy => {
1120 stopped_before_candidate = true;
1121 candidate_exceeds_batch_policy = true;
1122 break;
1123 }
1124 }
1125 scheduler.schedule_delete(
1126 CommitRowOp::new(
1127 entity_path,
1128 raw_key,
1129 Some(canonical_before.as_raw_row().as_bytes().to_vec()),
1130 None,
1131 catalog.fingerprint(),
1132 ),
1133 batch_input_ordinal,
1134 )?;
1135 let reader = StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
1136 canonical_before.as_raw_row(),
1137 &row_contract,
1138 )?;
1139 let values = if capture_output_values {
1140 let mut values = Vec::with_capacity(descriptor.fields().len());
1141 for field in descriptor.fields() {
1142 values.push(
1143 reader
1144 .required_cached_value(usize::from(field.slot().get()))?
1145 .clone(),
1146 );
1147 }
1148 values
1149 } else {
1150 Vec::new()
1151 };
1152 output.push(AcceptedStructuralMutationRow {
1153 values,
1154 logical_changed: true,
1155 });
1156 continue;
1157 };
1158 let mutation_context =
1159 mutation_diagnostic_context(identity.entity_tag(), mode, batch_input_ordinal);
1160 let (expected_key, preloaded_before, pre_key_insert, mut keyed_patch) = match target {
1161 AcceptedStructuralMutationTarget::ResolveFromAfterImage => {
1162 let candidate_ordinal =
1163 if identity_field.is_some() && matches!(mode, MutationMode::Insert) {
1164 identity_insert_ordinal
1165 } else {
1166 batch_input_ordinal
1167 };
1168 (
1169 None,
1170 None,
1171 Some(AcceptedPreKeyInsert::new(
1172 identity.entity_tag(),
1173 authored_patch,
1174 candidate_ordinal,
1175 )),
1176 None,
1177 )
1178 }
1179 AcceptedStructuralMutationTarget::Expected(key) => {
1180 (Some(*key), None, None, Some(authored_patch))
1181 }
1182 AcceptedStructuralMutationTarget::ExpectedLoaded(loaded) => {
1183 let (key, row) = loaded.into_parts();
1184 (Some(key), Some(row), None, Some(authored_patch))
1185 }
1186 };
1187 if matches!(mode, MutationMode::Replace)
1188 && let Some(key) = expected_key.as_ref()
1189 {
1190 let patch = keyed_patch
1191 .take()
1192 .ok_or_else(InternalError::executor_invariant)?;
1193 keyed_patch = Some(preserve_dynamic_replacement_identity(
1194 key, descriptor, patch,
1195 )?);
1196 }
1197 let patch = pre_key_insert
1198 .as_ref()
1199 .map(AcceptedPreKeyInsert::fields)
1200 .or(keyed_patch.as_ref())
1201 .ok_or_else(InternalError::executor_invariant)?;
1202 let before = match (expected_key.as_ref(), preloaded_before) {
1203 (Some(_), Some(row)) => Some(row),
1204 (Some(key), None) => validated_existing_row(store, key, &row_contract)?,
1205 (None, None) => None,
1206 (None, Some(_)) => return Err(InternalError::executor_invariant()),
1207 };
1208 match mode {
1209 MutationMode::Insert if before.is_some() => {
1210 return Err(mutation_key_exists_error());
1211 }
1212 MutationMode::Update if before.is_none() => {
1213 let key = expected_key
1214 .as_ref()
1215 .ok_or_else(InternalError::executor_invariant)?;
1216 return Err(InternalError::store_not_found(key));
1217 }
1218 MutationMode::Insert | MutationMode::Replace | MutationMode::Update => {}
1219 }
1220
1221 let identity_allocation = if let Some(identity_field) = identity_field.as_ref()
1222 && matches!(mode, MutationMode::Insert)
1223 && before.is_none()
1224 {
1225 let candidate = pre_key_insert.as_ref().ok_or_else(|| {
1226 InternalError::mutation_database_owned_field_explicit(
1227 mutation_context,
1228 identity_field.field_id.get(),
1229 )
1230 })?;
1231 if identity_cursor.is_none() {
1232 let incarnation = identity_incarnation
1233 .ok_or_else(InternalError::identity_state_corruption)?;
1234 identity_cursor = Some(store.with_schema(|schema_store| {
1235 schema_store.identity_statement_cursor(
1236 incarnation,
1237 identity.entity_tag(),
1238 identity_field.field_id,
1239 &identity_field.accepted_kind,
1240 )
1241 })?);
1242 }
1243 let allocation = identity_cursor
1244 .as_mut()
1245 .ok_or_else(InternalError::identity_state_corruption)?
1246 .allocate(identity_field.field_slot, candidate.input_ordinal())?;
1247 identity_insert_ordinal = identity_insert_ordinal
1248 .checked_add(1)
1249 .ok_or_else(InternalError::identity_candidate_count_exhausted)?;
1250 Some(allocation)
1251 } else if let Some(identity_field) = identity_field.as_ref()
1252 && matches!(mode, MutationMode::Replace)
1253 && before.is_none()
1254 {
1255 return Err(InternalError::mutation_database_owned_field_explicit(
1256 mutation_context,
1257 identity_field.field_id.get(),
1258 ));
1259 } else {
1260 None
1261 };
1262
1263 let resolved = match (mode, before.as_ref()) {
1264 (MutationMode::Insert | MutationMode::Replace, None) => {
1265 resolve_insert_structural_patch_with_accepted_contract(
1266 entity_path,
1267 row_decode_contract.clone(),
1268 catalog.fingerprint(),
1269 catalog.accepted_row_constraints(),
1270 patch,
1271 write_context,
1272 mutation_context,
1273 identity_allocation.as_ref(),
1274 )?
1275 }
1276 (MutationMode::Update, Some(before)) => {
1277 resolve_update_structural_patch_with_accepted_contract(
1278 entity_path,
1279 row_decode_contract.clone(),
1280 catalog.fingerprint(),
1281 catalog.accepted_row_constraints(),
1282 before,
1283 patch,
1284 write_context,
1285 mutation_context,
1286 )?
1287 }
1288 (MutationMode::Replace, Some(before)) => {
1289 resolve_existing_replace_structural_patch_with_accepted_contract(
1290 entity_path,
1291 row_decode_contract.clone(),
1292 catalog.fingerprint(),
1293 catalog.accepted_row_constraints(),
1294 before,
1295 patch,
1296 write_context,
1297 mutation_context,
1298 )?
1299 }
1300 (MutationMode::Insert, Some(_)) | (MutationMode::Update, None) => {
1301 return Err(InternalError::executor_invariant());
1302 }
1303 };
1304 let (after, provenance) = resolved.into_parts();
1305 let reader = StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
1306 after.as_raw_row(),
1307 &row_contract,
1308 )?;
1309 let data_key = match expected_key {
1310 Some(key) => {
1311 reader.validate_primary_key(&key)?;
1312 key
1313 }
1314 None => {
1315 data_key_from_row(identity.entity_tag(), &row_contract, after.as_raw_row())?
1316 }
1317 };
1318 if let Some(allocation) = identity_allocation.as_ref() {
1319 validate_identity_materialization(
1320 identity.entity_tag(),
1321 identity_field
1322 .as_ref()
1323 .ok_or_else(InternalError::identity_corruption)?,
1324 pre_key_insert
1325 .as_ref()
1326 .ok_or_else(InternalError::identity_corruption)?,
1327 allocation,
1328 &data_key,
1329 &reader,
1330 )?;
1331 }
1332 if matches!(mode, MutationMode::Insert)
1333 && validated_existing_row(store, &data_key, &row_contract)?.is_some()
1334 {
1335 return Err(insert_key_exists_after_generation(
1336 identity_allocation.is_some(),
1337 ));
1338 }
1339 let raw_key = data_key.to_raw()?;
1340 let canonical_before = before
1341 .as_ref()
1342 .map(|before| {
1343 canonical_row_from_raw_row_with_accepted_decode_contract(
1344 entity_path,
1345 row_decode_contract.clone(),
1346 before,
1347 )
1348 })
1349 .transpose()?;
1350 let logical_changed = canonical_before.as_ref().is_none_or(|before| {
1351 before.as_raw_row().as_bytes() != after.as_raw_row().as_bytes()
1352 });
1353 let physical_changed = before
1354 .as_ref()
1355 .is_none_or(|before| before.as_bytes() != after.as_raw_row().as_bytes());
1356 let admission = admit_structural_mutation_staged_charge(
1357 &mut staged_bytes,
1358 [
1359 raw_key.as_bytes().len(),
1360 canonical_before
1361 .as_ref()
1362 .map_or(0, |before| before.as_raw_row().as_bytes().len()),
1363 after.as_raw_row().as_bytes().len(),
1364 ],
1365 packing,
1366 )?;
1367 match admission {
1368 AcceptedStructuralMutationStagedAdmission::Admitted => {}
1369 AcceptedStructuralMutationStagedAdmission::PageFull => {
1370 stopped_before_candidate = true;
1371 break;
1372 }
1373 AcceptedStructuralMutationStagedAdmission::CandidateExceedsPolicy => {
1374 stopped_before_candidate = true;
1375 candidate_exceeds_batch_policy = true;
1376 break;
1377 }
1378 }
1379 let row_op = physical_changed.then(|| {
1380 CommitRowOp::new(
1381 entity_path,
1382 raw_key.clone(),
1383 canonical_before
1384 .as_ref()
1385 .map(|before| before.as_raw_row().as_bytes().to_vec()),
1386 Some(after.as_raw_row().as_bytes().to_vec()),
1387 catalog.fingerprint(),
1388 )
1389 });
1390 scheduler.schedule_save_after_image(
1391 mode,
1392 &data_key,
1393 after.as_raw_row(),
1394 provenance.as_slice(),
1395 row_op,
1396 batch_input_ordinal,
1397 )?;
1398 let values = if capture_output_values {
1399 let mut values = Vec::with_capacity(descriptor.fields().len());
1400 for field in descriptor.fields() {
1401 values.push(
1402 reader
1403 .required_cached_value(usize::from(field.slot().get()))?
1404 .clone(),
1405 );
1406 }
1407 values
1408 } else {
1409 Vec::new()
1410 };
1411 output.push(AcceptedStructuralMutationRow {
1412 values,
1413 logical_changed,
1414 });
1415 }
1416
1417 let report = AcceptedStructuralMutationPackingReport {
1418 admitted_mutations: output.len(),
1419 stopped_before_candidate,
1420 candidate_exceeds_batch_policy,
1421 };
1422 let batch = scheduler.finish();
1423 let (prepared, commit_directive) = precommit_preparation(output, report)?;
1424 let identity_ranges = identity_cursor
1425 .map(IdentityStatementCursor::into_range_advance)
1426 .transpose()?
1427 .into_iter()
1428 .flatten()
1429 .collect::<Vec<_>>();
1430 if !matches!(
1431 commit_directive,
1432 AcceptedStructuralMutationCommitDirective::Skip
1433 ) && batch.is_empty()
1434 && !identity_ranges.is_empty()
1435 {
1436 return Err(InternalError::identity_corruption());
1437 }
1438 match commit_directive {
1439 AcceptedStructuralMutationCommitDirective::Skip => {}
1440 AcceptedStructuralMutationCommitDirective::Standard if batch.is_empty() => {}
1441 AcceptedStructuralMutationCommitDirective::Standard => {
1442 commit_structural_row_ops_with_window_for_path(
1443 &self.db,
1444 entity_path,
1445 batch,
1446 identity_ranges,
1447 "accepted_structural_batch_apply",
1448 )?;
1449 }
1450 AcceptedStructuralMutationCommitDirective::WithMutationProgress(operation)
1451 if batch.is_empty() =>
1452 {
1453 let _ = operation;
1454 return Err(InternalError::executor_invariant());
1455 }
1456 AcceptedStructuralMutationCommitDirective::WithMutationProgress(operation) => {
1457 commit_structural_row_ops_with_mutation_progress_for_path(
1458 &self.db,
1459 entity_path,
1460 batch,
1461 identity_ranges,
1462 operation,
1463 "accepted_structural_batch_apply",
1464 )?;
1465 }
1466 }
1467 Ok(prepared)
1468 }
1469
1470 fn execute_one_accepted_save_mutation(
1471 &self,
1472 catalog: &AcceptedSchemaCatalogContext,
1473 descriptor: &AcceptedRowLayoutRuntimeContract<'_>,
1474 mode: MutationMode,
1475 target: AcceptedStructuralMutationTarget,
1476 patch: AcceptedMutationIntentPatch,
1477 ) -> Result<DynamicMutationResult, InternalError> {
1478 let identity = catalog.identity();
1479 let entity_path = identity.entity_path();
1480 let result = self.execute_accepted_structural_save_batch(
1481 catalog,
1482 descriptor,
1483 vec![AcceptedStructuralMutation::save(mode, target, patch)],
1484 Timestamp::now(),
1485 |rows| prepare_dynamic_mutation_result(catalog, descriptor, rows, false),
1486 )?;
1487 record(MetricsEvent::SaveMutation {
1488 entity_path: entity_path.into(),
1489 kind: match mode {
1490 MutationMode::Insert => SaveMutationKind::Insert,
1491 MutationMode::Replace => SaveMutationKind::Replace,
1492 MutationMode::Update => SaveMutationKind::Update,
1493 },
1494 rows_touched: u64::from(result.affected_rows),
1495 });
1496 Ok(result)
1497 }
1498
1499 pub fn execute_trusted_dynamic_mutation(
1506 &self,
1507 request: &DynamicMutation,
1508 ) -> Result<DynamicMutationResult, InternalError> {
1509 self.execute_trusted_dynamic_mutation_batch_with_result_policy(vec![request.clone()], false)
1510 }
1511
1512 pub fn execute_trusted_dynamic_mutation_batch(
1518 &self,
1519 requests: Vec<DynamicMutation>,
1520 ) -> Result<DynamicMutationResult, InternalError> {
1521 self.execute_trusted_dynamic_mutation_batch_with_result_policy(requests, true)
1522 }
1523
1524 fn execute_trusted_dynamic_mutation_batch_with_result_policy(
1525 &self,
1526 requests: Vec<DynamicMutation>,
1527 enforce_mixed_batch_result_bound: bool,
1528 ) -> Result<DynamicMutationResult, InternalError> {
1529 if requests.is_empty() {
1530 return Err(InternalError::mutation_batch_empty());
1531 }
1532 if requests.len() > MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS {
1533 return Err(InternalError::mutation_batch_too_many_items(
1534 requests.len(),
1535 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1536 ));
1537 }
1538 let first = requests
1539 .first()
1540 .ok_or_else(InternalError::mutation_batch_empty)?;
1541 if first.entity().is_empty() {
1542 return Err(InternalError::executor_unsupported());
1543 }
1544 let catalog = self.accepted_schema_catalog_context_for_entity_name(Some(first.entity()))?;
1545 let accepted_identity = catalog.identity();
1546 let descriptor =
1547 AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())?;
1548 let mut mutations = Vec::with_capacity(requests.len());
1549 let mut save_kinds = Vec::with_capacity(requests.len());
1550
1551 for (batch_position, request) in requests.iter().enumerate() {
1552 let batch_position = u32::try_from(batch_position).map_err(|_| {
1553 InternalError::mutation_batch_too_many_items(
1554 requests.len(),
1555 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
1556 )
1557 })?;
1558 if request.entity().is_empty() {
1559 return Err(InternalError::executor_unsupported());
1560 }
1561 let item_catalog =
1562 self.accepted_schema_catalog_context_for_entity_name(Some(request.entity()))?;
1563 if item_catalog.identity() != accepted_identity {
1564 return Err(InternalError::mutation_batch_entity_mismatch(
1565 batch_position,
1566 accepted_identity.entity_tag().value(),
1567 item_catalog.identity().entity_tag().value(),
1568 ));
1569 }
1570 let (mutation, save_kind) = lower_dynamic_mutation_intent(
1571 accepted_identity.entity_tag(),
1572 accepted_identity.entity_path(),
1573 &descriptor,
1574 request,
1575 batch_position,
1576 )?;
1577 mutations.push(mutation);
1578 save_kinds.push(save_kind);
1579 }
1580
1581 let entity_path = accepted_identity.entity_path_handle();
1582 let (result, metrics) = self.execute_accepted_structural_save_batch(
1583 &catalog,
1584 &descriptor,
1585 mutations,
1586 Timestamp::now(),
1587 |rows| {
1588 if rows.len() != save_kinds.len() {
1589 return Err(InternalError::executor_invariant());
1590 }
1591 let metrics = rows
1592 .iter()
1593 .zip(save_kinds)
1594 .filter_map(|(row, kind)| kind.map(|kind| (kind, row.logical_changed())))
1595 .collect::<Vec<_>>();
1596 let result = prepare_dynamic_mutation_result(
1597 &catalog,
1598 &descriptor,
1599 rows,
1600 enforce_mixed_batch_result_bound,
1601 )?;
1602 Ok((result, metrics))
1603 },
1604 )?;
1605 for (kind, logical_changed) in metrics {
1606 record(MetricsEvent::SaveMutation {
1607 entity_path: entity_path.clone(),
1608 kind,
1609 rows_touched: u64::from(logical_changed),
1610 });
1611 }
1612 Ok(result)
1613 }
1614
1615 #[doc(hidden)]
1618 pub fn execute_trusted_typed_mutation(
1619 &self,
1620 binding: &DynamicTypedEntityBinding,
1621 request: &DynamicTypedMutation,
1622 ) -> Result<Option<DynamicMutationResult>, InternalError> {
1623 let Some(catalog) = self.current_typed_entity_binding_catalog(binding)? else {
1624 return Ok(None);
1625 };
1626 let identity = catalog.identity();
1627 let descriptor =
1628 AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())?;
1629 let mode = dynamic_typed_mutation_mode(request);
1630 let (target, patch) = match request {
1631 DynamicTypedMutation::Insert { patch } => (
1632 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
1633 patch,
1634 ),
1635 DynamicTypedMutation::Update { key, patch }
1636 | DynamicTypedMutation::Replace { key, patch } => (
1637 AcceptedStructuralMutationTarget::expected(dynamic_key(
1638 identity.entity_tag(),
1639 key,
1640 )?),
1641 patch,
1642 ),
1643 };
1644 if !patch.is_bound_to(binding) {
1645 return Ok(None);
1646 }
1647 let patch = lower_typed_patch(
1648 &descriptor,
1649 patch,
1650 mode,
1651 mutation_diagnostic_context(identity.entity_tag(), mode, 0),
1652 )?;
1653 self.execute_one_accepted_save_mutation(&catalog, &descriptor, mode, target, patch)
1654 .map(Some)
1655 }
1656
1657 pub fn execute_trusted_dynamic_insert_batch(
1663 &self,
1664 entity: &str,
1665 patches: Vec<DynamicStructuralPatch>,
1666 ) -> Result<DynamicMutationResult, InternalError> {
1667 let mutations = patches
1668 .into_iter()
1669 .map(|patch| DynamicMutation::Insert {
1670 entity: entity.to_string(),
1671 patch,
1672 })
1673 .collect();
1674 self.execute_trusted_dynamic_mutation_batch_with_result_policy(mutations, false)
1675 }
1676}
1677
1678#[cfg(test)]
1679mod typed_adapter_tests {
1680 use super::{
1681 AcceptedFieldKind, DbSession, DynamicTypedBindingError, DynamicTypedFieldBindingRequest,
1682 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell, dynamic_typed_field_type,
1683 typed_adapter_field_kind_matches,
1684 };
1685 use crate::{
1686 db::{
1687 data::DataStore,
1688 index::IndexStore,
1689 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
1690 schema::{
1691 AcceptedSchemaRevision, FieldId, FieldStorageDecode, LeafCodec,
1692 PersistedFieldSnapshot, PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot,
1693 SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
1694 accepted_schema_candidate_with_field_bindings_for_tests,
1695 },
1696 },
1697 traits::{CanisterKind, Path},
1698 types::EntityTag,
1699 value::InputValue,
1700 };
1701 use icydb_schema::{EntitySourceKey, FieldSourceKey, ScalarType};
1702 use std::{cell::RefCell, collections::BTreeMap};
1703
1704 const STORE_PATH: &str = "session::write::typed_adapter_tests::Store";
1705 const ENTITY_SOURCE: &str = "session::write::typed_adapter_tests::Entity";
1706 const OTHER_ENTITY_SOURCE: &str = "session::write::typed_adapter_tests::OtherEntity";
1707 const ID_SOURCE: &str = "session::write::typed_adapter_tests::Entity::id";
1708 const VALUE_SOURCE: &str = "session::write::typed_adapter_tests::Entity::value";
1709 const REPLACEMENT_SOURCE: &str =
1710 "session::write::typed_adapter_tests::Entity::replacement_value";
1711 const OTHER_ID_SOURCE: &str = "session::write::typed_adapter_tests::OtherEntity::id";
1712
1713 struct TestCanister;
1714
1715 impl Path for TestCanister {
1716 const PATH: &'static str = "session::write::typed_adapter_tests::Canister";
1717 }
1718
1719 impl CanisterKind for TestCanister {
1720 const COMMIT_MEMORY_ID: u8 = 41;
1721 const COMMIT_STABLE_KEY: &'static str = "icydb.typed_adapter_tests.commit.v1";
1722 const STARTUP_MEMORY_ID: u8 = 49;
1723 const STARTUP_STABLE_KEY: &'static str = "icydb.typed_adapter_tests.startup.control.v1";
1724 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 42;
1725 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
1726 "icydb.typed_adapter_tests.integrity.progress.v1";
1727 }
1728
1729 thread_local! {
1730 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
1731 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
1732 static SCHEMA_STORE: RefCell<SchemaStore> =
1733 const { RefCell::new(SchemaStore::init_heap()) };
1734 static STORE_REGISTRY: StoreRegistry = {
1735 let mut registry = StoreRegistry::new();
1736 registry.register_store(
1737 STORE_PATH,
1738 &DATA_STORE,
1739 &INDEX_STORE,
1740 &SCHEMA_STORE,
1741 StoreAllocationIdentities::absent(),
1742 StoreRuntimeStorageCapabilities::heap(),
1743 ).expect("typed adapter test store should register");
1744 registry
1745 };
1746 }
1747
1748 fn nat64_field(id: u32, name: &str, slot: u16) -> PersistedFieldSnapshot {
1749 PersistedFieldSnapshot::new_initial(
1750 FieldId::new(id),
1751 name.to_string(),
1752 SchemaFieldSlot::new(slot),
1753 AcceptedFieldKind::Nat64,
1754 Vec::new(),
1755 false,
1756 SchemaInsertDefault::None,
1757 FieldStorageDecode::ByKind,
1758 LeafCodec::Scalar(ScalarCodec::Nat64),
1759 )
1760 }
1761
1762 fn snapshot(
1763 entity_source: &str,
1764 entity_name: &str,
1765 fields: Vec<PersistedFieldSnapshot>,
1766 ) -> PersistedSchemaSnapshot {
1767 let layout = SchemaRowLayout::initial(
1768 fields
1769 .iter()
1770 .map(|field| (field.id(), field.slot()))
1771 .collect(),
1772 );
1773 PersistedSchemaSnapshot::new(
1774 SchemaVersion::initial(),
1775 entity_source.to_string(),
1776 entity_name.to_string(),
1777 FieldId::new(1),
1778 layout,
1779 fields,
1780 )
1781 }
1782
1783 fn field_source(source: &str) -> FieldSourceKey {
1784 FieldSourceKey::try_new(source).expect("typed field source should admit")
1785 }
1786
1787 fn entity_source(source: &str) -> EntitySourceKey {
1788 EntitySourceKey::try_new(source).expect("typed entity source should admit")
1789 }
1790
1791 fn publish(
1792 session: &DbSession<TestCanister>,
1793 expected: AcceptedSchemaRevision,
1794 revision: AcceptedSchemaRevision,
1795 snapshots: BTreeMap<EntityTag, PersistedSchemaSnapshot>,
1796 fields: BTreeMap<(EntityTag, FieldSourceKey), FieldId>,
1797 ) {
1798 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
1799 STORE_PATH, revision, snapshots, fields,
1800 );
1801 let store = session
1802 .db
1803 .store_handle(STORE_PATH)
1804 .expect("typed adapter test store should resolve");
1805 crate::db::commit::publish_accepted_schema_candidate(
1806 STORE_PATH, store, expected, &candidate,
1807 )
1808 .expect("typed binding candidate should publish");
1809 }
1810
1811 fn request(source: &str) -> DynamicTypedFieldBindingRequest {
1812 DynamicTypedFieldBindingRequest::new(
1813 source.to_string(),
1814 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
1815 false,
1816 )
1817 }
1818
1819 fn assert_query_diagnostic(
1820 error: crate::db::QueryError,
1821 code: icydb_diagnostic_code::DiagnosticCode,
1822 origin: icydb_diagnostic_code::ErrorOrigin,
1823 detail: icydb_diagnostic_code::DiagnosticDetail,
1824 ) {
1825 let diagnostic = error.diagnostic();
1826 assert_eq!(diagnostic.code(), code);
1827 assert_eq!(diagnostic.origin(), origin);
1828 assert_eq!(diagnostic.detail(), Some(&detail));
1829 }
1830
1831 #[test]
1832 fn typed_adapter_kind_matching_is_exact_but_accepts_relation_key_wrappers() {
1833 let relation = AcceptedFieldKind::Relation {
1834 target_path: "test::Target".to_string(),
1835 target_entity_name: "Target".to_string(),
1836 target_entity_tag: EntityTag::new(7),
1837 target_store_path: "test::Store".to_string(),
1838 key_kind: Box::new(AcceptedFieldKind::Nat64),
1839 };
1840
1841 assert!(typed_adapter_field_kind_matches(
1842 &relation,
1843 &AcceptedFieldKind::Nat64,
1844 ));
1845 assert!(typed_adapter_field_kind_matches(
1846 &AcceptedFieldKind::List(Box::new(relation)),
1847 &AcceptedFieldKind::List(Box::new(AcceptedFieldKind::Nat64)),
1848 ));
1849 assert!(!typed_adapter_field_kind_matches(
1850 &AcceptedFieldKind::Nat64,
1851 &AcceptedFieldKind::Nat32,
1852 ));
1853 }
1854
1855 #[test]
1856 fn typed_adapter_field_contract_rejects_invalid_named_source_identity() {
1857 assert!(matches!(
1858 dynamic_typed_field_type(DynamicTypedFieldType::Named(String::new())),
1859 Err(DynamicTypedBindingError::FieldUnavailable),
1860 ));
1861 assert!(matches!(
1862 dynamic_typed_field_type(DynamicTypedFieldType::Scalar(ScalarType::Nat16)),
1863 Ok(icydb_schema::FieldType::Scalar(ScalarType::Nat16)),
1864 ));
1865 }
1866
1867 #[expect(clippy::too_many_lines)]
1870 #[test]
1871 fn typed_binding_uses_accepted_ids_and_slots_across_renames_and_name_reuse() {
1872 let entity_tag = EntityTag::new(91);
1873 let other_entity_tag = EntityTag::new(92);
1874 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
1875 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
1876 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
1877
1878 let session = DbSession::<TestCanister>::new(
1879 &STORE_REGISTRY,
1880 &crate::db::RequestExecutionRoot::__new_runtime_root(),
1881 );
1882 session
1883 .db
1884 .drive_startup_recovery_page()
1885 .expect("typed adapter test database should initialize");
1886 publish(
1887 &session,
1888 AcceptedSchemaRevision::NONE,
1889 AcceptedSchemaRevision::INITIAL,
1890 BTreeMap::from([(
1891 entity_tag,
1892 snapshot(
1893 ENTITY_SOURCE,
1894 "Entity",
1895 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
1896 ),
1897 )]),
1898 BTreeMap::from([
1899 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1900 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1901 ]),
1902 );
1903
1904 let initial_catalog = session
1905 .find_accepted_schema_catalog_context_for_entity_source_key(ENTITY_SOURCE)
1906 .expect("initial source catalog lookup should inspect")
1907 .expect("initial source catalog should exist");
1908 assert_eq!(initial_catalog.identity().entity_tag(), entity_tag);
1909 let initial = session
1910 .issue_typed_entity_binding(
1911 entity_source(ENTITY_SOURCE).as_str(),
1912 &[request(ID_SOURCE), request(VALUE_SOURCE)],
1913 )
1914 .expect("initial typed binding should issue");
1915 assert_eq!(initial.field_slot(ID_SOURCE), Some(0));
1916 assert_eq!(initial.field_slot(VALUE_SOURCE), Some(1));
1917 assert_eq!(initial.output_field_slot("value"), Some(1));
1918 let initial_patch = initial
1919 .bind_write_fields(vec![(
1920 VALUE_SOURCE.to_string(),
1921 DynamicWriteCell::Value(InputValue::Nat64(7)),
1922 )])
1923 .expect("source-bound patch should lower");
1924 assert_eq!(
1925 initial_patch.fields(),
1926 &[(2, 1, DynamicWriteCell::Value(InputValue::Nat64(7)))]
1927 );
1928
1929 publish(
1930 &session,
1931 AcceptedSchemaRevision::INITIAL,
1932 AcceptedSchemaRevision::new(2),
1933 BTreeMap::from([
1934 (
1935 entity_tag,
1936 snapshot(
1937 ENTITY_SOURCE,
1938 "RenamedEntity",
1939 vec![
1940 nat64_field(1, "id", 0),
1941 nat64_field(2, "renamed_value", 1),
1942 nat64_field(3, "value", 2),
1943 ],
1944 ),
1945 ),
1946 (
1947 other_entity_tag,
1948 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
1949 ),
1950 ]),
1951 BTreeMap::from([
1952 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
1953 ((entity_tag, field_source(VALUE_SOURCE)), FieldId::new(2)),
1954 (
1955 (entity_tag, field_source(REPLACEMENT_SOURCE)),
1956 FieldId::new(3),
1957 ),
1958 (
1959 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
1960 FieldId::new(1),
1961 ),
1962 ]),
1963 );
1964
1965 let stale_authority = session
1966 .ensure_accepted_schema_authority_is_current_for_store_path(
1967 STORE_PATH,
1968 initial_catalog.value_catalog_handle().authority(),
1969 )
1970 .expect_err("the initial accepted authority must be stale after revision two");
1971 assert_eq!(
1972 stale_authority.diagnostic_facts(),
1973 vec![
1974 (
1975 icydb_diagnostic_code::DiagnosticFactTag::ExpectedRevision,
1976 AcceptedSchemaRevision::INITIAL.get(),
1977 ),
1978 (
1979 icydb_diagnostic_code::DiagnosticFactTag::CurrentRevision,
1980 AcceptedSchemaRevision::new(2).get(),
1981 ),
1982 ],
1983 );
1984
1985 assert!(
1986 !session
1987 .typed_entity_binding_is_current(&initial)
1988 .expect("renamed binding currentness should inspect")
1989 );
1990 let renamed = session
1991 .issue_typed_entity_binding(ENTITY_SOURCE, &[request(ID_SOURCE), request(VALUE_SOURCE)])
1992 .expect("renamed source-bound adapter should rebind");
1993 assert_eq!(renamed.entity(), "RenamedEntity");
1994 assert_eq!(renamed.field_slot(VALUE_SOURCE), Some(1));
1995 assert_eq!(renamed.output_field_slot("renamed_value"), Some(1));
1996 assert_eq!(renamed.output_field_slot("value"), None);
1997
1998 publish(
1999 &session,
2000 AcceptedSchemaRevision::new(2),
2001 AcceptedSchemaRevision::new(3),
2002 BTreeMap::from([
2003 (
2004 entity_tag,
2005 snapshot(
2006 ENTITY_SOURCE,
2007 "RenamedEntity",
2008 vec![nat64_field(1, "id", 0), nat64_field(2, "value", 1)],
2009 ),
2010 ),
2011 (
2012 other_entity_tag,
2013 snapshot(OTHER_ENTITY_SOURCE, "Entity", vec![nat64_field(1, "id", 0)]),
2014 ),
2015 ]),
2016 BTreeMap::from([
2017 ((entity_tag, field_source(ID_SOURCE)), FieldId::new(1)),
2018 (
2019 (entity_tag, field_source(REPLACEMENT_SOURCE)),
2020 FieldId::new(2),
2021 ),
2022 (
2023 (other_entity_tag, field_source(OTHER_ID_SOURCE)),
2024 FieldId::new(1),
2025 ),
2026 ]),
2027 );
2028
2029 assert!(matches!(
2030 session.issue_typed_entity_binding(
2031 ENTITY_SOURCE,
2032 &[request(ID_SOURCE), request(VALUE_SOURCE)],
2033 ),
2034 Err(DynamicTypedBindingError::FieldUnavailable),
2035 ));
2036 assert!(
2037 !session
2038 .typed_entity_binding_is_current(&renamed)
2039 .expect("removed source binding should become stale")
2040 );
2041
2042 let replacement = session
2043 .issue_typed_entity_binding(
2044 ENTITY_SOURCE,
2045 &[request(ID_SOURCE), request(REPLACEMENT_SOURCE)],
2046 )
2047 .expect("explicit replacement source should bind");
2048 assert!(
2049 session
2050 .execute_trusted_typed_mutation(
2051 &replacement,
2052 &DynamicTypedMutation::Insert {
2053 patch: initial_patch
2054 },
2055 )
2056 .expect("cross-binding patch should fail closed")
2057 .is_none()
2058 );
2059 let patch = replacement
2060 .bind_write_fields(vec![
2061 (
2062 ID_SOURCE.to_string(),
2063 DynamicWriteCell::Value(InputValue::Nat64(1)),
2064 ),
2065 (
2066 REPLACEMENT_SOURCE.to_string(),
2067 DynamicWriteCell::Value(InputValue::Nat64(9)),
2068 ),
2069 ])
2070 .expect("replacement source write should bind by accepted IDs and slots");
2071 let result = session
2072 .execute_trusted_typed_mutation(&replacement, &DynamicTypedMutation::Insert { patch })
2073 .expect("typed insert should use the accepted mutation pipeline")
2074 .expect("replacement binding should remain current");
2075 assert_eq!(result.entity, "RenamedEntity");
2076 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
2077 assert_eq!(
2078 result.rows,
2079 vec![vec![
2080 crate::value::OutputValue::Nat64(1),
2081 crate::value::OutputValue::Nat64(9)
2082 ]]
2083 );
2084 assert_eq!(result.affected_rows, 1);
2085
2086 let second_patch = replacement
2087 .bind_write_fields(vec![
2088 (
2089 ID_SOURCE.to_string(),
2090 DynamicWriteCell::Value(InputValue::Nat64(2)),
2091 ),
2092 (
2093 REPLACEMENT_SOURCE.to_string(),
2094 DynamicWriteCell::Value(InputValue::Nat64(10)),
2095 ),
2096 ])
2097 .expect("second source-bound patch should lower");
2098 session
2099 .execute_trusted_typed_mutation(
2100 &replacement,
2101 &DynamicTypedMutation::Insert {
2102 patch: second_patch,
2103 },
2104 )
2105 .expect("second typed insert should use the accepted mutation pipeline")
2106 .expect("replacement binding should remain current");
2107
2108 {
2109 let query = crate::db::DynamicQuery::new("RenamedEntity")
2110 .select(["id", "value"])
2111 .order_by(crate::db::asc("id"))
2112 .limit(1);
2113 let result = session
2114 .execute_trusted_live_page(&query, None)
2115 .expect("SQL-free dynamic execution should use accepted authority");
2116 assert_eq!(result.entity, "RenamedEntity");
2117 assert_eq!(result.columns, vec!["id".to_string(), "value".to_string()]);
2118 assert_eq!(
2119 result.rows,
2120 vec![vec![
2121 crate::value::OutputValue::Nat64(1),
2122 crate::value::OutputValue::Nat64(9)
2123 ]]
2124 );
2125 assert_eq!(result.row_count, 1);
2126 assert_query_diagnostic(
2127 session
2128 .execute_trusted_live_page(&query.cursor("00"), None)
2129 .expect_err("scalar execution must reject grouped cursor state"),
2130 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
2131 icydb_diagnostic_code::ErrorOrigin::Query,
2132 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
2133 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
2134 },
2135 );
2136 assert_query_diagnostic(
2137 session
2138 .execute_public_dynamic_grouped_query(
2139 &crate::db::DynamicQuery::new("RenamedEntity").grouped_limits(1, 1024),
2140 )
2141 .expect_err("grouped execution must reject scalar query state"),
2142 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
2143 icydb_diagnostic_code::ErrorOrigin::Query,
2144 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
2145 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
2146 },
2147 );
2148
2149 let grouped_query = crate::db::DynamicQuery::new("RenamedEntity")
2150 .filter(crate::db::FieldRef::new("id").eq(1_u64))
2151 .group_by("value")
2152 .aggregate(crate::db::count())
2153 .grouped_limits(1, 16 * 1024)
2154 .limit(1);
2155 let grouped = session
2156 .execute_public_dynamic_grouped_query(&grouped_query)
2157 .expect("SQL-free grouped execution should use accepted authority");
2158 let typed_grouped = session
2159 .execute_public_dynamic_grouped_query_for_typed_binding(
2160 &replacement,
2161 &grouped_query,
2162 )
2163 .expect("typed grouped execution should inspect accepted authority")
2164 .expect("replacement binding should remain current");
2165 assert_eq!(typed_grouped, grouped);
2166 assert!(
2167 session
2168 .execute_public_dynamic_grouped_query_for_typed_binding(
2169 &renamed,
2170 &grouped_query,
2171 )
2172 .expect("stale grouped binding should inspect accepted authority")
2173 .is_none(),
2174 "stale typed grouped bindings must fail closed before execution"
2175 );
2176 assert_eq!(grouped.entity, "RenamedEntity");
2177 assert_eq!(grouped.row_count, 1);
2178 assert_eq!(grouped.rows.len(), 1);
2179 assert_eq!(
2180 grouped.rows[0].group_key(),
2181 &[crate::value::OutputValue::Nat64(9)]
2182 );
2183 assert_eq!(
2184 grouped.rows[0].aggregate_values(),
2185 &[crate::value::OutputValue::Nat64(1)]
2186 );
2187 assert_eq!(grouped.next_cursor, None);
2188
2189 let grouped_state_error = session
2190 .execute_trusted_dynamic_grouped_query(&grouped_query.clone().grouped_limits(1, 1))
2191 .expect_err("grouped retained state must respect its explicit byte ceiling");
2192 assert!(matches!(
2193 grouped_state_error.diagnostic().detail(),
2194 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
2195 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
2196 })
2197 ));
2198 assert_eq!(
2199 grouped_state_error.diagnostic_facts()[0],
2200 (
2201 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
2202 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::GroupDistinctStateBytes.raw(),
2203 ),
2204 );
2205
2206 assert_query_diagnostic(
2207 session
2208 .execute_public_dynamic_grouped_query(&grouped_query.clone().select(["value"]))
2209 .expect_err("grouped output must reject scalar selection"),
2210 icydb_diagnostic_code::DiagnosticCode::QueryIntent,
2211 icydb_diagnostic_code::ErrorOrigin::Query,
2212 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
2213 kind: icydb_diagnostic_code::QueryErrorKind::Intent,
2214 },
2215 );
2216 assert_query_diagnostic(
2217 session
2218 .execute_public_dynamic_grouped_query(
2219 &crate::db::DynamicQuery::new("RenamedEntity")
2220 .group_by("value")
2221 .aggregate(crate::db::count()),
2222 )
2223 .expect_err("public grouped execution must require explicit limits"),
2224 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
2225 icydb_diagnostic_code::ErrorOrigin::Query,
2226 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
2227 reason:
2228 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
2229 },
2230 );
2231 assert_query_diagnostic(
2232 session
2233 .execute_trusted_dynamic_grouped_query(
2234 &crate::db::DynamicQuery::new("RenamedEntity")
2235 .group_by("value")
2236 .aggregate(crate::db::count())
2237 .grouped_limits(0, 1024),
2238 )
2239 .expect_err("trusted grouped execution must reject zero limits"),
2240 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
2241 icydb_diagnostic_code::ErrorOrigin::Query,
2242 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
2243 reason:
2244 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryRequiresLimits,
2245 },
2246 );
2247 assert_query_diagnostic(
2248 session
2249 .execute_public_dynamic_grouped_query(&grouped_query.grouped_limits(101, 1024))
2250 .expect_err("public grouped execution must enforce its group budget"),
2251 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
2252 icydb_diagnostic_code::ErrorOrigin::Query,
2253 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
2254 reason:
2255 icydb_diagnostic_code::QueryReadAdmissionCode::GroupedQueryExceedsBudget,
2256 },
2257 );
2258
2259 let paged_query = crate::db::DynamicQuery::new("RenamedEntity")
2260 .group_by("value")
2261 .aggregate(crate::db::count())
2262 .grouped_limits(2, 16 * 1024)
2263 .limit(1);
2264 assert_query_diagnostic(
2265 session
2266 .execute_public_dynamic_grouped_query(&paged_query)
2267 .expect_err("public grouped execution must reject an unbounded full scan"),
2268 icydb_diagnostic_code::DiagnosticCode::QueryReadAdmission,
2269 icydb_diagnostic_code::ErrorOrigin::Query,
2270 icydb_diagnostic_code::DiagnosticDetail::QueryReadAdmission {
2271 reason:
2272 icydb_diagnostic_code::QueryReadAdmissionCode::UnboundedFullScanRejected,
2273 },
2274 );
2275 let first_page = session
2276 .execute_trusted_dynamic_grouped_query(&paged_query)
2277 .expect("SQL-free grouped first page should execute");
2278 assert_eq!(first_page.row_count, 1);
2279 assert_eq!(
2280 first_page.rows[0].group_key(),
2281 &[crate::value::OutputValue::Nat64(9)]
2282 );
2283 let cursor = first_page
2284 .next_cursor
2285 .expect("first grouped page should return a continuation cursor");
2286 assert_query_diagnostic(
2287 session
2288 .execute_trusted_dynamic_grouped_query(
2289 &paged_query.clone().cursor(format!("{cursor}0")),
2290 )
2291 .expect_err("tampered grouped cursor must fail closed"),
2292 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
2293 icydb_diagnostic_code::ErrorOrigin::Cursor,
2294 icydb_diagnostic_code::DiagnosticDetail::QueryKind {
2295 kind: icydb_diagnostic_code::QueryErrorKind::InvalidContinuationCursor,
2296 },
2297 );
2298 let second_page = session
2299 .execute_trusted_dynamic_grouped_query(&paged_query.cursor(cursor))
2300 .expect("SQL-free grouped continuation should execute");
2301 assert_eq!(second_page.row_count, 1);
2302 assert_eq!(
2303 second_page.rows[0].group_key(),
2304 &[crate::value::OutputValue::Nat64(10)]
2305 );
2306 assert_eq!(second_page.next_cursor, None);
2307 }
2308 }
2309}
2310
2311#[cfg(test)]
2312mod mixed_relation_batch_tests {
2313 use super::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell};
2314 use crate::{
2315 db::{
2316 DynamicQuery, asc,
2317 data::DataStore,
2318 desc,
2319 index::IndexStore,
2320 query::expr::FilterExpr,
2321 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
2322 schema::{
2323 AcceptedConstraintCatalog, AcceptedFieldKind, AcceptedSchemaRevision, FieldId,
2324 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2325 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2326 PersistedRelationEdgeSnapshot, PersistedSchemaSnapshot, RelationId, ScalarCodec,
2327 SchemaFieldSlot, SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore,
2328 SchemaVersion, accepted_schema_candidate_with_field_bindings_for_tests,
2329 },
2330 },
2331 error::ErrorClass,
2332 traits::{CanisterKind, Path},
2333 types::EntityTag,
2334 value::{InputValue, OutputValue},
2335 };
2336 use icydb_schema::FieldSourceKey;
2337 use std::{cell::RefCell, collections::BTreeMap};
2338
2339 const STORE_PATH: &str = "session::write::mixed_relation_batch_tests::Store";
2340 const ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node";
2341 const ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::id";
2342 const PARENT_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::parent_id";
2343 const CODE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Node::code";
2344 const ENTITY_NAME: &str = "MixedRelationNode";
2345 const ENTITY_TAG: EntityTag = EntityTag::new(94);
2346 const OTHER_ENTITY_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other";
2347 const OTHER_ID_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::id";
2348 const OTHER_VALUE_SOURCE: &str = "session::write::mixed_relation_batch_tests::Other::value";
2349 const OTHER_ENTITY_NAME: &str = "MixedRelationOther";
2350 const OTHER_ENTITY_TAG: EntityTag = EntityTag::new(95);
2351
2352 struct TestCanister;
2353
2354 impl Path for TestCanister {
2355 const PATH: &'static str = "session::write::mixed_relation_batch_tests::Canister";
2356 }
2357
2358 impl CanisterKind for TestCanister {
2359 const COMMIT_MEMORY_ID: u8 = 47;
2360 const COMMIT_STABLE_KEY: &'static str = "icydb.mixed_relation_batch_tests.commit.v1";
2361 const STARTUP_MEMORY_ID: u8 = 50;
2362 const STARTUP_STABLE_KEY: &'static str =
2363 "icydb.mixed_relation_batch_tests.startup.control.v1";
2364 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 48;
2365 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2366 "icydb.mixed_relation_batch_tests.integrity.progress.v1";
2367 }
2368
2369 thread_local! {
2370 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2371 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2372 static SCHEMA_STORE: RefCell<SchemaStore> =
2373 const { RefCell::new(SchemaStore::init_heap()) };
2374 static STORE_REGISTRY: StoreRegistry = {
2375 let mut registry = StoreRegistry::new();
2376 registry.register_store(
2377 STORE_PATH,
2378 &DATA_STORE,
2379 &INDEX_STORE,
2380 &SCHEMA_STORE,
2381 StoreAllocationIdentities::absent(),
2382 StoreRuntimeStorageCapabilities::heap(),
2383 ).expect("mixed relation test store should register");
2384 registry
2385 };
2386 }
2387
2388 fn source_key(source: &str) -> FieldSourceKey {
2389 FieldSourceKey::try_new(source).expect("mixed relation field source should admit")
2390 }
2391
2392 fn relation_snapshot() -> PersistedSchemaSnapshot {
2393 let fields = vec![
2394 PersistedFieldSnapshot::new_initial(
2395 FieldId::new(1),
2396 "id".to_string(),
2397 SchemaFieldSlot::new(0),
2398 AcceptedFieldKind::Nat64,
2399 Vec::new(),
2400 false,
2401 SchemaInsertDefault::None,
2402 FieldStorageDecode::ByKind,
2403 LeafCodec::Scalar(ScalarCodec::Nat64),
2404 ),
2405 PersistedFieldSnapshot::new_initial(
2406 FieldId::new(2),
2407 "parent_id".to_string(),
2408 SchemaFieldSlot::new(1),
2409 AcceptedFieldKind::Nat64,
2410 Vec::new(),
2411 true,
2412 SchemaInsertDefault::None,
2413 FieldStorageDecode::ByKind,
2414 LeafCodec::Scalar(ScalarCodec::Nat64),
2415 ),
2416 PersistedFieldSnapshot::new_initial(
2417 FieldId::new(3),
2418 "code".to_string(),
2419 SchemaFieldSlot::new(2),
2420 AcceptedFieldKind::Nat64,
2421 Vec::new(),
2422 false,
2423 SchemaInsertDefault::None,
2424 FieldStorageDecode::ByKind,
2425 LeafCodec::Scalar(ScalarCodec::Nat64),
2426 ),
2427 ];
2428 let relation = PersistedRelationEdgeSnapshot::new(
2429 RelationId::new(1).expect("mixed relation identity should be non-zero"),
2430 "parent".to_string(),
2431 ENTITY_SOURCE.to_string(),
2432 vec![FieldId::new(2)],
2433 );
2434 let snapshot = PersistedSchemaSnapshot::new_with_indexes(
2435 SchemaVersion::initial(),
2436 ENTITY_SOURCE.to_string(),
2437 ENTITY_NAME.to_string(),
2438 FieldId::new(1),
2439 SchemaRowLayout::initial(
2440 fields
2441 .iter()
2442 .map(|field| (field.id(), field.slot()))
2443 .collect(),
2444 ),
2445 fields,
2446 vec![PersistedIndexSnapshot::new(
2447 SchemaIndexId::new(1).expect("mixed unique index identity should be non-zero"),
2448 1,
2449 "by_code".to_string(),
2450 STORE_PATH.to_string(),
2451 true,
2452 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2453 FieldId::new(3),
2454 SchemaFieldSlot::new(2),
2455 vec!["code".to_string()],
2456 AcceptedFieldKind::Nat64,
2457 false,
2458 )]),
2459 None,
2460 )],
2461 )
2462 .with_relations(vec![relation]);
2463 let constraints = AcceptedConstraintCatalog::initial(
2464 snapshot.fields(),
2465 snapshot.indexes(),
2466 snapshot.relations(),
2467 )
2468 .expect("mixed relation constraints should close");
2469 snapshot.with_constraint_catalog(constraints)
2470 }
2471
2472 fn other_snapshot() -> PersistedSchemaSnapshot {
2473 let fields = vec![
2474 PersistedFieldSnapshot::new_initial(
2475 FieldId::new(1),
2476 "id".to_string(),
2477 SchemaFieldSlot::new(0),
2478 AcceptedFieldKind::Nat64,
2479 Vec::new(),
2480 false,
2481 SchemaInsertDefault::None,
2482 FieldStorageDecode::ByKind,
2483 LeafCodec::Scalar(ScalarCodec::Nat64),
2484 ),
2485 PersistedFieldSnapshot::new_initial(
2486 FieldId::new(2),
2487 "value".to_string(),
2488 SchemaFieldSlot::new(1),
2489 AcceptedFieldKind::Nat64,
2490 Vec::new(),
2491 false,
2492 SchemaInsertDefault::None,
2493 FieldStorageDecode::ByKind,
2494 LeafCodec::Scalar(ScalarCodec::Nat64),
2495 ),
2496 ];
2497 PersistedSchemaSnapshot::new(
2498 SchemaVersion::initial(),
2499 OTHER_ENTITY_SOURCE.to_string(),
2500 OTHER_ENTITY_NAME.to_string(),
2501 FieldId::new(1),
2502 SchemaRowLayout::initial(
2503 fields
2504 .iter()
2505 .map(|field| (field.id(), field.slot()))
2506 .collect(),
2507 ),
2508 fields,
2509 )
2510 }
2511
2512 fn initialize() -> DbSession<TestCanister> {
2513 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2514 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2515 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2516 let session = DbSession::<TestCanister>::new(
2517 &STORE_REGISTRY,
2518 &crate::db::RequestExecutionRoot::__new_runtime_root(),
2519 );
2520 session
2521 .db
2522 .drive_startup_recovery_page()
2523 .expect("mixed relation database should initialize");
2524 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2525 STORE_PATH,
2526 AcceptedSchemaRevision::INITIAL,
2527 BTreeMap::from([
2528 (ENTITY_TAG, relation_snapshot()),
2529 (OTHER_ENTITY_TAG, other_snapshot()),
2530 ]),
2531 BTreeMap::from([
2532 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2533 ((ENTITY_TAG, source_key(PARENT_SOURCE)), FieldId::new(2)),
2534 ((ENTITY_TAG, source_key(CODE_SOURCE)), FieldId::new(3)),
2535 (
2536 (OTHER_ENTITY_TAG, source_key(OTHER_ID_SOURCE)),
2537 FieldId::new(1),
2538 ),
2539 (
2540 (OTHER_ENTITY_TAG, source_key(OTHER_VALUE_SOURCE)),
2541 FieldId::new(2),
2542 ),
2543 ]),
2544 );
2545 let store = session
2546 .db
2547 .store_handle(STORE_PATH)
2548 .expect("mixed relation store should resolve");
2549 crate::db::commit::publish_accepted_schema_candidate(
2550 STORE_PATH,
2551 store,
2552 AcceptedSchemaRevision::NONE,
2553 &candidate,
2554 )
2555 .expect("mixed relation candidate should publish");
2556 session
2557 }
2558
2559 fn patch(id: Option<u64>, parent: Option<u64>, code: Option<u64>) -> DynamicStructuralPatch {
2560 let mut fields = Vec::new();
2561 if let Some(id) = id {
2562 fields.push((
2563 "id".to_string(),
2564 DynamicWriteCell::Value(InputValue::Nat64(id)),
2565 ));
2566 }
2567 fields.push((
2568 "parent_id".to_string(),
2569 parent.map_or(DynamicWriteCell::Null, |parent| {
2570 DynamicWriteCell::Value(InputValue::Nat64(parent))
2571 }),
2572 ));
2573 if let Some(code) = code {
2574 fields.push((
2575 "code".to_string(),
2576 DynamicWriteCell::Value(InputValue::Nat64(code)),
2577 ));
2578 }
2579 DynamicStructuralPatch::new(fields)
2580 }
2581
2582 fn insert(id: u64, parent: Option<u64>) -> DynamicMutation {
2583 insert_with_code(id, parent, id)
2584 }
2585
2586 fn insert_with_code(id: u64, parent: Option<u64>, code: u64) -> DynamicMutation {
2587 DynamicMutation::Insert {
2588 entity: ENTITY_NAME.to_string(),
2589 patch: patch(Some(id), parent, Some(code)),
2590 }
2591 }
2592
2593 fn update_parent(id: u64, parent: Option<u64>) -> DynamicMutation {
2594 DynamicMutation::Update {
2595 entity: ENTITY_NAME.to_string(),
2596 key: InputValue::Nat64(id),
2597 patch: patch(None, parent, None),
2598 }
2599 }
2600
2601 fn update_code(id: u64, code: u64) -> DynamicMutation {
2602 DynamicMutation::Update {
2603 entity: ENTITY_NAME.to_string(),
2604 key: InputValue::Nat64(id),
2605 patch: DynamicStructuralPatch::new(vec![(
2606 "code".to_string(),
2607 DynamicWriteCell::Value(InputValue::Nat64(code)),
2608 )]),
2609 }
2610 }
2611
2612 fn delete(id: u64) -> DynamicMutation {
2613 DynamicMutation::Delete {
2614 entity: ENTITY_NAME.to_string(),
2615 key: InputValue::Nat64(id),
2616 }
2617 }
2618
2619 fn expected_row(id: u64, parent: Option<u64>) -> Vec<OutputValue> {
2620 expected_row_with_code(id, parent, id)
2621 }
2622
2623 fn expected_row_with_code(id: u64, parent: Option<u64>, code: u64) -> Vec<OutputValue> {
2624 vec![
2625 OutputValue::Nat64(id),
2626 parent.map_or(OutputValue::Null, OutputValue::Nat64),
2627 OutputValue::Nat64(code),
2628 ]
2629 }
2630
2631 fn other_patch(id: Option<u64>, value: u64) -> DynamicStructuralPatch {
2632 let mut fields = Vec::new();
2633 if let Some(id) = id {
2634 fields.push((
2635 "id".to_string(),
2636 DynamicWriteCell::Value(InputValue::Nat64(id)),
2637 ));
2638 }
2639 fields.push((
2640 "value".to_string(),
2641 DynamicWriteCell::Value(InputValue::Nat64(value)),
2642 ));
2643 DynamicStructuralPatch::new(fields)
2644 }
2645
2646 fn assert_relation_violation(error: &crate::error::InternalError) {
2647 assert!(error.diagnostic_facts().contains(&(
2648 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
2649 icydb_diagnostic_code::DiagnosticConstraintKind::Relation.raw(),
2650 )));
2651 }
2652
2653 #[test]
2654 fn live_pages_resume_mixed_projection_from_authenticated_hidden_order_values() {
2655 let session = initialize();
2656 session
2657 .execute_trusted_dynamic_mutation_batch(vec![
2658 insert_with_code(1, None, 10),
2659 insert_with_code(2, Some(1), 20),
2660 insert_with_code(3, None, 30),
2661 ])
2662 .expect("live-page rows should insert");
2663 let query = DynamicQuery::new(ENTITY_NAME)
2664 .select(["id"])
2665 .order_by(desc("code"));
2666
2667 let first = session
2668 .execute_public_live_page(&query, None)
2669 .expect("initial live page should execute");
2670 assert_eq!(
2671 first.rows,
2672 vec![vec![OutputValue::Nat64(3)], vec![OutputValue::Nat64(2)]]
2673 );
2674 let cursor = first
2675 .continuation
2676 .as_deref()
2677 .expect("unreturned matching row should produce continuation");
2678 let second = session
2679 .execute_public_live_page(&query, Some(cursor))
2680 .expect("authenticated live continuation should resume");
2681 assert_eq!(second.rows, vec![vec![OutputValue::Nat64(1)]]);
2682 assert_eq!(second.continuation, None);
2683
2684 let total_limit = session
2685 .execute_public_live_page(&query.clone().limit(2), None)
2686 .expect("total live-page limit should execute");
2687 assert_eq!(
2688 total_limit.rows,
2689 vec![vec![OutputValue::Nat64(3)], vec![OutputValue::Nat64(2)]],
2690 );
2691 assert_eq!(
2692 total_limit.continuation, None,
2693 "query LIMIT is a total traversal window rather than a page size",
2694 );
2695
2696 let three_row_window = query.clone().limit(3);
2697 let limited_first = session
2698 .execute_public_live_page(&three_row_window, None)
2699 .expect("first total-window page should execute");
2700 let limited_cursor = limited_first
2701 .continuation
2702 .as_deref()
2703 .expect("a partially consumed total window should continue");
2704 let limited_second = session
2705 .execute_public_live_page(&three_row_window, Some(limited_cursor))
2706 .expect("remaining total window should preserve the plan signature");
2707 assert_eq!(limited_second.rows, vec![vec![OutputValue::Nat64(1)]]);
2708 assert_eq!(limited_second.continuation, None);
2709
2710 let mixed_order = DynamicQuery::new(ENTITY_NAME)
2711 .select(["id"])
2712 .order_by(desc("parent_id"))
2713 .order_by(asc("id"));
2714 let mixed_first = session
2715 .execute_trusted_live_page(&mixed_order, None)
2716 .expect("mixed-direction nullable order should execute");
2717 assert_eq!(
2718 mixed_first.rows,
2719 vec![vec![OutputValue::Nat64(2)], vec![OutputValue::Nat64(1)]],
2720 );
2721 let mixed_cursor = mixed_first
2722 .continuation
2723 .as_deref()
2724 .expect("duplicate null order values should retain continuation");
2725 let mixed_second = session
2726 .execute_trusted_live_page(&mixed_order, Some(mixed_cursor))
2727 .expect("mixed-direction nullable order should resume");
2728 assert_eq!(mixed_second.rows, vec![vec![OutputValue::Nat64(3)]]);
2729 assert_eq!(mixed_second.continuation, None);
2730
2731 let mismatched_window = session
2732 .execute_public_live_page(&query.clone().limit(3), Some(cursor))
2733 .expect_err("a changed total limit must invalidate the continuation");
2734 assert_eq!(
2735 mismatched_window.diagnostic_code(),
2736 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
2737 );
2738
2739 let mut tampered = cursor.as_bytes().to_vec();
2740 let last = tampered.len().saturating_sub(1);
2741 tampered[last] = if tampered[last] == b'0' { b'1' } else { b'0' };
2742 let tampered = String::from_utf8(tampered).expect("hex cursor should remain UTF-8");
2743 let error = session
2744 .execute_public_live_page(&query, Some(tampered.as_str()))
2745 .expect_err("tampered cursor must fail closed");
2746 assert_eq!(
2747 error.diagnostic_code(),
2748 icydb_diagnostic_code::DiagnosticCode::QueryInvalidContinuationCursor,
2749 );
2750 }
2751
2752 #[test]
2753 fn live_pages_resume_across_changed_output_work_envelopes() {
2754 let session = initialize();
2755 session
2756 .execute_trusted_dynamic_mutation_batch(vec![
2757 insert(1, None),
2758 insert(2, None),
2759 insert(3, None),
2760 ])
2761 .expect("output-envelope rows should insert");
2762 let query = DynamicQuery::new(ENTITY_NAME)
2763 .select(["id"])
2764 .order_by(desc("code"));
2765 let first = session
2766 .execute_trusted_live_page_with_result_bytes_limit_for_tests(&query, None, 32)
2767 .expect("small output envelope should publish the first bounded page");
2768 assert_eq!(first.rows, vec![vec![OutputValue::Nat64(3)]]);
2769 let continuation = first
2770 .continuation
2771 .expect("small output envelope should leave authenticated progress");
2772
2773 let second = session
2774 .execute_trusted_live_page_with_result_bytes_limit_for_tests(
2775 &query,
2776 Some(continuation.as_str()),
2777 64,
2778 )
2779 .unwrap_or_else(|error| {
2780 panic!(
2781 "larger output envelope should resume the same query: {error:?}, facts={:?}",
2782 error.diagnostic_facts(),
2783 )
2784 });
2785 assert_eq!(
2786 second.rows,
2787 vec![vec![OutputValue::Nat64(2)], vec![OutputValue::Nat64(1)]]
2788 );
2789 let second_continuation = second
2790 .continuation
2791 .as_deref()
2792 .expect("an exact-full page still needs to prove physical exhaustion");
2793 assert_ne!(first.work.envelope_identity, second.work.envelope_identity);
2794
2795 let terminal = session
2796 .execute_trusted_live_page_with_result_bytes_limit_for_tests(
2797 &query,
2798 Some(second_continuation),
2799 48,
2800 )
2801 .expect("a third finite envelope should prove exhaustion without replaying rows");
2802 assert!(terminal.rows.is_empty());
2803 assert_eq!(terminal.continuation, None);
2804 assert_ne!(
2805 second.work.envelope_identity,
2806 terminal.work.envelope_identity
2807 );
2808
2809 assert_eq!(
2810 [first.rows, second.rows, terminal.rows].concat(),
2811 vec![
2812 vec![OutputValue::Nat64(3)],
2813 vec![OutputValue::Nat64(2)],
2814 vec![OutputValue::Nat64(1)],
2815 ]
2816 );
2817 }
2818
2819 #[test]
2820 fn distinct_live_pages_resume_adjacent_groups_and_global_replay_end_to_end() {
2821 let session = initialize();
2822 session
2823 .execute_trusted_dynamic_mutation_batch(vec![
2824 insert(1, None),
2825 insert(2, None),
2826 insert(3, Some(1)),
2827 insert(4, Some(2)),
2828 insert(5, Some(1)),
2829 insert(6, Some(3)),
2830 insert(7, Some(2)),
2831 ])
2832 .expect("DISTINCT continuation rows should insert atomically");
2833
2834 let adjacent = DynamicQuery::new(ENTITY_NAME)
2835 .select(["parent_id"])
2836 .order_by(asc("parent_id"))
2837 .order_by(asc("id"))
2838 .distinct_for_internal_execution();
2839 let global = DynamicQuery::new(ENTITY_NAME)
2840 .select(["parent_id"])
2841 .order_by(asc("id"))
2842 .distinct_for_internal_execution();
2843
2844 let traverse = |query: &DynamicQuery, strategy: &str| {
2845 let mut continuation = None;
2846 let mut rows = Vec::new();
2847 let mut cursors = std::collections::BTreeSet::new();
2848 let mut pages = 0_u32;
2849 let mut entries_visited = 0_u64;
2850 loop {
2851 let page = session
2852 .execute_trusted_live_page(query, continuation.as_deref())
2853 .unwrap_or_else(|error| {
2854 panic!("{strategy} DISTINCT page should execute: {error:?}")
2855 });
2856 pages = pages.saturating_add(1);
2857 entries_visited = entries_visited.saturating_add(page.work.entries_visited);
2858 assert_eq!(page.row_count as usize, page.rows.len());
2859 assert_eq!(page.work.result_rows, page.row_count);
2860 rows.extend(page.rows);
2861 let Some(cursor) = page.continuation else {
2862 break;
2863 };
2864 assert!(
2865 cursors.insert(cursor.clone()),
2866 "{strategy} DISTINCT continuation must advance monotonically",
2867 );
2868 continuation = Some(cursor);
2869 assert!(pages < 8, "{strategy} DISTINCT traversal must terminate");
2870 }
2871
2872 (rows, pages, entries_visited)
2873 };
2874
2875 let expected = vec![
2876 vec![OutputValue::Null],
2877 vec![OutputValue::Nat64(1)],
2878 vec![OutputValue::Nat64(2)],
2879 vec![OutputValue::Nat64(3)],
2880 ];
2881 let (adjacent_rows, adjacent_pages, adjacent_entries) = traverse(&adjacent, "adjacent");
2882 let (global_rows, global_pages, global_entries) = traverse(&global, "global");
2883
2884 assert_eq!(adjacent_rows, expected);
2885 assert_eq!(global_rows, expected);
2886 assert_eq!(adjacent_pages, 2);
2887 assert_eq!(global_pages, 2);
2888 assert!(adjacent_entries > 0);
2889 assert!(global_entries > 0);
2890 }
2891
2892 #[test]
2893 fn selective_live_pages_publish_monotonic_empty_physical_progress() {
2894 let session = initialize();
2895 session
2896 .execute_trusted_dynamic_mutation_batch(
2897 (1..=9)
2898 .map(|id| {
2899 let parent = match id {
2900 1 => Some(2),
2901 9 => Some(1),
2902 _ => None,
2903 };
2904 insert(id, parent)
2905 })
2906 .collect(),
2907 )
2908 .expect("selective live-page rows should insert");
2909 let query = DynamicQuery::new(ENTITY_NAME)
2910 .select(["id"])
2911 .filter(FilterExpr::eq("parent_id", 1_u64))
2912 .order_by(asc("id"))
2913 .limit(1);
2914
2915 let first = session
2916 .execute_trusted_live_page(&query, None)
2917 .expect("first selective page should stop with physical progress");
2918 assert!(first.rows.is_empty());
2919 assert_eq!(first.work.entries_visited, 4);
2920 let first_cursor = first
2921 .continuation
2922 .expect("filtered physical progress must return a continuation");
2923
2924 let second = session
2925 .execute_trusted_live_page(&query, Some(first_cursor.as_str()))
2926 .expect("second selective page should resume after the first physical frontier");
2927 assert!(second.rows.is_empty());
2928 assert_eq!(second.work.entries_visited, 4);
2929 let second_cursor = second
2930 .continuation
2931 .expect("second filtered frontier must remain resumable");
2932 assert_ne!(second_cursor, first_cursor);
2933
2934 let third = session
2935 .execute_trusted_live_page(&query, Some(second_cursor.as_str()))
2936 .expect("final selective page should return the late match");
2937 assert_eq!(third.rows, vec![vec![OutputValue::Nat64(9)]]);
2938 assert_eq!(third.work.entries_visited, 1);
2939 assert_eq!(third.continuation, None);
2940
2941 let descending = DynamicQuery::new(ENTITY_NAME)
2942 .select(["id"])
2943 .filter(FilterExpr::eq("parent_id", 2_u64))
2944 .order_by(desc("id"))
2945 .limit(1);
2946 let descending_first = session
2947 .execute_trusted_live_page(&descending, None)
2948 .expect("descending selective page should stop with physical progress");
2949 assert!(descending_first.rows.is_empty());
2950 let descending_first_cursor = descending_first
2951 .continuation
2952 .expect("descending filtered progress must return a continuation");
2953 let descending_second = session
2954 .execute_trusted_live_page(&descending, Some(descending_first_cursor.as_str()))
2955 .expect("descending progress should resume after its physical frontier");
2956 assert!(descending_second.rows.is_empty());
2957 let descending_second_cursor = descending_second
2958 .continuation
2959 .expect("descending second frontier must remain resumable");
2960 assert_ne!(descending_second_cursor, descending_first_cursor);
2961 let descending_third = session
2962 .execute_trusted_live_page(&descending, Some(descending_second_cursor.as_str()))
2963 .expect("descending final page should return the late match");
2964 assert_eq!(descending_third.rows, vec![vec![OutputValue::Nat64(1)]]);
2965 assert_eq!(descending_third.continuation, None);
2966 }
2967
2968 #[test]
2969 fn accepted_relation_edges_drive_catalog_and_describe_introspection() {
2970 let session = initialize();
2971 let entities = session
2972 .show_entities()
2973 .expect("accepted entity catalog should resolve");
2974 let source = entities
2975 .iter()
2976 .find(|entity| entity.entity_name() == ENTITY_NAME)
2977 .expect("relation source should be listed");
2978 assert_eq!(source.relations(), 1);
2979
2980 let description = session
2981 .try_describe_entity_by_name(ENTITY_NAME)
2982 .expect("accepted relation source should describe");
2983 let [relation] = description.relations() else {
2984 panic!("accepted relation edge should produce one relation row");
2985 };
2986 assert_eq!(relation.field(), "parent_id");
2987 assert_eq!(relation.target_path(), ENTITY_SOURCE);
2988 assert_eq!(relation.target_entity_name(), ENTITY_NAME);
2989 assert_eq!(relation.target_store_path(), STORE_PATH);
2990 assert_eq!(
2991 relation.cardinality(),
2992 crate::db::EntityRelationCardinality::Single,
2993 );
2994 }
2995
2996 #[test]
2997 fn mixed_relation_validation_uses_the_complete_final_row_overlay() {
2998 let session = initialize();
2999 session
3000 .execute_trusted_dynamic_mutation_batch(vec![insert(1, None), insert(2, Some(1))])
3001 .expect("the initial relation should commit");
3002
3003 let blocked = session
3004 .execute_trusted_dynamic_mutation(&delete(1))
3005 .expect_err("an unaffected committed source must block target deletion");
3006 assert_relation_violation(&blocked);
3007
3008 let deleted = session
3009 .execute_trusted_dynamic_mutation_batch(vec![delete(2), delete(1)])
3010 .expect("a source and its target should delete atomically");
3011 assert_eq!(
3012 deleted.rows,
3013 vec![expected_row(2, Some(1)), expected_row(1, None)],
3014 );
3015
3016 session
3017 .execute_trusted_dynamic_mutation_batch(vec![insert(3, None), insert(4, Some(3))])
3018 .expect("the update-away fixture should commit");
3019 let updated_away = session
3020 .execute_trusted_dynamic_mutation_batch(vec![update_parent(4, None), delete(3)])
3021 .expect("an updated final source may release a deleted target");
3022 assert_eq!(
3023 updated_away.rows,
3024 vec![expected_row(4, None), expected_row(3, None)],
3025 );
3026
3027 session
3028 .execute_trusted_dynamic_mutation_batch(vec![insert(5, None), insert(6, Some(5))])
3029 .expect("the retained-reference fixture should commit");
3030 let retained = session
3031 .execute_trusted_dynamic_mutation_batch(vec![update_parent(6, Some(5)), delete(5)])
3032 .expect_err("a final updated source must still block target deletion");
3033 assert_relation_violation(&retained);
3034
3035 session
3036 .execute_trusted_dynamic_mutation(&insert(7, None))
3037 .expect("the inserted-reference fixture target should commit");
3038 let inserted_reference = session
3039 .execute_trusted_dynamic_mutation_batch(vec![insert(8, Some(7)), delete(7)])
3040 .expect_err("a final inserted source must not reference a deleted target");
3041 assert_relation_violation(&inserted_reference);
3042
3043 let inserted_target = session
3044 .execute_trusted_dynamic_mutation_batch(vec![insert(10, Some(9)), insert(9, None)])
3045 .expect("an inserted relation should see its batch-final target");
3046 assert_eq!(
3047 inserted_target.rows,
3048 vec![expected_row(10, Some(9)), expected_row(9, None)],
3049 );
3050
3051 session
3052 .execute_trusted_dynamic_mutation(&insert(11, None))
3053 .expect("the updated-reference fixture source should commit");
3054 let updated_target = session
3055 .execute_trusted_dynamic_mutation_batch(vec![
3056 update_parent(11, Some(12)),
3057 insert(12, None),
3058 ])
3059 .expect("an updated relation should see its batch-final target");
3060 assert_eq!(
3061 updated_target.rows,
3062 vec![expected_row(11, Some(12)), expected_row(12, None)],
3063 );
3064 }
3065
3066 #[test]
3067 fn mixed_batch_rejects_cross_entity_missing_and_collision_then_honors_replace() {
3068 let session = initialize();
3069 session
3070 .execute_trusted_dynamic_mutation(&insert(1, None))
3071 .expect("the primary mixed fixture row should commit");
3072 session
3073 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3074 entity: OTHER_ENTITY_NAME.to_string(),
3075 patch: other_patch(Some(1), 10),
3076 })
3077 .expect("the secondary mixed fixture row should commit");
3078
3079 let mixed_entity = session
3080 .execute_trusted_dynamic_mutation_batch(vec![
3081 update_code(1, 11),
3082 DynamicMutation::Update {
3083 entity: OTHER_ENTITY_NAME.to_string(),
3084 key: InputValue::Nat64(1),
3085 patch: other_patch(None, 11),
3086 },
3087 ])
3088 .expect_err("one atomic batch must not cross accepted entities");
3089 assert_eq!(mixed_entity.class(), ErrorClass::Conflict);
3090 assert_eq!(
3091 mixed_entity.diagnostic_facts(),
3092 vec![
3093 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3094 (
3095 icydb_diagnostic_code::DiagnosticFactTag::ExpectedEntityTag,
3096 ENTITY_TAG.value(),
3097 ),
3098 (
3099 icydb_diagnostic_code::DiagnosticFactTag::ActualEntityTag,
3100 OTHER_ENTITY_TAG.value(),
3101 ),
3102 ],
3103 );
3104
3105 let missing = session
3106 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 12), delete(99)])
3107 .expect_err("a late missing delete must reject the earlier staged update");
3108 assert_eq!(missing.class(), ErrorClass::NotFound);
3109
3110 session
3111 .execute_trusted_dynamic_mutation(&insert(2, None))
3112 .expect("the collision fixture should commit");
3113 let collision = session
3114 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 13), insert(2, None)])
3115 .expect_err("an insert collision must reject the earlier staged update");
3116 assert_eq!(collision.class(), ErrorClass::Conflict);
3117 let failures_unchanged = session
3118 .execute_trusted_dynamic_mutation(&update_code(1, 1))
3119 .expect("failed batches must preserve the original unique value");
3120 assert_eq!(failures_unchanged.affected_rows, 0);
3121
3122 let replaced = session
3123 .execute_trusted_dynamic_mutation_batch(vec![
3124 update_code(1, 14),
3125 DynamicMutation::Replace {
3126 entity: ENTITY_NAME.to_string(),
3127 key: InputValue::Nat64(99),
3128 patch: patch(None, None, Some(99)),
3129 },
3130 ])
3131 .expect("ordinary caller-key replace should insert its absent final row");
3132 assert_eq!(
3133 replaced.rows,
3134 vec![
3135 expected_row_with_code(1, None, 14),
3136 expected_row_with_code(99, None, 99),
3137 ],
3138 );
3139
3140 let unchanged = session
3141 .execute_trusted_dynamic_mutation(&update_code(1, 14))
3142 .expect("the successful mixed replace must publish its preceding update");
3143 assert_eq!(unchanged.affected_rows, 0);
3144 let other_unchanged = session
3145 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3146 entity: OTHER_ENTITY_NAME.to_string(),
3147 key: InputValue::Nat64(1),
3148 patch: other_patch(None, 10),
3149 })
3150 .expect("cross-entity rejection must preserve the secondary row");
3151 assert_eq!(other_unchanged.affected_rows, 0);
3152 }
3153
3154 #[test]
3155 fn mixed_batch_unique_swap_and_delete_release_use_the_final_overlay() {
3156 let session = initialize();
3157 session
3158 .execute_trusted_dynamic_mutation_batch(vec![
3159 insert_with_code(1, None, 10),
3160 insert_with_code(2, None, 20),
3161 ])
3162 .expect("the unique-overlay fixture should commit");
3163
3164 let swapped = session
3165 .execute_trusted_dynamic_mutation_batch(vec![update_code(1, 20), update_code(2, 10)])
3166 .expect("two final rows should atomically swap unique memberships");
3167 assert_eq!(
3168 swapped.rows,
3169 vec![
3170 expected_row_with_code(1, None, 20),
3171 expected_row_with_code(2, None, 10),
3172 ],
3173 );
3174
3175 let released = session
3176 .execute_trusted_dynamic_mutation_batch(vec![delete(1), insert_with_code(3, None, 20)])
3177 .expect("a delete should release unique membership to a final inserted row");
3178 assert_eq!(
3179 released.rows,
3180 vec![
3181 expected_row_with_code(1, None, 20),
3182 expected_row_with_code(3, None, 20),
3183 ],
3184 );
3185 }
3186}
3187
3188#[cfg(test)]
3189mod identity_pre_key_tests {
3190 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3191 use super::DynamicTypedEntityBinding;
3192 use super::{
3193 AcceptedMutationIntentPatch, AcceptedRowLayoutRuntimeContract, AcceptedStructuralMutation,
3194 AcceptedStructuralMutationPacking, AcceptedStructuralMutationStagedAdmission,
3195 AcceptedStructuralMutationTarget, DbSession, DynamicMutation, DynamicStructuralPatch,
3196 DynamicTypedFieldBindingRequest, DynamicTypedFieldType, DynamicTypedMutation,
3197 DynamicWriteCell, FieldSlot, MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS,
3198 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES,
3199 MutationProgressRecordOp, add_structural_mutation_staged_bytes,
3200 admit_structural_mutation_staged_charge, checked_pre_key_candidate_count,
3201 insert_key_exists_after_generation, structural_mutation_staged_charge,
3202 validate_structural_mutation_result_bytes,
3203 };
3204 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3205 use crate::db::data::DecodedDataStoreKey;
3206 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3207 use crate::db::executor::budget::{
3208 HardExecutionBudget, HardExecutionContext, HardExecutionFailureHeadroom,
3209 with_query_execution_budget_for_tests,
3210 };
3211 use crate::db::mutation_job::{MutationJobRecord, MutationJobTransition};
3212 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3213 use crate::db::{
3214 CompareProofAndAdvanceError, DynamicQuery, ExhaustiveReadError, RawDataStoreKey,
3215 ReadSetRevisionError, ResumableJobAdvance, ResumableJobAdvanceRequest,
3216 ResumableJobAdvanceStatus, ResumableJobError, ResumableJobId, ResumableJobIdempotencyKey,
3217 ResumableJobStatus, asc,
3218 };
3219 use crate::{
3220 db::{
3221 GeneratedStartupDriverStep, MutationJobAdvanceRequest, MutationJobId,
3222 MutationJobIdempotencyKey, MutationJobPhase, MutationJobStatus,
3223 commit::{
3224 database_incarnation_id, forget_recovered_domain_for_tests,
3225 install_startup_recovery_wakeup,
3226 },
3227 data::DataStore,
3228 drive_generated_startup_recovery_page,
3229 executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
3230 index::{IndexId, IndexKey, IndexKeyKind, IndexStore, IndexStoreVisit},
3231 integrity::{
3232 InsertMutationJobResult, PhysicalUnitCheckpoint, QuickIntegrityStatus,
3233 RowInspectionLimits, execute_quick_integrity, execute_row_integrity_page,
3234 with_mutation_progress_store,
3235 },
3236 journal::{
3237 JournalBatch, JournalRecord, JournalSequence, JournalTailControl, JournalTailStore,
3238 encode_journal_batch,
3239 },
3240 registry::{
3241 StoreAllocationIdentities, StoreAllocationIdentity, StoreHandle, StoreRegistry,
3242 StoreRuntimeStorageCapabilities,
3243 },
3244 schema::{
3245 AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
3246 FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
3247 PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
3248 PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
3249 SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
3250 accepted_schema_candidate_with_field_bindings_for_tests,
3251 cardinality_build::{
3252 CardinalityBuildAuthority, CardinalityGenerationPageOutcome,
3253 drive_cardinality_generation_page,
3254 },
3255 cardinality_generation::{CardinalityGenerationHeader, CardinalityGenerationState},
3256 },
3257 write_context::MutationMode,
3258 },
3259 error::{ErrorClass, ErrorOrigin, InternalError},
3260 testing::test_memory,
3261 traits::{CanisterKind, Path},
3262 types::{EntityTag, Timestamp},
3263 value::{InputValue, OutputValue, Value},
3264 };
3265 use icydb_schema::{FieldSourceKey, ScalarType};
3266 use std::{
3267 cell::{Cell, RefCell},
3268 collections::BTreeMap,
3269 time::Instant,
3270 };
3271
3272 const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
3273 const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
3274 const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
3275 const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
3276 const ENTITY_NAME: &str = "IdentityRow";
3277 const ENTITY_TAG: EntityTag = EntityTag::new(93);
3278 const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
3279 const UNRELATED_STORE_PATH: &str = "session::write::identity_pre_key_tests::UnrelatedStore";
3280
3281 struct TestCanister;
3282
3283 impl Path for TestCanister {
3284 const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
3285 }
3286
3287 impl CanisterKind for TestCanister {
3288 const COMMIT_MEMORY_ID: u8 = 45;
3289 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
3290 const STARTUP_MEMORY_ID: u8 = 49;
3291 const STARTUP_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.startup.control.v1";
3292 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
3293 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3294 "icydb.identity_pre_key_tests.integrity.progress.v1";
3295 }
3296
3297 thread_local! {
3298 static STARTUP_WAKEUPS: Cell<u32> = const { Cell::new(0) };
3299 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
3300 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
3301 static SCHEMA_STORE: RefCell<SchemaStore> =
3302 const { RefCell::new(SchemaStore::init_heap()) };
3303 static UNRELATED_DATA_STORE: RefCell<DataStore> =
3304 const { RefCell::new(DataStore::init_heap()) };
3305 static UNRELATED_INDEX_STORE: RefCell<IndexStore> =
3306 const { RefCell::new(IndexStore::init_heap()) };
3307 static UNRELATED_SCHEMA_STORE: RefCell<SchemaStore> =
3308 const { RefCell::new(SchemaStore::init_heap()) };
3309 static STORE_REGISTRY: StoreRegistry = {
3310 let mut registry = StoreRegistry::new();
3311 registry.register_store(
3312 STORE_PATH,
3313 &DATA_STORE,
3314 &INDEX_STORE,
3315 &SCHEMA_STORE,
3316 StoreAllocationIdentities::absent(),
3317 StoreRuntimeStorageCapabilities::heap(),
3318 ).expect("identity pre-key test store should register");
3319 registry.register_store(
3320 UNRELATED_STORE_PATH,
3321 &UNRELATED_DATA_STORE,
3322 &UNRELATED_INDEX_STORE,
3323 &UNRELATED_SCHEMA_STORE,
3324 StoreAllocationIdentities::absent(),
3325 StoreRuntimeStorageCapabilities::heap(),
3326 ).expect("unrelated identity test store should register");
3327 registry
3328 };
3329 static JOURNALED_DATA_STORE: RefCell<DataStore> =
3330 RefCell::new(DataStore::init_journaled(test_memory(186)));
3331 static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
3332 RefCell::new(IndexStore::init_journaled(test_memory(187)));
3333 static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
3334 RefCell::new(SchemaStore::init_journaled(test_memory(188)));
3335 static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
3336 RefCell::new(JournalTailStore::init(test_memory(189)));
3337 static JOURNALED_STORE_REGISTRY: StoreRegistry = {
3338 let mut registry = StoreRegistry::new();
3339 registry.register_journaled_store(
3340 JOURNALED_STORE_PATH,
3341 &JOURNALED_DATA_STORE,
3342 &JOURNALED_INDEX_STORE,
3343 &JOURNALED_SCHEMA_STORE,
3344 &JOURNALED_TAIL_STORE,
3345 StoreAllocationIdentities::new_journaled(
3346 StoreAllocationIdentity::new(186, "icydb.test.identity_range.data.v1"),
3347 StoreAllocationIdentity::new(187, "icydb.test.identity_range.index.v1"),
3348 StoreAllocationIdentity::new(188, "icydb.test.identity_range.schema.v1"),
3349 StoreAllocationIdentity::new(189, "icydb.test.identity_range.journal.v1"),
3350 ),
3351 StoreRuntimeStorageCapabilities::journaled(),
3352 ).expect("identity range journaled store should register");
3353 registry
3354 };
3355 }
3356
3357 fn record_startup_wakeup() {
3358 STARTUP_WAKEUPS.with(|wakeups| wakeups.set(wakeups.get().saturating_add(1)));
3359 }
3360
3361 struct JournaledTestCanister;
3362
3363 impl Path for JournaledTestCanister {
3364 const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
3365 }
3366
3367 impl CanisterKind for JournaledTestCanister {
3368 const COMMIT_MEMORY_ID: u8 = 190;
3369 const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
3370 const STARTUP_MEMORY_ID: u8 = 192;
3371 const STARTUP_STABLE_KEY: &'static str = "icydb.identity_range_tests.startup.control.v1";
3372 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
3373 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3374 "icydb.identity_range_tests.integrity.progress.v1";
3375 }
3376
3377 fn source_key(source: &str) -> FieldSourceKey {
3378 FieldSourceKey::try_new(source).expect("identity test field source should admit")
3379 }
3380
3381 fn identity_snapshot(store_path: &str, payload_unique: bool) -> PersistedSchemaSnapshot {
3382 let fields = vec![
3383 PersistedFieldSnapshot::new_initial_with_write_policy(
3384 FieldId::new(1),
3385 "id".to_string(),
3386 SchemaFieldSlot::new(0),
3387 AcceptedFieldKind::Nat64,
3388 Vec::new(),
3389 false,
3390 SchemaInsertDefault::None,
3391 SchemaFieldWritePolicy::from_model_policies(
3392 Some(FieldInsertGeneration::Identity),
3393 None,
3394 ),
3395 FieldStorageDecode::ByKind,
3396 LeafCodec::Scalar(ScalarCodec::Nat64),
3397 ),
3398 PersistedFieldSnapshot::new_initial(
3399 FieldId::new(2),
3400 "payload".to_string(),
3401 SchemaFieldSlot::new(1),
3402 AcceptedFieldKind::Nat64,
3403 Vec::new(),
3404 false,
3405 SchemaInsertDefault::None,
3406 FieldStorageDecode::ByKind,
3407 LeafCodec::Scalar(ScalarCodec::Nat64),
3408 ),
3409 ];
3410 PersistedSchemaSnapshot::new_with_indexes(
3411 SchemaVersion::initial(),
3412 ENTITY_SOURCE.to_string(),
3413 ENTITY_NAME.to_string(),
3414 FieldId::new(1),
3415 SchemaRowLayout::initial(
3416 fields
3417 .iter()
3418 .map(|field| (field.id(), field.slot()))
3419 .collect(),
3420 ),
3421 fields,
3422 vec![PersistedIndexSnapshot::new(
3423 SchemaIndexId::new(1).expect("identity test index ID should admit"),
3424 1,
3425 "by_payload".to_string(),
3426 store_path.to_string(),
3427 payload_unique,
3428 PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
3429 FieldId::new(2),
3430 SchemaFieldSlot::new(1),
3431 vec!["payload".to_string()],
3432 AcceptedFieldKind::Nat64,
3433 false,
3434 )]),
3435 None,
3436 )],
3437 )
3438 }
3439
3440 fn initialize() -> DbSession<TestCanister> {
3441 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
3442 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
3443 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
3444 UNRELATED_DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
3445 UNRELATED_INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
3446 UNRELATED_SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
3447 let session = DbSession::<TestCanister>::new(
3448 &STORE_REGISTRY,
3449 &crate::db::RequestExecutionRoot::__new_runtime_root(),
3450 );
3451 session
3452 .db
3453 .drive_startup_recovery_page()
3454 .expect("identity pre-key test database should initialize");
3455 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
3456 STORE_PATH,
3457 AcceptedSchemaRevision::INITIAL,
3458 BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH, false))]),
3459 BTreeMap::from([
3460 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
3461 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
3462 ]),
3463 );
3464 let store = session
3465 .db
3466 .store_handle(STORE_PATH)
3467 .expect("identity pre-key test store should resolve");
3468 crate::db::commit::publish_accepted_schema_candidate(
3469 STORE_PATH,
3470 store,
3471 AcceptedSchemaRevision::NONE,
3472 &candidate,
3473 )
3474 .expect("identity candidate should publish with explicit zero state");
3475 session
3476 }
3477
3478 fn initialize_journaled_with_root_and_payload_uniqueness(
3479 payload_unique: bool,
3480 ) -> (
3481 DbSession<JournaledTestCanister>,
3482 crate::db::RequestExecutionRoot,
3483 ) {
3484 let root = crate::db::RequestExecutionRoot::__new_runtime_root();
3485 let session = DbSession::<JournaledTestCanister>::new(&JOURNALED_STORE_REGISTRY, &root);
3486 session
3487 .db
3488 .drive_startup_recovery_page()
3489 .expect("journaled identity database should initialize");
3490 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
3491 JOURNALED_STORE_PATH,
3492 AcceptedSchemaRevision::INITIAL,
3493 BTreeMap::from([(
3494 ENTITY_TAG,
3495 identity_snapshot(JOURNALED_STORE_PATH, payload_unique),
3496 )]),
3497 BTreeMap::from([
3498 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
3499 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
3500 ]),
3501 );
3502 let store = session
3503 .db
3504 .store_handle(JOURNALED_STORE_PATH)
3505 .expect("journaled identity store should resolve");
3506 crate::db::commit::publish_accepted_schema_candidate(
3507 JOURNALED_STORE_PATH,
3508 store,
3509 AcceptedSchemaRevision::NONE,
3510 &candidate,
3511 )
3512 .expect("journaled identity candidate should publish");
3513 (session, root)
3514 }
3515
3516 fn initialize_journaled_with_root() -> (
3517 DbSession<JournaledTestCanister>,
3518 crate::db::RequestExecutionRoot,
3519 ) {
3520 initialize_journaled_with_root_and_payload_uniqueness(false)
3521 }
3522
3523 fn initialize_journaled() -> DbSession<JournaledTestCanister> {
3524 initialize_journaled_with_root().0
3525 }
3526
3527 fn initialize_journaled_with_unique_payload() -> DbSession<JournaledTestCanister> {
3528 initialize_journaled_with_root_and_payload_uniqueness(true).0
3529 }
3530
3531 fn drive_journaled_recovery_to_completion(session: &DbSession<JournaledTestCanister>) {
3532 for _ in 0..8 {
3533 if session
3534 .db
3535 .drive_startup_recovery_page()
3536 .expect("dedicated driver recovery should remain valid")
3537 {
3538 return;
3539 }
3540 }
3541 panic!("dedicated driver recovery should quiesce within eight complete batches");
3542 }
3543
3544 fn drive_journaled_cardinality_to_ready(session: &DbSession<JournaledTestCanister>) {
3545 let handle = session
3546 .db
3547 .store_handle(JOURNALED_STORE_PATH)
3548 .expect("journaled cardinality store should resolve");
3549 for _ in 0..8 {
3550 let outcome = handle
3551 .with_data(|data| {
3552 handle.with_index(|index| {
3553 handle.with_schema_mut(|schema| {
3554 drive_cardinality_generation_page(data, index, schema, |schema| {
3555 let watermark = JOURNALED_TAIL_STORE
3556 .with(|tail| tail.borrow().fold_watermark())?;
3557 CardinalityBuildAuthority::derive(
3558 schema,
3559 database_incarnation_id()?,
3560 handle.allocation_identities(),
3561 watermark,
3562 )
3563 })
3564 })
3565 })
3566 })
3567 .expect("bounded cardinality generation should advance");
3568 if outcome == CardinalityGenerationPageOutcome::Quiescent {
3569 return;
3570 }
3571 }
3572 panic!("cardinality generation should become Ready within eight bounded pages");
3573 }
3574
3575 fn journaled_user_index_prefix() -> (IndexId, Vec<Vec<u8>>) {
3576 JOURNALED_INDEX_STORE.with(|store| {
3577 let mut selected = None;
3578 store
3579 .borrow()
3580 .visit_entries(|raw_key, _value| {
3581 let key = IndexKey::try_from_raw(raw_key)
3582 .expect("accepted user index key should decode");
3583 if key.key_kind() != IndexKeyKind::User {
3584 return Ok::<_, InternalError>(IndexStoreVisit::Continue);
3585 }
3586 let components = (0..key.component_count())
3587 .map(|index| {
3588 key.component(index)
3589 .expect("accepted index component should exist")
3590 .to_vec()
3591 })
3592 .collect::<Vec<_>>();
3593 selected = Some((*key.index_id(), components));
3594 Ok(IndexStoreVisit::Stop)
3595 })
3596 .expect("accepted user index should be inspectable");
3597 selected.expect("the cardinality fixture should contain one user index entry")
3598 })
3599 }
3600
3601 fn reset_journaled_cardinality_projections() -> u64 {
3602 JOURNALED_DATA_STORE.with(|store| {
3603 store
3604 .borrow_mut()
3605 .reset_journaled_live_projection()
3606 .expect("row projection should reset without a count scan");
3607 });
3608 let data_generation = JOURNALED_DATA_STORE.with(|store| store.borrow().generation());
3609 let fold_watermark = JOURNALED_TAIL_STORE
3610 .with(|store| store.borrow().fold_watermark())
3611 .expect("journal watermark should remain current-form");
3612 JOURNALED_INDEX_STORE.with(|store| {
3613 store
3614 .borrow_mut()
3615 .reset_journaled_live_projection(data_generation, fold_watermark)
3616 .expect("index projection should reset without a count scan");
3617 });
3618 data_generation
3619 }
3620
3621 fn assert_journaled_cardinality(
3622 handle: StoreHandle,
3623 index_id: IndexId,
3624 prefix_components: &[Vec<u8>],
3625 expected: u64,
3626 ) {
3627 assert_eq!(handle.exact_entity_count(ENTITY_TAG), Some(expected));
3628 let data_generation = JOURNALED_DATA_STORE.with(|store| store.borrow().generation());
3629 assert_eq!(
3630 handle.exact_user_index_prefix_count(
3631 data_generation,
3632 IndexKeyKind::User,
3633 index_id,
3634 prefix_components,
3635 ),
3636 Some(expected),
3637 );
3638 }
3639
3640 fn mark_journaled_cardinality_building() {
3641 let current = JOURNALED_SCHEMA_STORE.with(|store| {
3642 store
3643 .borrow()
3644 .cardinality_generation_header()
3645 .expect("Ready header should decode")
3646 .expect("Ready header should exist")
3647 });
3648 JOURNALED_SCHEMA_STORE.with(|store| {
3649 store
3650 .borrow_mut()
3651 .write_cardinality_generation_header(CardinalityGenerationHeader::new(
3652 current.generation(),
3653 CardinalityGenerationState::Building,
3654 current.slot(),
3655 current.source(),
3656 ))
3657 .expect("Building fallback fixture should persist");
3658 });
3659 }
3660
3661 fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
3662 AcceptedMutationIntentPatch::new()
3663 .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
3664 }
3665
3666 fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
3667 DynamicStructuralPatch::new(vec![(
3668 "payload".to_string(),
3669 DynamicWriteCell::Value(InputValue::Nat64(value)),
3670 )])
3671 }
3672
3673 fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
3674 vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
3675 }
3676
3677 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3678 fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
3679 session
3680 .issue_typed_entity_binding(
3681 ENTITY_SOURCE,
3682 &[
3683 DynamicTypedFieldBindingRequest::new(
3684 ID_SOURCE.to_string(),
3685 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3686 false,
3687 ),
3688 DynamicTypedFieldBindingRequest::new(
3689 PAYLOAD_SOURCE.to_string(),
3690 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3691 false,
3692 ),
3693 ],
3694 )
3695 .expect("exact-key test binding should issue")
3696 }
3697
3698 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3699 fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
3700 let output = session
3701 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3702 entity: ENTITY_NAME.to_string(),
3703 patch: dynamic_payload_patch(payload),
3704 })
3705 .expect("exact-key fixture insert should commit");
3706 match output.rows.as_slice() {
3707 [row] => match row.as_slice() {
3708 [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
3709 if *actual_payload == payload =>
3710 {
3711 *id
3712 }
3713 _ => panic!("exact-key fixture should return its identity and payload"),
3714 },
3715 _ => panic!("exact-key fixture insert should return one row"),
3716 }
3717 }
3718
3719 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3720 fn identity_row_stored_bytes<C: CanisterKind>(
3721 session: &DbSession<C>,
3722 store_path: &'static str,
3723 key: u64,
3724 ) -> u64 {
3725 let data_key = DecodedDataStoreKey::try_from_structural_key(ENTITY_TAG, &Value::Nat64(key))
3726 .expect("identity row key should encode");
3727 let raw_key = data_key.to_raw().expect("identity raw key should encode");
3728 let store = session
3729 .db
3730 .recovered_store(store_path)
3731 .expect("identity store should resolve");
3732 store.with_data(|data_store| {
3733 u64::try_from(
3734 data_store
3735 .get(&raw_key)
3736 .expect("inserted identity row should exist")
3737 .len(),
3738 )
3739 .expect("bounded row length should fit u64")
3740 })
3741 }
3742
3743 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3744 fn with_stored_bytes_limit<T>(
3745 limit: u64,
3746 shape_fingerprint_prefix: u64,
3747 operation: impl FnOnce() -> Result<T, crate::db::query::intent::QueryError>,
3748 ) -> Result<T, crate::db::query::intent::QueryError> {
3749 let budget = HardExecutionBudget::uniform_for_tests(
3750 u64::MAX,
3751 HardExecutionFailureHeadroom::new(500, 256),
3752 )
3753 .with_limit_for_tests(
3754 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::StoredBytesRead,
3755 limit,
3756 );
3757 let context = HardExecutionContext::new(
3758 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3759 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3760 shape_fingerprint_prefix,
3761 );
3762
3763 with_query_execution_budget_for_tests(budget, context, operation)
3764 }
3765
3766 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3767 fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
3768 let first = insert_exact_key_fixture(session, 41);
3769 let second = insert_exact_key_fixture(session, 42);
3770 let missing = u64::MAX;
3771 let binding = exact_key_binding(session);
3772 let gets_before = DataStore::current_get_call_count();
3773 let result = session
3774 .execute_public_exact_key_batch_for_typed_binding(
3775 &binding,
3776 &[second, missing, first, second],
3777 )
3778 .expect("exact-key batch should execute")
3779 .expect("exact-key binding should remain current");
3780
3781 assert_eq!(result.positions, vec![0, 1, 2, 0]);
3782 assert_eq!(
3783 result.distinct_rows,
3784 vec![
3785 Some(expected_dynamic_row(second, 42)),
3786 None,
3787 Some(expected_dynamic_row(first, 41)),
3788 ],
3789 );
3790 assert_eq!(
3791 DataStore::current_get_call_count().saturating_sub(gets_before),
3792 3,
3793 "four input positions with one duplicate must perform three physical reads",
3794 );
3795 }
3796
3797 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3798 #[test]
3799 fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
3800 assert_exact_key_batch(&initialize());
3801 assert_exact_key_batch(&initialize_journaled());
3802 }
3803
3804 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3805 fn assert_primary_range_materialization_fetches_once<C: CanisterKind>(
3806 session: &DbSession<C>,
3807 store_path: &'static str,
3808 ) {
3809 let key = insert_exact_key_fixture(session, 41);
3810 let stored_bytes = identity_row_stored_bytes(session, store_path, key);
3811
3812 let scalar = DynamicQuery::new(ENTITY_NAME)
3813 .select(["id", "payload"])
3814 .order_by(asc("id"))
3815 .limit(1);
3816 let gets_before = DataStore::current_get_call_count();
3817 let scalar_page = with_stored_bytes_limit(stored_bytes, 0x7072_696d_6172_792d, || {
3818 session.execute_trusted_live_page(&scalar, None)
3819 })
3820 .expect("one scalar primary-range row should fit one payload-read allowance");
3821 assert_eq!(scalar_page.row_count, 1);
3822 assert_eq!(
3823 DataStore::current_get_call_count().saturating_sub(gets_before),
3824 1,
3825 "scalar primary traversal should fetch its emitted row exactly once",
3826 );
3827
3828 let grouped = DynamicQuery::new(ENTITY_NAME)
3829 .group_by("payload")
3830 .aggregate(crate::db::count())
3831 .grouped_limits(10, 16 * 1_024)
3832 .limit(1);
3833 let gets_before = DataStore::current_get_call_count();
3834 let grouped_page = with_stored_bytes_limit(stored_bytes, 0x6772_6f75_7065_642d, || {
3835 session.execute_trusted_dynamic_grouped_query(&grouped)
3836 })
3837 .expect("one grouped primary-range row should fit one payload-read allowance");
3838 assert_eq!(grouped_page.row_count, 1);
3839 assert_eq!(
3840 DataStore::current_get_call_count().saturating_sub(gets_before),
3841 1,
3842 "grouped primary traversal should fetch its source row exactly once",
3843 );
3844 }
3845
3846 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3847 #[test]
3848 fn row_materialization_fetches_each_required_payload_at_most_once() {
3849 assert_primary_range_materialization_fetches_once(&initialize(), STORE_PATH);
3850 assert_primary_range_materialization_fetches_once(
3851 &initialize_journaled(),
3852 JOURNALED_STORE_PATH,
3853 );
3854 }
3855
3856 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3857 #[test]
3858 fn ordered_grouped_pages_close_a_group_spanning_physical_refills_before_resume() {
3859 let session = initialize();
3860 let mut patches = Vec::new();
3861 for _ in 0..70 {
3862 patches.push(dynamic_payload_patch(10));
3863 }
3864 for _ in 0..3 {
3865 patches.push(dynamic_payload_patch(20));
3866 }
3867 patches.push(dynamic_payload_patch(30));
3868 let inserted = session
3869 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, patches)
3870 .expect("ordered grouped continuation rows should insert");
3871 assert_eq!(inserted.rows.len(), 74);
3872
3873 let query = DynamicQuery::new(ENTITY_NAME)
3874 .group_by("payload")
3875 .aggregate(crate::db::count())
3876 .aggregate(crate::db::sum("id"))
3877 .order_by(asc("payload"))
3878 .grouped_limits(4, 16 * 1_024)
3879 .limit(1);
3880 let expected = [
3881 (10_u64, 70_u64, crate::types::Decimal::new(2_485, 0)),
3882 (20, 3, crate::types::Decimal::new(216, 0)),
3883 (30, 1, crate::types::Decimal::new(74, 0)),
3884 ];
3885 let mut continuation: Option<String> = None;
3886 let mut seen_cursors = std::collections::BTreeSet::new();
3887
3888 for (page_index, (group_key, row_count, id_sum)) in expected.into_iter().enumerate() {
3889 let request = continuation.as_ref().map_or_else(
3890 || query.clone(),
3891 |cursor| query.clone().cursor(cursor.clone()),
3892 );
3893 let entries_before = IndexStore::current_entry_read_count();
3894 let rows_before = DataStore::current_get_call_count();
3895 let page = session
3896 .execute_trusted_dynamic_grouped_query(&request)
3897 .unwrap_or_else(|error| {
3898 panic!("ordered grouped page {page_index} should execute: {error:?}")
3899 });
3900 let entries_read =
3901 IndexStore::current_entry_read_count().saturating_sub(entries_before);
3902 let rows_read = DataStore::current_get_call_count().saturating_sub(rows_before);
3903
3904 assert_eq!(page.row_count, 1);
3905 let [row] = page.rows.as_slice() else {
3906 panic!("ordered grouped page must contain exactly one closed group")
3907 };
3908 assert_eq!(row.group_key(), &[OutputValue::Nat64(group_key)]);
3909 assert_eq!(
3910 row.aggregate_values(),
3911 &[OutputValue::Nat64(row_count), OutputValue::Decimal(id_sum),],
3912 );
3913 if page_index == 0 {
3914 assert!(
3915 entries_read.saturating_add(rows_read) >= 70,
3916 "the first closed group must span the maintained 64-entry physical refill",
3917 );
3918 }
3919
3920 continuation = page.next_cursor;
3921 if page_index + 1 < expected.len() {
3922 let cursor = continuation
3923 .as_ref()
3924 .expect("another closed group should retain continuation");
3925 assert!(
3926 seen_cursors.insert(cursor.clone()),
3927 "ordered grouped continuation must advance monotonically",
3928 );
3929 } else {
3930 assert_eq!(continuation, None);
3931 }
3932 }
3933 }
3934
3935 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3936 #[test]
3937 fn exhaustive_pages_require_and_recompare_the_complete_source_proof() {
3938 let session = initialize();
3939 let first = insert_exact_key_fixture(&session, 41);
3940 let second = insert_exact_key_fixture(&session, 42);
3941 let third = insert_exact_key_fixture(&session, 43);
3942 let query = DynamicQuery::new(ENTITY_NAME)
3943 .select(["id", "payload"])
3944 .order_by(asc("id"));
3945
3946 let page = session
3947 .execute_trusted_exhaustive_page(&query, None, None)
3948 .expect("initial exhaustive page should capture its source proof");
3949 assert_eq!(
3950 page.rows,
3951 vec![
3952 expected_dynamic_row(first, 41),
3953 expected_dynamic_row(second, 42),
3954 ],
3955 );
3956 let continuation = page
3957 .continuation
3958 .as_deref()
3959 .expect("unreturned row should retain exhaustive continuation");
3960 assert!(matches!(
3961 session.execute_trusted_exhaustive_page(&query, Some(continuation), None),
3962 Err(ExhaustiveReadError::Revision(
3963 ReadSetRevisionError::ResumeProofRequired
3964 )),
3965 ));
3966 let resumed = session
3967 .execute_trusted_exhaustive_page(&query, Some(continuation), Some(&page.proof))
3968 .expect("unchanged proof should resume exhaustive traversal");
3969 assert_eq!(resumed.rows, vec![expected_dynamic_row(third, 43)]);
3970 assert_eq!(resumed.continuation, None);
3971
3972 let stale_page = session
3973 .execute_trusted_exhaustive_page(&query, None, None)
3974 .expect("fresh exhaustive page should capture current revision");
3975 let stale_continuation = stale_page
3976 .continuation
3977 .as_deref()
3978 .expect("fresh three-row traversal should retain continuation");
3979 let _ = insert_exact_key_fixture(&session, 44);
3980 assert!(matches!(
3981 session.execute_trusted_exhaustive_page(
3982 &query,
3983 Some(stale_continuation),
3984 Some(&stale_page.proof),
3985 ),
3986 Err(ExhaustiveReadError::Revision(
3987 ReadSetRevisionError::StoreDataChanged { .. }
3988 )),
3989 ));
3990 }
3991
3992 #[cfg(all(feature = "sql", feature = "diagnostics"))]
3993 #[test]
3994 fn heap_sources_cannot_back_durable_resumable_jobs() {
3995 let session = initialize();
3996 let proof = session
3997 .capture_read_set_revision_proof(&[ENTITY_NAME])
3998 .expect("heap source proof should capture for one-call exhaustive reads");
3999 let job_id = ResumableJobId::try_from_bytes([70; 32])
4000 .expect("nonzero heap test job identity should admit");
4001
4002 assert!(matches!(
4003 session.start_resumable_job(job_id, proof, Vec::new()),
4004 Err(ResumableJobError::SourceProof(
4005 ReadSetRevisionError::DurableStoreRequired { .. }
4006 )),
4007 ));
4008 }
4009
4010 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4011 #[test]
4012 fn proof_and_progress_controls_charge_one_shared_request_scope() {
4013 let (session, root) = initialize_journaled_with_root();
4014 let resource = icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions;
4015 let before = root.observed(resource);
4016 let proof = session
4017 .capture_read_set_revision_proof(&[ENTITY_NAME])
4018 .expect("proof capture should use the retained request scope");
4019 let job_id = ResumableJobId::try_from_bytes([75; 32])
4020 .expect("nonzero accounting job identity should admit");
4021 session
4022 .start_resumable_job(job_id, proof, Vec::new())
4023 .expect("job start should use the same retained request scope");
4024 let _ = session
4025 .resumable_job_state(job_id)
4026 .expect("job load should use the same retained request scope");
4027
4028 assert_eq!(root.observed(resource).saturating_sub(before), 3);
4029 }
4030
4031 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4032 #[test]
4033 fn source_proofs_ignore_unrelated_stores_but_bind_access_state_changes() {
4034 let session = initialize();
4035 let proof = session
4036 .capture_read_set_revision_proof(&[ENTITY_NAME])
4037 .expect("source proof should cover only the entity's physical store");
4038 let shared_store_proof = session
4039 .capture_read_set_revision_proof(&[ENTITY_NAME, ENTITY_NAME])
4040 .expect("entities sharing one physical source should deduplicate");
4041 assert_eq!(shared_store_proof, proof);
4042 assert_eq!(shared_store_proof.stores().len(), 1);
4043 let unrelated = session
4044 .db
4045 .store_handle(UNRELATED_STORE_PATH)
4046 .expect("unrelated registered store should resolve");
4047 unrelated.with_data_mut(|store| {
4048 let _ = store.remove(&RawDataStoreKey::from_persisted_bytes(vec![1]));
4049 });
4050 session
4051 .verify_read_set_revision_proof(&proof)
4052 .expect("a nonparticipating store mutation must not invalidate the proof");
4053
4054 let source = session
4055 .db
4056 .store_handle(STORE_PATH)
4057 .expect("participating source store should resolve");
4058 source
4059 .mark_index_building()
4060 .expect("source access-state transition should advance its revision");
4061 assert!(matches!(
4062 session.verify_read_set_revision_proof(&proof),
4063 Err(ExhaustiveReadError::Revision(
4064 ReadSetRevisionError::StoreAccessChanged { .. }
4065 )),
4066 ));
4067 }
4068
4069 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4070 #[expect(
4071 clippy::too_many_lines,
4072 reason = "one lifecycle test proves successful replay plus pre-page and post-page source invalidation without sharing progress state across tests"
4073 )]
4074 #[test]
4075 fn journaled_job_advance_is_idempotent_and_revision_checked_on_both_sides() {
4076 let session = initialize_journaled();
4077 let proof = session
4078 .capture_read_set_revision_proof(&[ENTITY_NAME])
4079 .expect("journaled source proof should capture");
4080 let job_id =
4081 ResumableJobId::try_from_bytes([71; 32]).expect("nonzero job identity should admit");
4082 session
4083 .start_resumable_job(job_id, proof, vec![0])
4084 .expect("journaled job should start outside its protected source revision");
4085 let request = ResumableJobAdvanceRequest::new(
4086 job_id,
4087 0,
4088 ResumableJobIdempotencyKey::new("page-0")
4089 .expect("bounded idempotency key should admit"),
4090 );
4091 let calls = Cell::new(0_u8);
4092 let receipt = session
4093 .compare_proof_and_advance(&request, |state| {
4094 calls.set(calls.get() + 1);
4095 assert_eq!(state.application_state, vec![0]);
4096 Ok::<_, ()>(
4097 ResumableJobAdvance::new(Some("cursor-1".to_string()), vec![1], vec![9])
4098 .expect("bounded application advance should admit"),
4099 )
4100 })
4101 .expect("unchanged source should advance exactly once");
4102 assert_eq!(calls.get(), 1);
4103 assert_eq!(receipt.status, ResumableJobAdvanceStatus::Advanced);
4104 assert_eq!(receipt.committed_sequence, 1);
4105
4106 let replay = session
4107 .compare_proof_and_advance::<()>(&request, |_| {
4108 panic!("lost-response replay must not execute application work")
4109 })
4110 .expect("same request identity should return its persisted receipt");
4111 assert_eq!(replay, receipt);
4112 let retained = session
4113 .resumable_job_state(job_id)
4114 .expect("advanced state should remain durable");
4115 assert_eq!(retained.sequence, 1);
4116 assert_eq!(retained.application_state, vec![1]);
4117
4118 let _ = insert_exact_key_fixture(&session, 51);
4119 let pre_change_request = ResumableJobAdvanceRequest::new(
4120 job_id,
4121 1,
4122 ResumableJobIdempotencyKey::new("page-1")
4123 .expect("bounded idempotency key should admit"),
4124 );
4125 let pre_change_calls = Cell::new(0_u8);
4126 let invalidated = session
4127 .compare_proof_and_advance::<()>(&pre_change_request, |_| {
4128 pre_change_calls.set(pre_change_calls.get() + 1);
4129 unreachable!("pre-page proof failure must reject before application work")
4130 })
4131 .expect("source drift should persist one replayable invalidation receipt");
4132 assert_eq!(pre_change_calls.get(), 0);
4133 assert_eq!(invalidated.status, ResumableJobAdvanceStatus::Invalidated);
4134 let invalidated_state = session
4135 .resumable_job_state(job_id)
4136 .expect("invalidated job should remain inspectable");
4137 assert_eq!(invalidated_state.status, ResumableJobStatus::Invalidated);
4138 assert_eq!(invalidated_state.continuation, None);
4139 assert_eq!(invalidated_state.application_state, vec![1]);
4140 assert_eq!(
4141 session
4142 .compare_proof_and_advance::<()>(&pre_change_request, |_| {
4143 panic!("invalidation replay must not execute application work")
4144 })
4145 .expect("lost invalidation reply should replay exactly"),
4146 invalidated,
4147 );
4148
4149 let post_proof = session
4150 .capture_read_set_revision_proof(&[ENTITY_NAME])
4151 .expect("post-change journaled proof should capture");
4152 let post_job_id = ResumableJobId::try_from_bytes([72; 32])
4153 .expect("nonzero post-change job identity should admit");
4154 session
4155 .start_resumable_job(post_job_id, post_proof, vec![7])
4156 .expect("post-change journaled job should start");
4157 let post_request = ResumableJobAdvanceRequest::new(
4158 post_job_id,
4159 0,
4160 ResumableJobIdempotencyKey::new("post-page-0")
4161 .expect("bounded idempotency key should admit"),
4162 );
4163 let post_receipt = session
4164 .compare_proof_and_advance::<()>(&post_request, |_| {
4165 let _ = insert_exact_key_fixture(&session, 52);
4166 Ok(ResumableJobAdvance::new(None, vec![8], vec![10])
4167 .expect("bounded post-change candidate should admit"))
4168 })
4169 .expect("post-page drift should discard the candidate and persist invalidation");
4170 assert_eq!(post_receipt.status, ResumableJobAdvanceStatus::Invalidated);
4171 let post_state = session
4172 .resumable_job_state(post_job_id)
4173 .expect("post-page invalidation should remain inspectable");
4174 assert_eq!(post_state.status, ResumableJobStatus::Invalidated);
4175 assert_eq!(post_state.application_state, vec![7]);
4176 session
4177 .acknowledge_resumable_job(post_job_id, post_state.sequence)
4178 .expect("terminal job acknowledgement should remove retained progress");
4179 session
4180 .acknowledge_resumable_job(post_job_id, post_state.sequence)
4181 .expect("lost acknowledgement reply should be safely replayable");
4182 assert_eq!(
4183 session.resumable_job_state(post_job_id),
4184 Err(ResumableJobError::NotFound),
4185 );
4186
4187 let completed_job_id = ResumableJobId::try_from_bytes([74; 32])
4188 .expect("nonzero completed job identity should admit");
4189 let completed_proof = session
4190 .capture_read_set_revision_proof(&[ENTITY_NAME])
4191 .expect("completed-job source proof should capture");
4192 session
4193 .start_resumable_job(completed_job_id, completed_proof, Vec::new())
4194 .expect("completed-job fixture should start");
4195 let completed_request = ResumableJobAdvanceRequest::new(
4196 completed_job_id,
4197 0,
4198 ResumableJobIdempotencyKey::new("complete")
4199 .expect("bounded completion key should admit"),
4200 );
4201 let completed_receipt = session
4202 .compare_proof_and_advance::<()>(&completed_request, |_| {
4203 Ok(ResumableJobAdvance::new(None, vec![99], vec![100])
4204 .expect("bounded terminal advance should admit"))
4205 })
4206 .expect("null continuation should commit terminal completion");
4207 let completed_state = session
4208 .resumable_job_state(completed_job_id)
4209 .expect("completed state should remain replayable before acknowledgement");
4210 assert_eq!(completed_state.status, ResumableJobStatus::Completed);
4211 assert_eq!(
4212 session
4213 .compare_proof_and_advance::<()>(&completed_request, |_| {
4214 panic!("completed request replay must not execute application work")
4215 })
4216 .expect("completed request should replay until acknowledgement"),
4217 completed_receipt,
4218 );
4219 let after_completion = ResumableJobAdvanceRequest::new(
4220 completed_job_id,
4221 1,
4222 ResumableJobIdempotencyKey::new("after-complete")
4223 .expect("bounded post-completion key should admit"),
4224 );
4225 assert!(matches!(
4226 session.compare_proof_and_advance::<()>(&after_completion, |_| {
4227 panic!("completed jobs cannot execute another page")
4228 }),
4229 Err(CompareProofAndAdvanceError::Protocol(
4230 ResumableJobError::Completed
4231 )),
4232 ));
4233 session
4234 .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
4235 .expect("completed job should acknowledge and free capacity");
4236 session
4237 .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
4238 .expect("completion acknowledgement should be idempotent");
4239
4240 let stale_job_id = ResumableJobId::try_from_bytes([73; 32])
4241 .expect("nonzero stale-sequence job identity should admit");
4242 let stale_proof = session
4243 .capture_read_set_revision_proof(&[ENTITY_NAME])
4244 .expect("stale-sequence source proof should capture");
4245 session
4246 .start_resumable_job(stale_job_id, stale_proof, Vec::new())
4247 .expect("stale-sequence job should start");
4248 let stale_request = ResumableJobAdvanceRequest::new(
4249 stale_job_id,
4250 4,
4251 ResumableJobIdempotencyKey::new("stale").expect("bounded idempotency key should admit"),
4252 );
4253 assert!(matches!(
4254 session.compare_proof_and_advance::<()>(&stale_request, |_| {
4255 panic!("stale sequence must reject before application work")
4256 }),
4257 Err(CompareProofAndAdvanceError::Protocol(
4258 ResumableJobError::StaleSequence {
4259 expected: 4,
4260 actual: 0,
4261 }
4262 )),
4263 ));
4264 assert_eq!(
4265 session.acknowledge_resumable_job(stale_job_id, 0),
4266 Err(ResumableJobError::NotTerminal),
4267 );
4268 }
4269
4270 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4271 #[test]
4272 fn exact_key_batch_uses_typed_hard_execution_budget() {
4273 let session = initialize();
4274 let binding = exact_key_binding(&session);
4275 let budget =
4276 HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
4277 let error = session
4278 .execute_exact_key_batch_with_hard_budget_for_tests(&binding, &[u64::MAX], &budget)
4279 .expect_err("zero query budget should reject the exact-key route");
4280
4281 assert!(matches!(
4282 error.diagnostic().detail(),
4283 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4284 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
4285 })
4286 ));
4287 let facts = error.diagnostic_facts();
4288 assert_eq!(
4289 &facts[..5],
4290 &[
4291 (
4292 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
4293 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions.raw(),
4294 ),
4295 (icydb_diagnostic_code::DiagnosticFactTag::Limit, 0),
4296 (icydb_diagnostic_code::DiagnosticFactTag::Actual, 1),
4297 (
4298 icydb_diagnostic_code::DiagnosticFactTag::ExecutionBudgetScope,
4299 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution.raw(),
4300 ),
4301 (
4302 icydb_diagnostic_code::DiagnosticFactTag::ExecutionLane,
4303 icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead.raw(),
4304 ),
4305 ],
4306 );
4307 assert_eq!(
4308 facts[5].0,
4309 icydb_diagnostic_code::DiagnosticFactTag::QueryShapeFingerprintPrefix,
4310 );
4311 assert_ne!(facts[5].1, 0);
4312 }
4313
4314 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4315 fn assert_planned_query_exhausts(
4316 session: &DbSession<TestCanister>,
4317 query: &crate::db::DynamicQuery,
4318 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
4319 ) {
4320 let budget = HardExecutionBudget::uniform_for_tests(
4321 u64::MAX,
4322 HardExecutionFailureHeadroom::new(500, 256),
4323 )
4324 .with_limit_for_tests(resource, 0);
4325 let context = HardExecutionContext::new(
4326 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
4327 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
4328 0x7068_7973_6963_616c,
4329 );
4330 let error = with_query_execution_budget_for_tests(budget, context, || {
4331 session.execute_trusted_live_page(query, None)
4332 })
4333 .expect_err("the injected zero resource allowance should reject planned execution");
4334
4335 assert!(matches!(
4336 error.diagnostic().detail(),
4337 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4338 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
4339 })
4340 ));
4341 assert_eq!(
4342 error.diagnostic_facts()[0],
4343 (
4344 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
4345 resource.raw(),
4346 ),
4347 );
4348 }
4349
4350 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4351 fn assert_grouped_query_exhausts(
4352 session: &DbSession<TestCanister>,
4353 query: &crate::db::DynamicQuery,
4354 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
4355 ) {
4356 let budget = HardExecutionBudget::uniform_for_tests(
4357 u64::MAX,
4358 HardExecutionFailureHeadroom::new(500, 256),
4359 )
4360 .with_limit_for_tests(resource, 0);
4361 let context = HardExecutionContext::new(
4362 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
4363 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
4364 0x6772_6f75_7065_642d,
4365 );
4366 let error = with_query_execution_budget_for_tests(budget, context, || {
4367 session.execute_trusted_dynamic_grouped_query(query)
4368 })
4369 .expect_err("the injected zero resource allowance should reject grouped execution");
4370
4371 assert!(matches!(
4372 error.diagnostic().detail(),
4373 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4374 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
4375 })
4376 ));
4377 assert_eq!(
4378 error.diagnostic_facts()[0],
4379 (
4380 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
4381 resource.raw(),
4382 ),
4383 );
4384 }
4385
4386 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4387 fn assert_sql_query_exhausts(
4388 session: &DbSession<TestCanister>,
4389 sql: &str,
4390 resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
4391 ) {
4392 let budget = HardExecutionBudget::uniform_for_tests(
4393 u64::MAX,
4394 HardExecutionFailureHeadroom::new(500, 256),
4395 )
4396 .with_limit_for_tests(resource, 0);
4397 let context = HardExecutionContext::new(
4398 icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
4399 icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
4400 0x7371_6c2d_736f_7274,
4401 );
4402 let error = with_query_execution_budget_for_tests(budget, context, || {
4403 session.execute_trusted_sql_query(sql)
4404 })
4405 .expect_err("the injected zero resource allowance should reject SQL execution");
4406
4407 assert!(matches!(
4408 error.diagnostic().detail(),
4409 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4410 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
4411 })
4412 ));
4413 assert_eq!(
4414 error.diagnostic_facts()[0],
4415 (
4416 icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
4417 resource.raw(),
4418 ),
4419 );
4420 }
4421
4422 #[cfg(all(feature = "sql", feature = "diagnostics"))]
4423 #[test]
4424 fn planned_read_routes_share_physical_resource_accounting() {
4425 let session = initialize();
4426 let first = insert_exact_key_fixture(&session, 41);
4427 insert_exact_key_fixture(&session, 42);
4428
4429 let fallback = crate::db::DynamicQuery::new(ENTITY_NAME)
4430 .filter(crate::db::FieldRef::new("id").eq(first))
4431 .select(["id", "payload"])
4432 .order_by(crate::db::asc("id"))
4433 .limit(1);
4434 assert_eq!(
4435 session
4436 .execute_trusted_live_page(&fallback, None)
4437 .expect("bounded fallback execution should preserve its result")
4438 .row_count,
4439 1,
4440 );
4441 assert_planned_query_exhausts(
4442 &session,
4443 &fallback,
4444 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::RowsVisited,
4445 );
4446
4447 let covering = crate::db::DynamicQuery::new(ENTITY_NAME)
4448 .filter(crate::db::FieldRef::new("payload").eq(41_u64))
4449 .select(["payload"])
4450 .order_by(crate::db::asc("payload"))
4451 .limit(1);
4452 assert_eq!(
4453 session
4454 .execute_trusted_live_page(&covering, None)
4455 .expect("bounded covering execution should preserve its result")
4456 .row_count,
4457 1,
4458 );
4459 assert_planned_query_exhausts(
4460 &session,
4461 &covering,
4462 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
4463 );
4464
4465 let residual = crate::db::DynamicQuery::new(ENTITY_NAME)
4466 .filter(crate::db::FieldRef::new("payload").eq_field("id"))
4467 .select(["id"])
4468 .order_by(crate::db::asc("id"))
4469 .limit(1);
4470 assert_eq!(
4471 session
4472 .execute_trusted_live_page(&residual, None)
4473 .expect("bounded residual execution should preserve its result")
4474 .row_count,
4475 0,
4476 );
4477 assert_planned_query_exhausts(
4478 &session,
4479 &residual,
4480 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::PredicateExpressionSteps,
4481 );
4482
4483 assert_planned_query_exhausts(
4484 &session,
4485 &fallback,
4486 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::ResultBytes,
4487 );
4488
4489 let grouped = crate::db::DynamicQuery::new(ENTITY_NAME)
4490 .group_by("payload")
4491 .aggregate(crate::db::count())
4492 .order_by(crate::db::asc("payload"))
4493 .grouped_limits(10, 16 * 1_024)
4494 .limit(1);
4495 let grouped_result = session
4496 .execute_trusted_dynamic_grouped_query(&grouped)
4497 .expect("bounded grouped execution should preserve its result");
4498 assert_eq!(grouped_result.row_count, 1);
4499 assert!(grouped_result.next_cursor.is_some());
4500 assert_grouped_query_exhausts(
4501 &session,
4502 &grouped,
4503 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::GroupDistinctEntries,
4504 );
4505 assert_grouped_query_exhausts(
4506 &session,
4507 &grouped,
4508 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::CursorSteps,
4509 );
4510
4511 assert_sql_query_exhausts(
4512 &session,
4513 "SELECT payload, COUNT(*) AS row_count FROM IdentityRow \
4514 GROUP BY payload ORDER BY row_count DESC, payload ASC LIMIT 1",
4515 icydb_diagnostic_code::DiagnosticExecutionBudgetResource::SortEntries,
4516 );
4517 }
4518
4519 fn assert_dynamic_payload<C: CanisterKind>(
4520 session: &DbSession<C>,
4521 key: u64,
4522 expected_payload: u64,
4523 ) {
4524 let unchanged = session
4525 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
4526 entity: ENTITY_NAME.to_string(),
4527 key: InputValue::Nat64(key),
4528 patch: dynamic_payload_patch(expected_payload),
4529 })
4530 .expect("the expected row should remain readable through a no-op update");
4531 assert_eq!(unchanged.affected_rows, 0);
4532 assert_eq!(
4533 unchanged.rows,
4534 vec![expected_dynamic_row(key, expected_payload)],
4535 );
4536 }
4537
4538 fn assert_exact_batch_backlog_pressure(
4539 pressure: &InternalError,
4540 before: JournalTailControl,
4541 next_sequence: u64,
4542 ) {
4543 assert_eq!(
4544 pressure.diagnostic().error_code(),
4545 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_CONVERGENCE_BACKLOG_PRESSURE,
4546 );
4547 assert_eq!(
4548 pressure.diagnostic_facts(),
4549 vec![
4550 (
4551 icydb_diagnostic_code::DiagnosticFactTag::BacklogResource,
4552 icydb_diagnostic_code::DiagnosticBacklogResource::Batches.raw(),
4553 ),
4554 (icydb_diagnostic_code::DiagnosticFactTag::CurrentCount, 38),
4555 (icydb_diagnostic_code::DiagnosticFactTag::ProposedCount, 1),
4556 (icydb_diagnostic_code::DiagnosticFactTag::Limit, 38),
4557 ],
4558 );
4559 assert_eq!(
4560 crate::db::commit::next_database_commit_sequence()
4561 .expect("pressure must leave the database sequence readable"),
4562 next_sequence,
4563 );
4564 assert!(matches!(
4565 crate::db::commit::observe_commit_control()
4566 .expect("pressure must leave commit control observable"),
4567 crate::db::commit::CommitControlObservation::Present {
4568 marker_present: false,
4569 ..
4570 },
4571 ));
4572 assert_eq!(
4573 JOURNALED_TAIL_STORE.with(|tail| {
4574 tail.borrow()
4575 .current_tail_control()
4576 .expect("pressure must preserve the exact tail control")
4577 }),
4578 before,
4579 );
4580 }
4581
4582 fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
4583 values
4584 .iter()
4585 .map(|value| {
4586 AcceptedStructuralMutation::save(
4587 MutationMode::Insert,
4588 AcceptedStructuralMutationTarget::ResolveFromAfterImage,
4589 payload_patch(*value),
4590 )
4591 })
4592 .collect()
4593 }
4594
4595 fn atomic_progress_fixture(
4596 identity_byte: u8,
4597 ) -> (
4598 MutationJobRecord,
4599 MutationJobRecord,
4600 MutationProgressRecordOp,
4601 ) {
4602 let job_id = MutationJobId::try_from_bytes([identity_byte; 32])
4603 .expect("nonzero atomic progress job id should admit");
4604 let before = MutationJobRecord::new(job_id, vec![1, identity_byte], vec![2])
4605 .expect("atomic progress predecessor should admit");
4606 let request = MutationJobAdvanceRequest::new(
4607 job_id,
4608 0,
4609 MutationJobIdempotencyKey::new(format!("atomic-{identity_byte}"))
4610 .expect("atomic progress replay key should admit"),
4611 );
4612 let (after, _) = before
4613 .apply_transition(
4614 &request,
4615 MutationJobTransition::new(
4616 MutationJobStatus::Active,
4617 MutationJobPhase::Forward,
4618 vec![3],
4619 1,
4620 1,
4621 0,
4622 ),
4623 )
4624 .expect("atomic progress successor should admit");
4625 let operation = MutationProgressRecordOp::replace(&before, &after)
4626 .expect("atomic progress replacement should admit");
4627 (before, after, operation)
4628 }
4629
4630 fn assert_identity_boundary(error: &InternalError) {
4631 assert_eq!(error.class(), ErrorClass::Unsupported);
4632 assert_eq!(error.origin(), ErrorOrigin::Identity);
4633 }
4634
4635 #[test]
4636 fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
4637 let generated = insert_key_exists_after_generation(true);
4638 assert_eq!(generated.class(), ErrorClass::Corruption);
4639 assert_eq!(generated.origin(), ErrorOrigin::Identity);
4640
4641 let ordinary = insert_key_exists_after_generation(false);
4642 assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
4643 }
4644
4645 #[cfg(target_pointer_width = "64")]
4646 #[test]
4647 fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
4648 let error = checked_pre_key_candidate_count(
4649 usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
4650 )
4651 .expect_err("candidate counts beyond u32 must reject");
4652 assert_identity_boundary(&error);
4653 }
4654
4655 #[test]
4656 #[expect(
4657 clippy::too_many_lines,
4658 reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
4659 )]
4660 fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
4661 let session = initialize();
4662 let seeded = session
4663 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
4664 .expect("seed rows should commit");
4665 assert_eq!(seeded.affected_rows, 1);
4666
4667 let split = session
4668 .execute_trusted_dynamic_mutation_batch(vec![
4669 DynamicMutation::Update {
4670 entity: ENTITY_NAME.to_string(),
4671 key: InputValue::Nat64(1),
4672 patch: dynamic_payload_patch(60),
4673 },
4674 DynamicMutation::Insert {
4675 entity: ENTITY_NAME.to_string(),
4676 patch: dynamic_payload_patch(40),
4677 },
4678 ])
4679 .expect("one holding should split atomically");
4680 assert_eq!(split.affected_rows, 2);
4681 assert_eq!(
4682 split.rows,
4683 vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
4684 "split after-images must retain input order and exact quantity",
4685 );
4686
4687 let rejected_split = session
4688 .execute_trusted_dynamic_mutation_batch(vec![
4689 DynamicMutation::Update {
4690 entity: ENTITY_NAME.to_string(),
4691 key: InputValue::Nat64(1),
4692 patch: dynamic_payload_patch(50),
4693 },
4694 DynamicMutation::Insert {
4695 entity: ENTITY_NAME.to_string(),
4696 patch: DynamicStructuralPatch::new(Vec::new()),
4697 },
4698 ])
4699 .expect_err("an invalid split output must reject the staged source update");
4700 assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
4701 assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
4702 assert_eq!(
4703 rejected_split.diagnostic_facts(),
4704 vec![
4705 (
4706 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4707 ENTITY_TAG.value(),
4708 ),
4709 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
4710 (
4711 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4712 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4713 ),
4714 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
4715 ],
4716 );
4717 assert_dynamic_payload(&session, 1, 60);
4718 assert_dynamic_payload(&session, 2, 40);
4719
4720 let transfer = session
4721 .execute_trusted_dynamic_mutation_batch(vec![
4722 DynamicMutation::Update {
4723 entity: ENTITY_NAME.to_string(),
4724 key: InputValue::Nat64(1),
4725 patch: dynamic_payload_patch(70),
4726 },
4727 DynamicMutation::Update {
4728 entity: ENTITY_NAME.to_string(),
4729 key: InputValue::Nat64(2),
4730 patch: dynamic_payload_patch(30),
4731 },
4732 ])
4733 .expect("distinct transfer patches should share one atomic batch");
4734 assert_eq!(
4735 transfer.rows,
4736 vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
4737 "the transfer must preserve the exact total quantity",
4738 );
4739
4740 let merge = session
4741 .execute_trusted_dynamic_mutation_batch(vec![
4742 DynamicMutation::Delete {
4743 entity: ENTITY_NAME.to_string(),
4744 key: InputValue::Nat64(2),
4745 },
4746 DynamicMutation::Update {
4747 entity: ENTITY_NAME.to_string(),
4748 key: InputValue::Nat64(1),
4749 patch: dynamic_payload_patch(100),
4750 },
4751 ])
4752 .expect("two holdings should merge atomically");
4753 assert_eq!(
4754 merge.rows,
4755 vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
4756 "delete before-images and update after-images must retain input order",
4757 );
4758
4759 let resplit = session
4760 .execute_trusted_dynamic_mutation_batch(vec![
4761 DynamicMutation::Update {
4762 entity: ENTITY_NAME.to_string(),
4763 key: InputValue::Nat64(1),
4764 patch: dynamic_payload_patch(60),
4765 },
4766 DynamicMutation::Insert {
4767 entity: ENTITY_NAME.to_string(),
4768 patch: dynamic_payload_patch(40),
4769 },
4770 ])
4771 .expect("the merged holding should split again");
4772 assert_eq!(
4773 resplit.rows,
4774 vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
4775 );
4776
4777 let rejected_merge = session
4778 .execute_trusted_dynamic_mutation_batch(vec![
4779 DynamicMutation::Delete {
4780 entity: ENTITY_NAME.to_string(),
4781 key: InputValue::Nat64(3),
4782 },
4783 DynamicMutation::Update {
4784 entity: ENTITY_NAME.to_string(),
4785 key: InputValue::Nat64(99),
4786 patch: dynamic_payload_patch(100),
4787 },
4788 ])
4789 .expect_err("a late missing merge target must preserve the earlier staged delete");
4790 assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
4791 assert_dynamic_payload(&session, 1, 60);
4792 assert_dynamic_payload(&session, 3, 40);
4793
4794 SCHEMA_STORE.with(|store| {
4795 let cursor = store
4796 .borrow()
4797 .identity_statement_cursor(
4798 database_incarnation_id().expect("database incarnation should remain readable"),
4799 ENTITY_TAG,
4800 FieldId::new(1),
4801 &AcceptedFieldKind::Nat64,
4802 )
4803 .expect("mixed Identity state should remain readable");
4804 assert_eq!(cursor.expected_high_water(), 3);
4805 assert!(!cursor.has_allocations());
4806 });
4807 }
4808
4809 #[test]
4810 fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
4811 let session = initialize();
4812 session
4813 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
4814 .expect("the holding fixture should initialize");
4815
4816 let duplicate = session
4817 .execute_trusted_dynamic_mutation_batch(vec![
4818 DynamicMutation::Update {
4819 entity: ENTITY_NAME.to_string(),
4820 key: InputValue::Nat64(1),
4821 patch: dynamic_payload_patch(60),
4822 },
4823 DynamicMutation::Delete {
4824 entity: ENTITY_NAME.to_string(),
4825 key: InputValue::Nat64(1),
4826 },
4827 ])
4828 .expect_err("duplicate targets across operation kinds must reject");
4829 assert!(matches!(
4830 duplicate.diagnostic().detail(),
4831 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4832 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
4833 }),
4834 ));
4835 assert_eq!(
4836 duplicate.diagnostic_facts(),
4837 vec![
4838 (
4839 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4840 ENTITY_TAG.value(),
4841 ),
4842 (
4843 icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
4844 0,
4845 ),
4846 (
4847 icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
4848 1,
4849 ),
4850 ],
4851 );
4852 assert_dynamic_payload(&session, 1, 100);
4853 }
4854
4855 #[test]
4856 fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
4857 let session = initialize();
4858 let empty = session
4859 .execute_trusted_dynamic_mutation_batch(Vec::new())
4860 .expect_err("an empty public batch must reject");
4861 assert!(matches!(
4862 empty.diagnostic().detail(),
4863 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4864 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
4865 }),
4866 ));
4867 assert_eq!(
4868 empty.diagnostic_facts(),
4869 vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
4870 );
4871
4872 let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
4873 .map(|_| DynamicMutation::Delete {
4874 entity: ENTITY_NAME.to_string(),
4875 key: InputValue::Nat64(1),
4876 })
4877 .collect();
4878 let over_bound = session
4879 .execute_trusted_dynamic_mutation_batch(requests)
4880 .expect_err("operation cap plus one must reject before row resolution");
4881 assert!(matches!(
4882 over_bound.diagnostic().detail(),
4883 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4884 boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
4885 }),
4886 ));
4887 assert_eq!(
4888 over_bound.diagnostic_facts(),
4889 vec![
4890 (
4891 icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
4892 (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
4893 ),
4894 (
4895 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4896 MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
4897 ),
4898 ],
4899 );
4900 }
4901
4902 #[test]
4903 fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
4904 assert_eq!(
4905 structural_mutation_staged_charge([11, 13, 17])
4906 .expect("the writer-owned formula should sum all three row-image components"),
4907 41,
4908 );
4909 let mut exact = 0;
4910 add_structural_mutation_staged_bytes(
4911 &mut exact,
4912 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
4913 )
4914 .expect("the exact staged-byte boundary should admit");
4915 assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
4916
4917 let error = add_structural_mutation_staged_bytes(&mut exact, [1])
4918 .expect_err("one byte above the staged-byte boundary must reject");
4919 assert!(matches!(
4920 error.diagnostic().detail(),
4921 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4922 boundary:
4923 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
4924 }),
4925 ));
4926 assert_eq!(
4927 error.diagnostic_facts(),
4928 vec![
4929 (
4930 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4931 (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
4932 ),
4933 (
4934 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4935 MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
4936 ),
4937 ],
4938 );
4939
4940 let mut prefix = 0;
4941 assert_eq!(
4942 admit_structural_mutation_staged_charge(
4943 &mut prefix,
4944 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
4945 AcceptedStructuralMutationPacking::BoundedPrefix,
4946 )
4947 .expect("the exact prefix boundary should calculate"),
4948 AcceptedStructuralMutationStagedAdmission::Admitted,
4949 );
4950 assert_eq!(prefix, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
4951 assert_eq!(
4952 admit_structural_mutation_staged_charge(
4953 &mut prefix,
4954 [1],
4955 AcceptedStructuralMutationPacking::BoundedPrefix,
4956 )
4957 .expect("the next prefix candidate should calculate"),
4958 AcceptedStructuralMutationStagedAdmission::PageFull,
4959 );
4960 assert_eq!(prefix, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
4961
4962 let mut empty_prefix = 0;
4963 assert_eq!(
4964 admit_structural_mutation_staged_charge(
4965 &mut empty_prefix,
4966 [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1],
4967 AcceptedStructuralMutationPacking::BoundedPrefix,
4968 )
4969 .expect("one oversized candidate should classify without mutating the prefix"),
4970 AcceptedStructuralMutationStagedAdmission::CandidateExceedsPolicy,
4971 );
4972 assert_eq!(empty_prefix, 0);
4973
4974 validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
4975 .expect("the exact result-byte boundary should admit");
4976 let error = validate_structural_mutation_result_bytes(
4977 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
4978 )
4979 .expect_err("one byte above the result-byte boundary must reject");
4980 assert!(matches!(
4981 error.diagnostic().detail(),
4982 Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4983 boundary:
4984 icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
4985 }),
4986 ));
4987 assert_eq!(
4988 error.diagnostic_facts(),
4989 vec![
4990 (
4991 icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4992 (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
4993 ),
4994 (
4995 icydb_diagnostic_code::DiagnosticFactTag::Limit,
4996 MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
4997 ),
4998 ],
4999 );
5000 }
5001
5002 #[expect(
5003 clippy::too_many_lines,
5004 reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
5005 )]
5006 #[test]
5007 fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
5008 let session = initialize();
5009 let catalog = session
5010 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
5011 .expect("identity catalog should resolve");
5012 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
5013 .expect("identity row layout should build");
5014 let initial_description = session
5015 .try_describe_entity_by_name(ENTITY_NAME)
5016 .expect("accepted Identity description should resolve");
5017 assert_eq!(
5018 initial_description.entity_tag(),
5019 catalog.identity().entity_tag().value()
5020 );
5021 assert_eq!(
5022 initial_description.accepted_schema_fingerprint_method(),
5023 catalog.fingerprint_method_version()
5024 );
5025 assert_eq!(
5026 initial_description.accepted_schema_fingerprint(),
5027 catalog.fingerprint()
5028 );
5029 let initial_identity = initial_description
5030 .identity()
5031 .expect("accepted Identity policy should be described");
5032 assert_eq!(initial_identity.field(), "id");
5033 assert_eq!(initial_identity.generator(), "Identity::next");
5034 assert_eq!(initial_identity.accepted_kind(), "nat64");
5035 assert_eq!(initial_identity.minimum(), 1);
5036 assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
5037 assert_eq!(initial_identity.high_water(), 0);
5038 assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
5039 assert!(!initial_identity.exhausted());
5040
5041 let rejected = session
5042 .execute_accepted_structural_save_batch(
5043 &catalog,
5044 &descriptor,
5045 batch(&[1_000, 2_000]),
5046 Timestamp::from_millis(6),
5047 |_| Err::<(), _>(InternalError::executor_unsupported()),
5048 )
5049 .expect_err("a rejected precommit result must not publish its tentative range");
5050 assert_eq!(rejected.class(), ErrorClass::Unsupported);
5051 assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
5052
5053 let rows = session
5054 .execute_accepted_structural_save_batch(
5055 &catalog,
5056 &descriptor,
5057 batch(&[10, 20, 30]),
5058 Timestamp::from_millis(7),
5059 Ok,
5060 )
5061 .expect("one accepted batch should commit rows and one identity range");
5062 assert_eq!(
5063 rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
5064 vec![
5065 vec![Value::Nat64(1), Value::Nat64(10)],
5066 vec![Value::Nat64(2), Value::Nat64(20)],
5067 vec![Value::Nat64(3), Value::Nat64(30)],
5068 ],
5069 );
5070
5071 let dynamic = session
5072 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5073 entity: ENTITY_NAME.to_string(),
5074 patch: DynamicStructuralPatch::new(vec![(
5075 "payload".to_string(),
5076 DynamicWriteCell::Value(InputValue::Nat64(40)),
5077 )]),
5078 })
5079 .expect("dynamic omission should commit through shared Identity generation");
5080 assert_eq!(dynamic.affected_rows, 1);
5081
5082 for (request, operation) in [
5083 (
5084 DynamicMutation::Insert {
5085 entity: ENTITY_NAME.to_string(),
5086 patch: DynamicStructuralPatch::new(vec![
5087 (
5088 "id".to_string(),
5089 DynamicWriteCell::Value(InputValue::Nat64(41)),
5090 ),
5091 (
5092 "payload".to_string(),
5093 DynamicWriteCell::Value(InputValue::Nat64(42)),
5094 ),
5095 ]),
5096 },
5097 icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
5098 ),
5099 (
5100 DynamicMutation::Update {
5101 entity: ENTITY_NAME.to_string(),
5102 key: InputValue::Nat64(1),
5103 patch: DynamicStructuralPatch::new(vec![(
5104 "id".to_string(),
5105 DynamicWriteCell::Default,
5106 )]),
5107 },
5108 icydb_diagnostic_code::DiagnosticMutationOperation::Update,
5109 ),
5110 ] {
5111 let error = session
5112 .execute_trusted_dynamic_mutation(&request)
5113 .expect_err("structural Identity authorship and regeneration must reject");
5114 assert_eq!(error.class(), ErrorClass::Unsupported);
5115 assert_eq!(error.origin(), ErrorOrigin::Executor);
5116 assert_eq!(
5117 error.diagnostic_facts(),
5118 vec![
5119 (
5120 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
5121 ENTITY_TAG.value(),
5122 ),
5123 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
5124 (
5125 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
5126 operation.raw(),
5127 ),
5128 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
5129 ],
5130 );
5131 }
5132
5133 let binding = session
5134 .issue_typed_entity_binding(
5135 ENTITY_SOURCE,
5136 &[
5137 DynamicTypedFieldBindingRequest::new(
5138 ID_SOURCE.to_string(),
5139 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
5140 false,
5141 ),
5142 DynamicTypedFieldBindingRequest::new(
5143 PAYLOAD_SOURCE.to_string(),
5144 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
5145 false,
5146 ),
5147 ],
5148 )
5149 .expect("typed output should bind the Identity field");
5150 let typed_patch = binding
5151 .bind_write_fields(vec![(
5152 PAYLOAD_SOURCE.to_string(),
5153 DynamicWriteCell::Value(InputValue::Nat64(50)),
5154 )])
5155 .expect("typed payload should lower");
5156 let typed = session
5157 .execute_trusted_typed_mutation(
5158 &binding,
5159 &DynamicTypedMutation::Insert { patch: typed_patch },
5160 )
5161 .expect("typed omission should commit through shared Identity generation");
5162 assert_eq!(
5163 typed
5164 .expect("typed insert should return one mutation result")
5165 .affected_rows,
5166 1,
5167 );
5168 let explicit_typed_patch = binding
5169 .bind_write_fields(vec![
5170 (
5171 ID_SOURCE.to_string(),
5172 DynamicWriteCell::Value(InputValue::Nat64(51)),
5173 ),
5174 (
5175 PAYLOAD_SOURCE.to_string(),
5176 DynamicWriteCell::Value(InputValue::Nat64(52)),
5177 ),
5178 ])
5179 .expect("the low-level binding should retain exact authored intent");
5180 let explicit_typed_error = session
5181 .execute_trusted_typed_mutation(
5182 &binding,
5183 &DynamicTypedMutation::Insert {
5184 patch: explicit_typed_patch,
5185 },
5186 )
5187 .expect_err("typed Identity authorship must reject before allocation");
5188 assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
5189 assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
5190 assert_eq!(
5191 explicit_typed_error.diagnostic_facts(),
5192 vec![
5193 (
5194 icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
5195 ENTITY_TAG.value(),
5196 ),
5197 (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
5198 (
5199 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
5200 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
5201 ),
5202 (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
5203 ],
5204 );
5205
5206 let replace_error = session
5207 .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
5208 entity: ENTITY_NAME.to_string(),
5209 key: InputValue::Nat64(99),
5210 patch: DynamicStructuralPatch::new(vec![(
5211 "payload".to_string(),
5212 DynamicWriteCell::Value(InputValue::Nat64(60)),
5213 )]),
5214 })
5215 .expect_err("save-as-insert with a chosen Identity must reject");
5216 assert_eq!(replace_error.class(), ErrorClass::Unsupported);
5217 assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
5218
5219 #[cfg(feature = "sql")]
5220 {
5221 for sql in [
5222 "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
5223 "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
5224 ] {
5225 let _result = session
5226 .execute_trusted_sql_mutation(sql)
5227 .expect("SQL omission and DEFAULT should commit Identity generation");
5228 }
5229
5230 let error = session
5231 .execute_trusted_sql_mutation(
5232 "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
5233 )
5234 .expect_err("an explicit SQL Identity value must reject before allocation");
5235 let diagnostic = error.diagnostic();
5236 assert_eq!(
5237 diagnostic.code(),
5238 icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
5239 );
5240 assert!(matches!(
5241 diagnostic.detail(),
5242 Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
5243 boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
5244 }),
5245 ));
5246 }
5247
5248 let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
5249 assert_eq!(
5250 DATA_STORE.with(|store| store.borrow().len()),
5251 expected_committed
5252 );
5253 SCHEMA_STORE.with(|store| {
5254 let cursor = store
5255 .borrow()
5256 .identity_statement_cursor(
5257 database_incarnation_id().expect("database incarnation should remain readable"),
5258 ENTITY_TAG,
5259 FieldId::new(1),
5260 &AcceptedFieldKind::Nat64,
5261 )
5262 .expect("committed writes must leave active state readable");
5263 assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
5264 assert!(!cursor.has_allocations());
5265 });
5266 let committed_description = session
5267 .try_describe_entity_by_name(ENTITY_NAME)
5268 .expect("committed Identity description should resolve");
5269 let committed_identity = committed_description
5270 .identity()
5271 .expect("accepted Identity policy should remain described");
5272 assert_eq!(
5273 committed_identity.high_water(),
5274 u128::from(expected_committed),
5275 );
5276 assert_eq!(
5277 committed_identity.remaining(),
5278 u128::from(u64::MAX - expected_committed),
5279 );
5280 assert!(!committed_identity.exhausted());
5281 }
5282
5283 #[test]
5284 #[expect(
5285 clippy::too_many_lines,
5286 reason = "one ordered scenario proves target/progress atomicity, every interruption wake-up, state-only admission, and successful no-op wake-up behavior"
5287 )]
5288 fn mutation_progress_and_target_rows_recover_as_one_marker_transition() {
5289 let session = initialize_journaled();
5290 install_startup_recovery_wakeup(record_startup_wakeup);
5291 let catalog = session
5292 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
5293 .expect("journaled atomic-progress catalog should resolve");
5294 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
5295 .expect("journaled atomic-progress row layout should build");
5296
5297 for (ordinal, interruption) in [
5298 MutationCommitInterruption::MarkerPersisted,
5299 MutationCommitInterruption::JournalPublished,
5300 MutationCommitInterruption::RowsPublished,
5301 MutationCommitInterruption::ProgressReplaced,
5302 ]
5303 .into_iter()
5304 .enumerate()
5305 {
5306 let identity_byte = 31 + u8::try_from(ordinal).expect("small ordinal should fit");
5307 let (before, after, operation) = atomic_progress_fixture(identity_byte);
5308 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5309 match store.insert_mutation(&before)? {
5310 InsertMutationJobResult::Inserted => Ok(()),
5311 InsertMutationJobResult::Occupied(_) => {
5312 Err(crate::db::MutationJobError::IdentityConflict)
5313 }
5314 }
5315 })
5316 .expect("atomic predecessor should insert once");
5317
5318 let wakeups_before = STARTUP_WAKEUPS.with(Cell::get);
5319 interrupt_next_mutation_commit_for_tests(interruption);
5320 let interrupted = session.execute_accepted_structural_update_with_mutation_progress(
5321 &catalog,
5322 &descriptor,
5323 batch(&[700 + u64::try_from(ordinal).expect("small ordinal should fit")]),
5324 Timestamp::from_millis(17),
5325 operation,
5326 );
5327 assert!(
5328 interrupted.is_err(),
5329 "selected atomic boundary should interrupt"
5330 );
5331 assert_eq!(
5332 STARTUP_WAKEUPS.with(Cell::get),
5333 wakeups_before.saturating_add(1),
5334 "a normally returned retained-marker error must register its wake-up",
5335 );
5336
5337 forget_recovered_domain_for_tests(&session.db)
5338 .expect("interruption should reset volatile recovery ownership");
5339 let retained_before =
5340 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5341 store.load_mutation(before.state().job_id)
5342 })
5343 .expect("pre-driver progress should load");
5344 let row_count_before = JOURNALED_DATA_STORE.with(|store| store.borrow().len());
5345 let pending = session
5346 .db
5347 .ensure_recovered_state()
5348 .expect_err("ordinary admission must not drive retained-marker recovery");
5349 assert_eq!(
5350 pending.diagnostic().error_code(),
5351 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_DATABASE_STARTUP_RECOVERY_PENDING,
5352 );
5353 assert_eq!(
5354 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5355 store.load_mutation(before.state().job_id)
5356 })
5357 .expect("post-admission progress should load"),
5358 retained_before,
5359 );
5360 assert_eq!(
5361 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
5362 row_count_before,
5363 "state-only admission must not mutate target rows",
5364 );
5365 assert!(
5366 session
5367 .db
5368 .drive_startup_recovery_page()
5369 .expect("dedicated driver should finish target and progress together"),
5370 );
5371 let retained = with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5372 store.load_mutation(before.state().job_id)
5373 })
5374 .expect("recovered successor should load");
5375 assert_eq!(retained, after);
5376 assert_eq!(
5377 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
5378 u64::try_from(ordinal + 1).expect("small row count should fit"),
5379 );
5380 }
5381
5382 let (before, after, operation) = atomic_progress_fixture(39);
5383 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5384 match store.insert_mutation(&before)? {
5385 InsertMutationJobResult::Inserted => Ok(()),
5386 InsertMutationJobResult::Occupied(_) => {
5387 Err(crate::db::MutationJobError::IdentityConflict)
5388 }
5389 }
5390 })
5391 .expect("final predecessor should insert once");
5392 let wakeups_before_success = STARTUP_WAKEUPS.with(Cell::get);
5393 session
5394 .execute_accepted_structural_update_with_mutation_progress(
5395 &catalog,
5396 &descriptor,
5397 batch(&[799]),
5398 Timestamp::from_millis(18),
5399 operation,
5400 )
5401 .expect("uninterrupted atomic transition should clear its marker");
5402 assert_eq!(
5403 STARTUP_WAKEUPS.with(Cell::get),
5404 wakeups_before_success.saturating_add(1),
5405 "a successful retained commit must request online convergence",
5406 );
5407 let retained = with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5408 store.load_mutation(before.state().job_id)
5409 })
5410 .expect("final successor should load");
5411 assert_eq!(retained, after);
5412 forget_recovered_domain_for_tests(&session.db)
5413 .expect("post-clear recovery ownership should reset");
5414 let pending = session
5415 .db
5416 .ensure_recovered_state()
5417 .expect_err("an upgrade epoch must remain gated until its driver runs");
5418 assert_eq!(
5419 pending.diagnostic().error_code(),
5420 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_DATABASE_STARTUP_RECOVERY_PENDING,
5421 );
5422 assert!(
5423 session
5424 .db
5425 .drive_startup_recovery_page()
5426 .expect("post-clear driver recovery should fold the retained batch"),
5427 );
5428 }
5429
5430 #[test]
5431 fn mutation_progress_neither_side_mismatch_blocks_recovery() {
5432 let session = initialize_journaled();
5433 let catalog = session
5434 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
5435 .expect("journaled corruption catalog should resolve");
5436 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
5437 .expect("journaled corruption row layout should build");
5438 let (before, _after, operation) = atomic_progress_fixture(41);
5439 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5440 match store.insert_mutation(&before)? {
5441 InsertMutationJobResult::Inserted => Ok(()),
5442 InsertMutationJobResult::Occupied(_) => {
5443 Err(crate::db::MutationJobError::IdentityConflict)
5444 }
5445 }
5446 })
5447 .expect("corruption predecessor should insert once");
5448
5449 interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::MarkerPersisted);
5450 assert!(
5451 session
5452 .execute_accepted_structural_update_with_mutation_progress(
5453 &catalog,
5454 &descriptor,
5455 batch(&[811]),
5456 Timestamp::from_millis(19),
5457 operation,
5458 )
5459 .is_err(),
5460 "marker interruption should retain recovery authority",
5461 );
5462 let (unexpected, _) = before
5463 .apply_transition(
5464 &MutationJobAdvanceRequest::new(
5465 before.state().job_id,
5466 0,
5467 MutationJobIdempotencyKey::new("unexpected-third-state")
5468 .expect("unexpected replay key should admit"),
5469 ),
5470 MutationJobTransition::new(
5471 MutationJobStatus::Active,
5472 MutationJobPhase::Forward,
5473 vec![99],
5474 2,
5475 0,
5476 0,
5477 ),
5478 )
5479 .expect("unexpected but valid progress state should admit");
5480 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5481 store.replace_mutation(&unexpected)
5482 })
5483 .expect("test should install the neither-side state");
5484
5485 forget_recovered_domain_for_tests(&session.db)
5486 .expect("corrupt recovery ownership should reset");
5487 let error = session
5488 .db
5489 .drive_startup_recovery_page()
5490 .expect_err("neither-side progress must block recovery");
5491 assert_eq!(error.class(), ErrorClass::Corruption);
5492 assert_eq!(error.origin(), ErrorOrigin::Recovery);
5493 assert_eq!(
5494 with_mutation_progress_store::<JournaledTestCanister, _>(|store| {
5495 store.load_mutation(before.state().job_id)
5496 })
5497 .expect("unexpected state should remain inspectable to the test"),
5498 unexpected,
5499 );
5500 assert!(
5501 session.db.drive_startup_recovery_page().is_err(),
5502 "a retained corrupt marker must continue blocking database access",
5503 );
5504 }
5505
5506 #[test]
5507 #[expect(
5508 clippy::too_many_lines,
5509 reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
5510 )]
5511 fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
5512 let session = initialize_journaled();
5513 let catalog = session
5514 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
5515 .expect("journaled identity catalog should resolve");
5516 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
5517 .expect("journaled identity row layout should build");
5518
5519 for (ordinal, interruption) in [
5520 MutationCommitInterruption::MarkerPersisted,
5521 MutationCommitInterruption::JournalPublished,
5522 MutationCommitInterruption::RowsPublished,
5523 MutationCommitInterruption::StateMaterialized,
5524 ]
5525 .into_iter()
5526 .enumerate()
5527 {
5528 interrupt_next_mutation_commit_for_tests(interruption);
5529 let interrupted = session.execute_accepted_structural_save_batch(
5530 &catalog,
5531 &descriptor,
5532 batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
5533 Timestamp::from_millis(8),
5534 Ok,
5535 );
5536 assert!(
5537 interrupted.is_err(),
5538 "the selected durable boundary should interrupt",
5539 );
5540
5541 let Err(pending) = session.execute_accepted_structural_save_batch(
5542 &catalog,
5543 &descriptor,
5544 batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
5545 Timestamp::from_millis(9),
5546 Ok,
5547 ) else {
5548 panic!("ordinary mutation must not drive retained-marker recovery");
5549 };
5550 assert_eq!(
5551 pending.diagnostic().error_code(),
5552 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_DATABASE_STARTUP_RECOVERY_PENDING,
5553 );
5554 drive_journaled_recovery_to_completion(&session);
5555
5556 let committed = session
5557 .execute_accepted_structural_save_batch(
5558 &catalog,
5559 &descriptor,
5560 batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
5561 Timestamp::from_millis(9),
5562 Ok,
5563 )
5564 .expect("the next mutation must recover before allocating");
5565 let expected_high_water =
5566 u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
5567 assert_eq!(
5568 committed
5569 .into_iter()
5570 .map(|row| row.values)
5571 .collect::<Vec<_>>(),
5572 vec![vec![
5573 Value::Nat64(expected_high_water),
5574 Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
5575 ]],
5576 );
5577 assert_eq!(
5578 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
5579 expected_high_water,
5580 );
5581 JOURNALED_SCHEMA_STORE.with(|store| {
5582 let cursor = store
5583 .borrow()
5584 .identity_statement_cursor(
5585 database_incarnation_id()
5586 .expect("database incarnation should remain readable"),
5587 ENTITY_TAG,
5588 FieldId::new(1),
5589 &AcceptedFieldKind::Nat64,
5590 )
5591 .expect("guarded recovery must leave quiescent active state");
5592 assert_eq!(
5593 cursor.expected_high_water(),
5594 u128::from(expected_high_water),
5595 );
5596 assert!(!cursor.has_allocations());
5597 });
5598 }
5599
5600 for (ordinal, (interruption, deleted_key)) in [
5601 (MutationCommitInterruption::MarkerPersisted, 2),
5602 (MutationCommitInterruption::JournalPublished, 4),
5603 (MutationCommitInterruption::RowPrefixPublished, 6),
5604 (MutationCommitInterruption::RowsPublished, 8),
5605 (MutationCommitInterruption::StateMaterialized, 7),
5606 ]
5607 .into_iter()
5608 .enumerate()
5609 {
5610 let expected_payload =
5611 501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
5612 interrupt_next_mutation_commit_for_tests(interruption);
5613 let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
5614 DynamicMutation::Update {
5615 entity: ENTITY_NAME.to_string(),
5616 key: InputValue::Nat64(1),
5617 patch: dynamic_payload_patch(expected_payload),
5618 },
5619 DynamicMutation::Delete {
5620 entity: ENTITY_NAME.to_string(),
5621 key: InputValue::Nat64(deleted_key),
5622 },
5623 ]);
5624 assert!(
5625 interrupted.is_err(),
5626 "the selected caller-key mixed publication boundary should interrupt",
5627 );
5628 let pending = session
5629 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
5630 entity: ENTITY_NAME.to_string(),
5631 key: InputValue::Nat64(1),
5632 patch: dynamic_payload_patch(expected_payload),
5633 })
5634 .expect_err("ordinary update must not drive retained-marker recovery");
5635 assert_eq!(
5636 pending.diagnostic().error_code(),
5637 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_DATABASE_STARTUP_RECOVERY_PENDING,
5638 );
5639 drive_journaled_recovery_to_completion(&session);
5640 let recovered_update = session
5641 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
5642 entity: ENTITY_NAME.to_string(),
5643 key: InputValue::Nat64(1),
5644 patch: dynamic_payload_patch(expected_payload),
5645 })
5646 .expect("guarded reentry should complete the marker-authorized mixed batch");
5647 assert_eq!(
5648 recovered_update.affected_rows, 0,
5649 "the recovered update must already expose its admitted final image",
5650 );
5651 let recovered_delete = session
5652 .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
5653 entity: ENTITY_NAME.to_string(),
5654 key: InputValue::Nat64(deleted_key),
5655 })
5656 .expect_err("the recovered delete must already be materialized");
5657 assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
5658 JOURNALED_SCHEMA_STORE.with(|store| {
5659 let cursor = store
5660 .borrow()
5661 .identity_statement_cursor(
5662 database_incarnation_id()
5663 .expect("database incarnation should remain readable"),
5664 ENTITY_TAG,
5665 FieldId::new(1),
5666 &AcceptedFieldKind::Nat64,
5667 )
5668 .expect("caller-key recovery must preserve active Identity state");
5669 assert_eq!(cursor.expected_high_water(), 8);
5670 assert!(!cursor.has_allocations());
5671 });
5672 }
5673
5674 forget_recovered_domain_for_tests(&session.db)
5675 .expect("the final journal tail should remain recoverable");
5676 session
5677 .db
5678 .drive_startup_recovery_page()
5679 .expect("derived rebuild must not allocate another identity");
5680
5681 let data_generation = JOURNALED_DATA_STORE.with(|store| store.borrow().generation());
5682 let index_generation = JOURNALED_INDEX_STORE.with(|store| store.borrow().generation());
5683 let data_len = JOURNALED_DATA_STORE.with(|store| store.borrow().len());
5684 let index_len = JOURNALED_INDEX_STORE.with(|store| store.borrow().len());
5685 forget_recovered_domain_for_tests(&session.db)
5686 .expect("an empty-tail upgrade should reset recovery ownership");
5687 session
5688 .db
5689 .drive_startup_recovery_page()
5690 .expect("an empty-tail upgrade should admit without rebuilding stored rows or indexes");
5691 assert_eq!(
5692 JOURNALED_DATA_STORE.with(|store| store.borrow().generation()),
5693 data_generation
5694 .checked_add(1)
5695 .expect("test generation should advance once"),
5696 "empty-tail recovery must reset the disposable row projection exactly once",
5697 );
5698 assert_eq!(
5699 JOURNALED_INDEX_STORE.with(|store| store.borrow().generation()),
5700 index_generation
5701 .checked_add(1)
5702 .expect("test generation should advance once"),
5703 "empty-tail recovery must reset the disposable index projection exactly once",
5704 );
5705 assert_eq!(
5706 JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
5707 data_len,
5708 "empty-tail recovery must not rebuild or remove authoritative rows",
5709 );
5710 assert_eq!(
5711 JOURNALED_INDEX_STORE.with(|store| store.borrow().len()),
5712 index_len,
5713 "empty-tail recovery must not clear or rebuild canonical secondary indexes",
5714 );
5715
5716 let quick = execute_quick_integrity(
5717 &session.db,
5718 catalog.inspection_plan(),
5719 catalog.runtime_root_identity().database_incarnation(),
5720 )
5721 .expect("quiescent Identity control inventory should be inspectable");
5722 assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
5723 let row_page = execute_row_integrity_page(
5724 &session.db,
5725 catalog.inspection_plan(),
5726 PhysicalUnitCheckpoint::BeforeFirst,
5727 RowInspectionLimits::standard(),
5728 )
5729 .expect("Identity rows should remain within committed high-water");
5730 assert!(row_page.exhausted());
5731 assert!(row_page.findings().is_empty());
5732
5733 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
5734 assert!(
5735 JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
5736 "derived index rebuild should restore witnesses without allocating identities",
5737 );
5738 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
5739 JOURNALED_SCHEMA_STORE.with(|store| {
5740 let cursor = store
5741 .borrow()
5742 .identity_statement_cursor(
5743 database_incarnation_id().expect("database incarnation should remain readable"),
5744 ENTITY_TAG,
5745 FieldId::new(1),
5746 &AcceptedFieldKind::Nat64,
5747 )
5748 .expect("folded identity state should reopen without allocating");
5749 assert_eq!(cursor.expected_high_water(), 8);
5750 assert!(!cursor.has_allocations());
5751 });
5752 }
5753
5754 #[test]
5755 fn journaled_online_convergence_drains_the_full_backlog_in_complete_batch_callbacks_without_reallocating_ids()
5756 {
5757 const SUBMISSION: &str = "generated/8899aabbccddeeff";
5758 let session = initialize_journaled();
5759 let catalog = session
5760 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
5761 .expect("journaled identity catalog should resolve");
5762 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
5763 .expect("journaled identity row layout should build");
5764
5765 for payload in 0_u64..38 {
5766 session
5767 .execute_accepted_structural_save_batch(
5768 &catalog,
5769 &descriptor,
5770 batch(&[payload]),
5771 Timestamp::from_millis(8),
5772 Ok,
5773 )
5774 .unwrap_or_else(|error| {
5775 panic!("journaled identity fixture row {payload} should commit: {error:?}")
5776 });
5777 }
5778
5779 let before = JOURNALED_TAIL_STORE.with(|tail| {
5780 tail.borrow()
5781 .current_tail_control()
5782 .expect("online backlog control should remain valid")
5783 });
5784 assert_eq!(before.batch_count(), 38);
5785 let next_sequence = crate::db::commit::next_database_commit_sequence()
5786 .expect("database sequence preview should remain readable");
5787 let Err(pressure) = session.execute_accepted_structural_save_batch(
5788 &catalog,
5789 &descriptor,
5790 batch(&[38]),
5791 Timestamp::from_millis(8),
5792 Ok,
5793 ) else {
5794 panic!("the exact cumulative batch ceiling should reject one more batch")
5795 };
5796 assert_exact_batch_backlog_pressure(&pressure, before, next_sequence);
5797
5798 for folded_batches in 1..=38 {
5799 let complete = session
5800 .db
5801 .drive_startup_recovery_page()
5802 .expect("online complete-batch callback should commit");
5803 assert_eq!(complete, folded_batches == 38);
5804 }
5805
5806 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
5807 session
5808 .execute_accepted_structural_save_batch(
5809 &catalog,
5810 &descriptor,
5811 batch(&[38]),
5812 Timestamp::from_millis(8),
5813 Ok,
5814 )
5815 .expect("drain should make the rejected mutation retryable");
5816 assert!(
5817 session
5818 .db
5819 .drive_startup_recovery_page()
5820 .expect("the retry tail should converge"),
5821 );
5822
5823 assert_eq!(
5824 drive_generated_startup_recovery_page(&session, &JOURNALED_STORE_REGISTRY, SUBMISSION,)
5825 .expect("online convergence should commit"),
5826 GeneratedStartupDriverStep::Terminal,
5827 "the quiescent generated driver should stop",
5828 );
5829
5830 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 39);
5831 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
5832 assert_dynamic_payload(&session, 1, 0);
5833 assert_dynamic_payload(&session, 39, 38);
5834 JOURNALED_SCHEMA_STORE.with(|store| {
5835 let cursor = store
5836 .borrow()
5837 .identity_statement_cursor(
5838 database_incarnation_id().expect("database incarnation should remain readable"),
5839 ENTITY_TAG,
5840 FieldId::new(1),
5841 &AcceptedFieldKind::Nat64,
5842 )
5843 .expect("online convergence must preserve active Identity state");
5844 assert_eq!(cursor.expected_high_water(), 39);
5845 assert!(!cursor.has_allocations());
5846 });
5847 }
5848
5849 #[test]
5850 fn journaled_online_convergence_reconstructs_same_key_batches_from_canonical_predecessors() {
5851 let session = initialize_journaled();
5852 session
5853 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5854 entity: ENTITY_NAME.to_string(),
5855 patch: dynamic_payload_patch(10),
5856 })
5857 .expect("the initial positioned row should commit");
5858 for payload in [20, 30] {
5859 session
5860 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
5861 entity: ENTITY_NAME.to_string(),
5862 key: InputValue::Nat64(1),
5863 patch: dynamic_payload_patch(payload),
5864 })
5865 .unwrap_or_else(|error| {
5866 panic!("the positioned same-key update should commit: {error:?}")
5867 });
5868 }
5869
5870 assert_dynamic_payload(&session, 1, 30);
5871 assert_eq!(
5872 JOURNALED_INDEX_STORE.with(|store| store.borrow().len()),
5873 1,
5874 "the newest live index effect should hide every predecessor",
5875 );
5876 for folded_batches in 1..=3 {
5877 let complete = session
5878 .db
5879 .drive_startup_recovery_page()
5880 .expect("the positioned same-key batch should converge");
5881 assert_eq!(complete, folded_batches == 3);
5882 }
5883
5884 assert_dynamic_payload(&session, 1, 30);
5885 assert_eq!(
5886 JOURNALED_INDEX_STORE.with(|store| store.borrow().len()),
5887 1,
5888 "canonical derived state must contain only the newest membership",
5889 );
5890 assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
5891 }
5892
5893 #[test]
5894 fn ready_cardinality_combines_durable_base_with_exact_live_delta_and_fold_maintenance() {
5895 let session = initialize_journaled();
5896 session
5897 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5898 entity: ENTITY_NAME.to_string(),
5899 patch: dynamic_payload_patch(10),
5900 })
5901 .expect("initial cardinality row should commit");
5902 assert!(
5903 session
5904 .db
5905 .drive_startup_recovery_page()
5906 .expect("initial cardinality row should fold"),
5907 );
5908 drive_journaled_cardinality_to_ready(&session);
5909 let handle = session
5910 .db
5911 .store_handle(JOURNALED_STORE_PATH)
5912 .expect("journaled cardinality store should resolve");
5913 let (index_id, prefix_components) = journaled_user_index_prefix();
5914 reset_journaled_cardinality_projections();
5915 assert_eq!(
5916 JOURNALED_DATA_STORE.with(|store| store.borrow().exact_entity_count(ENTITY_TAG)),
5917 None,
5918 "the reopened-style volatile full count must remain unavailable",
5919 );
5920 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 1);
5921
5922 session
5923 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5924 entity: ENTITY_NAME.to_string(),
5925 patch: dynamic_payload_patch(10),
5926 })
5927 .expect("post-Ready row should commit into the live overlay");
5928 for payload in [20, 10] {
5929 session
5930 .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
5931 entity: ENTITY_NAME.to_string(),
5932 key: InputValue::Nat64(2),
5933 patch: dynamic_payload_patch(payload),
5934 })
5935 .expect("same-key post-Ready overlay should commit");
5936 }
5937 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 2);
5938 for folded in 1..=3 {
5939 let complete = session
5940 .db
5941 .drive_startup_recovery_page()
5942 .expect("post-Ready row should fold with exact maintenance");
5943 assert_eq!(complete, folded == 3);
5944 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 2);
5945 }
5946 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 2);
5947 session
5948 .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
5949 entity: ENTITY_NAME.to_string(),
5950 key: InputValue::Nat64(2),
5951 })
5952 .expect("post-Ready delete should commit into the live overlay");
5953 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 1);
5954 assert!(
5955 session
5956 .db
5957 .drive_startup_recovery_page()
5958 .expect("post-Ready delete should fold with exact maintenance"),
5959 );
5960 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 1);
5961 mark_journaled_cardinality_building();
5962 assert_eq!(
5963 handle.exact_entity_count(ENTITY_TAG),
5964 None,
5965 "non-Ready evidence must select the conservative path",
5966 );
5967 }
5968
5969 #[test]
5970 fn journaled_cardinality_rejects_volatile_counts_and_unfolded_accepted_root_drift() {
5971 let session = initialize_journaled();
5972 session
5973 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5974 entity: ENTITY_NAME.to_string(),
5975 patch: dynamic_payload_patch(10),
5976 })
5977 .expect("cardinality fixture row should commit");
5978 assert!(
5979 session
5980 .db
5981 .drive_startup_recovery_page()
5982 .expect("cardinality fixture row should fold"),
5983 );
5984 drive_journaled_cardinality_to_ready(&session);
5985 let handle = session
5986 .db
5987 .store_handle(JOURNALED_STORE_PATH)
5988 .expect("journaled cardinality store should resolve");
5989 let (index_id, prefix_components) = journaled_user_index_prefix();
5990 let data_generation = JOURNALED_DATA_STORE.with(|store| store.borrow().generation());
5991
5992 assert_eq!(
5993 JOURNALED_DATA_STORE.with(|store| store.borrow().exact_entity_count(ENTITY_TAG)),
5994 Some(1),
5995 "the live full-count cache should be populated before accepted-root drift",
5996 );
5997 assert_eq!(
5998 JOURNALED_INDEX_STORE.with(|store| {
5999 store.borrow().exact_prefix_cardinality(
6000 data_generation,
6001 IndexKeyKind::User,
6002 index_id,
6003 prefix_components.as_slice(),
6004 )
6005 }),
6006 Some(1),
6007 "the live prefix-count cache should be populated before accepted-root drift",
6008 );
6009 assert_eq!(
6010 JOURNALED_INDEX_STORE.with(|store| {
6011 store.borrow().exact_child_prefixes_for_parent_set(
6012 data_generation,
6013 IndexKeyKind::User,
6014 index_id,
6015 [prefix_components.as_slice()],
6016 8,
6017 )
6018 }),
6019 Some(Vec::new()),
6020 "the volatile child-prefix cache should demonstrate the bypass fixture",
6021 );
6022 assert_eq!(
6023 handle.exact_user_index_child_prefixes_for_parent_set(
6024 data_generation,
6025 index_id,
6026 [prefix_components.as_slice()],
6027 8,
6028 ),
6029 None,
6030 "journaled child enumeration must use its conservative route instead of volatile authority",
6031 );
6032 assert_journaled_cardinality(handle, index_id, prefix_components.as_slice(), 1);
6033
6034 let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
6035 JOURNALED_STORE_PATH,
6036 AcceptedSchemaRevision::new(2),
6037 BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH, false))]),
6038 BTreeMap::from([
6039 ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
6040 ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
6041 ]),
6042 );
6043 crate::db::commit::publish_accepted_schema_candidate(
6044 JOURNALED_STORE_PATH,
6045 handle,
6046 AcceptedSchemaRevision::INITIAL,
6047 &candidate,
6048 )
6049 .expect("a successor accepted root should publish into the live overlay");
6050
6051 assert_eq!(
6052 handle.exact_entity_count(ENTITY_TAG),
6053 None,
6054 "an unfolded accepted root must invalidate durable evidence immediately",
6055 );
6056 assert_eq!(
6057 handle.exact_user_index_prefix_count(
6058 data_generation,
6059 IndexKeyKind::User,
6060 index_id,
6061 prefix_components.as_slice(),
6062 ),
6063 None,
6064 "journaled consumers must not fall back to a populated volatile prefix cache",
6065 );
6066 }
6067
6068 #[test]
6069 fn journaled_convergence_uses_final_batch_rows_for_unique_release() {
6070 let session = initialize_journaled_with_unique_payload();
6071 let inserted = session
6072 .execute_trusted_dynamic_insert_batch(
6073 ENTITY_NAME,
6074 vec![dynamic_payload_patch(10), dynamic_payload_patch(20)],
6075 )
6076 .expect("the unique journal fixture should commit");
6077 assert_eq!(
6078 inserted.rows,
6079 vec![expected_dynamic_row(1, 10), expected_dynamic_row(2, 20)],
6080 );
6081 assert!(
6082 session
6083 .db
6084 .drive_startup_recovery_page()
6085 .expect("the unique fixture should become canonical"),
6086 );
6087
6088 let swapped = session
6089 .execute_trusted_dynamic_mutation_batch(vec![
6090 DynamicMutation::Update {
6091 entity: ENTITY_NAME.to_string(),
6092 key: InputValue::Nat64(1),
6093 patch: dynamic_payload_patch(20),
6094 },
6095 DynamicMutation::Update {
6096 entity: ENTITY_NAME.to_string(),
6097 key: InputValue::Nat64(2),
6098 patch: dynamic_payload_patch(10),
6099 },
6100 ])
6101 .expect("one journal batch should admit a final-row unique swap");
6102 assert_eq!(
6103 swapped.rows,
6104 vec![expected_dynamic_row(1, 20), expected_dynamic_row(2, 10)],
6105 );
6106 assert!(
6107 session
6108 .db
6109 .drive_startup_recovery_page()
6110 .expect("the unique swap should converge in one complete batch"),
6111 );
6112
6113 let released = session
6114 .execute_trusted_dynamic_mutation_batch(vec![
6115 DynamicMutation::Delete {
6116 entity: ENTITY_NAME.to_string(),
6117 key: InputValue::Nat64(1),
6118 },
6119 DynamicMutation::Insert {
6120 entity: ENTITY_NAME.to_string(),
6121 patch: dynamic_payload_patch(20),
6122 },
6123 ])
6124 .expect("a journaled delete should release its unique value to the final insert");
6125 assert_eq!(
6126 released.rows,
6127 vec![expected_dynamic_row(1, 20), expected_dynamic_row(3, 20)],
6128 );
6129 assert!(
6130 session
6131 .db
6132 .drive_startup_recovery_page()
6133 .expect("the delete and unique reuse should converge together"),
6134 );
6135
6136 assert_dynamic_payload(&session, 2, 10);
6137 assert_dynamic_payload(&session, 3, 20);
6138 assert_eq!(JOURNALED_INDEX_STORE.with(|store| store.borrow().len()), 2);
6139 assert!(
6140 session
6141 .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(20)],)
6142 .is_err(),
6143 "the converged unique index must remain authoritative",
6144 );
6145 }
6146
6147 #[test]
6148 fn journaled_startup_recovery_completes_one_large_batch_atomically() {
6149 let session = initialize_journaled();
6150 let catalog = session
6151 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
6152 .expect("journaled identity catalog should resolve");
6153 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
6154 .expect("journaled identity row layout should build");
6155 let payloads = (0_u64..129).collect::<Vec<_>>();
6156 session
6157 .execute_accepted_structural_save_batch(
6158 &catalog,
6159 &descriptor,
6160 batch(&payloads),
6161 Timestamp::from_millis(9),
6162 Ok,
6163 )
6164 .expect("one large journal batch should commit");
6165
6166 forget_recovered_domain_for_tests(&session.db)
6167 .expect("upgrade should reset recovery ownership");
6168 assert!(
6169 session
6170 .db
6171 .drive_startup_recovery_page()
6172 .expect("the complete batch recovery page should commit"),
6173 );
6174
6175 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 129);
6176 JOURNALED_TAIL_STORE.with(|tail| {
6177 let tail = tail.borrow();
6178 assert!(!tail.has_stored_batch());
6179 });
6180 assert_dynamic_payload(&session, 1, 0);
6181 assert_dynamic_payload(&session, 129, 128);
6182 }
6183
6184 #[test]
6185 fn complete_batch_validation_rejects_a_late_record_before_canonical_writes() {
6186 let session = initialize_journaled();
6187 let catalog = session
6188 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
6189 .expect("journaled identity catalog should resolve");
6190 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
6191 .expect("journaled identity row layout should build");
6192 session
6193 .execute_accepted_structural_save_batch(
6194 &catalog,
6195 &descriptor,
6196 batch(&[7]),
6197 Timestamp::from_millis(9),
6198 Ok,
6199 )
6200 .expect("journal batch predecessor should commit");
6201
6202 JOURNALED_TAIL_STORE.with(|tail| {
6203 let mut tail = tail.borrow_mut();
6204 let original = tail
6205 .next_batch_after(JournalSequence::new(0))
6206 .expect("journal batch should decode")
6207 .expect("journal batch should exist");
6208 let mut records = original.records().to_vec();
6209 records.push(
6210 JournalRecord::schema_put(JOURNALED_STORE_PATH, vec![0xff; 8])
6211 .expect("bounded semantic corruption should build"),
6212 );
6213 let corrupted = JournalBatch::new_with_database_commit_sequence(
6214 original.batch_id(),
6215 original.commit_marker_id(),
6216 original.journal_sequence(),
6217 original.database_commit_sequence(),
6218 records,
6219 )
6220 .expect("current corrupt batch shape should build");
6221 let encoded = encode_journal_batch(&corrupted)
6222 .expect("current corrupt batch envelope should encode");
6223 tail.clear_batches_through(original.journal_sequence());
6224 tail.insert_raw_batch_for_tests(original.journal_sequence(), encoded)
6225 .expect("corrupt persisted batch should replace the predecessor");
6226 });
6227
6228 forget_recovered_domain_for_tests(&session.db)
6229 .expect("upgrade should reset recovery ownership");
6230 let error = session
6231 .db
6232 .drive_startup_recovery_page()
6233 .expect_err("late semantic corruption must fail before fold apply");
6234 assert_eq!(error.class(), ErrorClass::Corruption);
6235 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 0);
6236 JOURNALED_TAIL_STORE.with(|tail| {
6237 let tail = tail.borrow();
6238 assert_eq!(
6239 tail.fold_watermark()
6240 .expect("watermark should remain readable")
6241 .highest_folded_journal_sequence(),
6242 JournalSequence::new(0),
6243 );
6244 assert!(tail.has_stored_batch());
6245 });
6246 }
6247
6248 #[test]
6249 fn prepared_batch_row_evidence_rejects_a_late_malformed_row_before_canonical_writes() {
6250 let session = initialize_journaled();
6251 let catalog = session
6252 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
6253 .expect("journaled identity catalog should resolve");
6254 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
6255 .expect("journaled identity row layout should build");
6256 session
6257 .execute_accepted_structural_save_batch(
6258 &catalog,
6259 &descriptor,
6260 batch(&[7, 8]),
6261 Timestamp::from_millis(9),
6262 Ok,
6263 )
6264 .expect("two-row journal batch should commit");
6265
6266 JOURNALED_TAIL_STORE.with(|tail| {
6267 let mut tail = tail.borrow_mut();
6268 let original = tail
6269 .next_batch_after(JournalSequence::new(0))
6270 .expect("journal batch should decode")
6271 .expect("journal batch should exist");
6272 let mut records = original.records().to_vec();
6273 let mut row_ordinal = 0_u8;
6274 for record in &mut records {
6275 if let JournalRecord::RowPut { row_bytes, .. } = record {
6276 row_ordinal = row_ordinal.saturating_add(1);
6277 if row_ordinal == 2 {
6278 *row_bytes = vec![0xff; 8];
6279 break;
6280 }
6281 }
6282 }
6283 assert_eq!(row_ordinal, 2, "the late row record should be present");
6284 let corrupted = JournalBatch::new_with_database_commit_sequence(
6285 original.batch_id(),
6286 original.commit_marker_id(),
6287 original.journal_sequence(),
6288 original.database_commit_sequence(),
6289 records,
6290 )
6291 .expect("current corrupt batch shape should build");
6292 let encoded = encode_journal_batch(&corrupted)
6293 .expect("current corrupt batch envelope should encode");
6294 tail.clear_batches_through(original.journal_sequence());
6295 tail.insert_raw_batch_for_tests(original.journal_sequence(), encoded)
6296 .expect("corrupt persisted batch should replace the predecessor");
6297 });
6298
6299 forget_recovered_domain_for_tests(&session.db)
6300 .expect("upgrade should reset recovery ownership");
6301 let error = session
6302 .db
6303 .drive_startup_recovery_page()
6304 .expect_err("late malformed row must fail during complete batch preparation");
6305 assert_eq!(error.class(), ErrorClass::Corruption);
6306 assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 0);
6307 JOURNALED_TAIL_STORE.with(|tail| {
6308 let tail = tail.borrow();
6309 assert_eq!(
6310 tail.fold_watermark()
6311 .expect("watermark should remain readable")
6312 .highest_folded_journal_sequence(),
6313 JournalSequence::new(0),
6314 );
6315 assert!(tail.has_stored_batch());
6316 });
6317 }
6318
6319 #[test]
6320 #[ignore = "release-closeout native timing probe for one marker-authorized driver recovery"]
6321 fn identity_recovery_closeout_reports_driver_time() {
6322 let session = initialize_journaled();
6323 let catalog = session
6324 .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
6325 .expect("journaled identity catalog should resolve");
6326 let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
6327 .expect("journaled identity row layout should build");
6328
6329 interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
6330 let interrupted = session.execute_accepted_structural_save_batch(
6331 &catalog,
6332 &descriptor,
6333 batch(&[1]),
6334 Timestamp::from_millis(10),
6335 Ok,
6336 );
6337 assert!(
6338 interrupted.is_err(),
6339 "the selected publication boundary should interrupt",
6340 );
6341
6342 let start = Instant::now();
6343 assert!(
6344 session
6345 .db
6346 .drive_startup_recovery_page()
6347 .expect("dedicated driver should recover before allocation"),
6348 );
6349 let committed = session
6350 .execute_accepted_structural_save_batch(
6351 &catalog,
6352 &descriptor,
6353 batch(&[2]),
6354 Timestamp::from_millis(11),
6355 Ok,
6356 )
6357 .expect("post-recovery allocation should commit");
6358 let elapsed = start.elapsed();
6359 assert_eq!(
6360 committed
6361 .into_iter()
6362 .map(|row| row.values)
6363 .collect::<Vec<_>>(),
6364 vec![vec![Value::Nat64(2), Value::Nat64(2)]],
6365 );
6366
6367 println!(
6368 "identity recovery closeout: driver_nanos={}",
6369 elapsed.as_nanos(),
6370 );
6371 }
6372}
6373
6374#[cfg(test)]
6375mod targeted_rule_mutation_tests {
6376 use super::{
6377 DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
6378 DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
6379 };
6380 use crate::{
6381 db::{
6382 data::{DataStore, encode_input_value_for_candidate_field_contract},
6383 index::IndexStore,
6384 registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
6385 schema::{
6386 AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
6387 AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
6388 AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
6389 ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
6390 PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
6391 ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
6392 SchemaRowLayout, SchemaStore, SchemaVersion,
6393 accepted_schema_candidate_with_catalogs_for_tests,
6394 build_record_newtype_composite_catalog_for_tests,
6395 empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
6396 },
6397 },
6398 error::InternalError,
6399 traits::{CanisterKind, Path},
6400 types::EntityTag,
6401 value::InputValue,
6402 };
6403 use icydb_schema::{
6404 ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
6405 };
6406 use std::{cell::RefCell, collections::BTreeMap};
6407
6408 const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
6409 const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
6410 const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
6411 const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
6412 const UPDATED_AT_SOURCE: &str =
6413 "session::write::targeted_rule_mutation_tests::Entity::updated_at";
6414 const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
6415 const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
6416 const DEGREE_MEMBER_SOURCE: &str =
6417 "session::write::targeted_rule_mutation_tests::Profile::degree";
6418 const DEGREE_RULE_SOURCE: &str =
6419 "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
6420
6421 struct TestCanister;
6422
6423 impl Path for TestCanister {
6424 const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
6425 }
6426
6427 impl CanisterKind for TestCanister {
6428 const COMMIT_MEMORY_ID: u8 = 43;
6429 const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
6430 const STARTUP_MEMORY_ID: u8 = 49;
6431 const STARTUP_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.startup.control.v1";
6432 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
6433 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
6434 "icydb.targeted_mutation_tests.integrity.progress.v1";
6435 }
6436
6437 thread_local! {
6438 static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
6439 static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
6440 static SCHEMA_STORE: RefCell<SchemaStore> =
6441 const { RefCell::new(SchemaStore::init_heap()) };
6442 static STORE_REGISTRY: StoreRegistry = {
6443 let mut registry = StoreRegistry::new();
6444 registry.register_store(
6445 STORE_PATH,
6446 &DATA_STORE,
6447 &INDEX_STORE,
6448 &SCHEMA_STORE,
6449 StoreAllocationIdentities::absent(),
6450 StoreRuntimeStorageCapabilities::heap(),
6451 ).expect("targeted mutation test store should register");
6452 registry
6453 };
6454 }
6455
6456 fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
6457 parse(raw.to_string()).expect("test source identity should admit")
6458 }
6459
6460 fn profile_input(degree: u64) -> InputValue {
6461 InputValue::Map(vec![(
6462 InputValue::Text("degree".to_string()),
6463 InputValue::Nat64(degree),
6464 )])
6465 }
6466
6467 fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
6468 DynamicStructuralPatch::new(vec![
6469 (
6470 "id".to_string(),
6471 DynamicWriteCell::Value(InputValue::Nat64(id)),
6472 ),
6473 (
6474 "profile".to_string(),
6475 DynamicWriteCell::Value(profile_input(degree)),
6476 ),
6477 ])
6478 }
6479
6480 fn encoded_value(
6481 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
6482 composite_catalog: &AcceptedCompositeCatalog,
6483 name: &str,
6484 kind: &AcceptedFieldKind,
6485 storage_decode: FieldStorageDecode,
6486 leaf_codec: LeafCodec,
6487 value: InputValue,
6488 ) -> Vec<u8> {
6489 let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
6490 encode_input_value_for_candidate_field_contract(
6491 enum_catalog,
6492 composite_catalog,
6493 field,
6494 value,
6495 &mut ValueAdmissionBudget::standard(),
6496 )
6497 .expect("test accepted value should encode")
6498 }
6499
6500 fn nat64_literal(
6501 enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
6502 composite_catalog: &AcceptedCompositeCatalog,
6503 value: u64,
6504 ) -> AcceptedCheckLiteralV1 {
6505 let kind = AcceptedFieldKind::Nat64;
6506 AcceptedCheckLiteralV1::from_accepted_parts(
6507 kind.clone(),
6508 FieldStorageDecode::ByKind,
6509 LeafCodec::Scalar(ScalarCodec::Nat64),
6510 encoded_value(
6511 enum_catalog,
6512 composite_catalog,
6513 "degree_bound",
6514 &kind,
6515 FieldStorageDecode::ByKind,
6516 LeafCodec::Scalar(ScalarCodec::Nat64),
6517 InputValue::Nat64(value),
6518 ),
6519 )
6520 }
6521
6522 fn targeted_constraint_id(error: &InternalError) -> u32 {
6523 let facts = error.diagnostic_facts();
6524 assert!(facts.contains(&(
6525 icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
6526 icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
6527 )));
6528 assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
6529 assert!(facts.contains(&(
6530 icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
6531 icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
6532 )));
6533 assert_eq!(
6534 facts
6535 .iter()
6536 .filter(|(tag, _)| matches!(
6537 tag,
6538 icydb_diagnostic_code::DiagnosticFactTag::RootField
6539 | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
6540 ))
6541 .copied()
6542 .collect::<Vec<_>>(),
6543 vec![
6544 (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
6545 (
6546 icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
6547 icydb_diagnostic_code::pack_u32_pair(1, 1),
6548 ),
6549 ]
6550 );
6551 let value = facts
6552 .iter()
6553 .find_map(|(tag, value)| {
6554 (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
6555 })
6556 .expect("targeted mutation should retain its accepted constraint ID");
6557 u32::try_from(value).expect("accepted constraint ID fits u32")
6558 }
6559
6560 #[expect(
6561 clippy::too_many_lines,
6562 reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
6563 )]
6564 #[test]
6565 fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
6566 DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
6567 INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
6568 SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
6569
6570 let entity_tag = EntityTag::new(93);
6571 let enum_catalog = empty_accepted_enum_catalog_for_tests();
6572 let (composite_catalog, profile_type, degree_type, degree_member) =
6573 build_record_newtype_composite_catalog_for_tests(
6574 "tests::TargetedProfile".to_string(),
6575 "degree".to_string(),
6576 "tests::TargetedDegree".to_string(),
6577 AcceptedFieldKind::Nat64,
6578 &enum_catalog,
6579 )
6580 .expect("targeted mutation composites should close");
6581 let profile_kind = AcceptedFieldKind::Composite {
6582 type_id: profile_type,
6583 };
6584 let profile_default = encoded_value(
6585 &enum_catalog,
6586 &composite_catalog,
6587 "profile",
6588 &profile_kind,
6589 FieldStorageDecode::CatalogValue,
6590 LeafCodec::Structural,
6591 profile_input(12),
6592 );
6593 let fields = vec![
6594 PersistedFieldSnapshot::new_initial(
6595 FieldId::new(1),
6596 "id".to_string(),
6597 SchemaFieldSlot::new(0),
6598 AcceptedFieldKind::Nat64,
6599 Vec::new(),
6600 false,
6601 SchemaInsertDefault::None,
6602 FieldStorageDecode::ByKind,
6603 LeafCodec::Scalar(ScalarCodec::Nat64),
6604 ),
6605 PersistedFieldSnapshot::new_initial(
6606 FieldId::new(2),
6607 "profile".to_string(),
6608 SchemaFieldSlot::new(1),
6609 profile_kind,
6610 vec![PersistedNestedLeafSnapshot::new(
6611 vec!["degree".to_string()],
6612 AcceptedFieldKind::Composite {
6613 type_id: degree_type,
6614 },
6615 false,
6616 )],
6617 false,
6618 SchemaInsertDefault::SlotPayload(profile_default),
6619 FieldStorageDecode::CatalogValue,
6620 LeafCodec::Structural,
6621 ),
6622 PersistedFieldSnapshot::new_initial_with_write_policy(
6623 FieldId::new(3),
6624 "updated_at".to_string(),
6625 SchemaFieldSlot::new(2),
6626 AcceptedFieldKind::Timestamp,
6627 Vec::new(),
6628 false,
6629 SchemaInsertDefault::None,
6630 SchemaFieldWritePolicy::from_model_policies(
6631 None,
6632 Some(FieldWriteManagement::UpdatedAt),
6633 ),
6634 FieldStorageDecode::ByKind,
6635 LeafCodec::Scalar(ScalarCodec::Timestamp),
6636 ),
6637 ];
6638 let mut snapshot = PersistedSchemaSnapshot::new(
6639 SchemaVersion::initial(),
6640 ENTITY_SOURCE.to_string(),
6641 "TargetedMutation".to_string(),
6642 FieldId::new(1),
6643 SchemaRowLayout::initial(
6644 fields
6645 .iter()
6646 .map(|field| (field.id(), field.slot()))
6647 .collect(),
6648 ),
6649 fields,
6650 );
6651 let constraint_catalog = snapshot
6652 .constraint_catalog()
6653 .clone()
6654 .with_added_targeted_rule(
6655 "profile_degree_multiple".to_string(),
6656 ConstraintOrigin::Generated,
6657 AcceptedRuleTarget::new(
6658 FieldId::new(2),
6659 AcceptedNamedTypeIdentity::Composite(degree_type),
6660 ),
6661 AcceptedRuleOperation::MultipleOf {
6662 divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
6663 },
6664 )
6665 .expect("targeted mutation rule should allocate");
6666 let targeted_rule_id = constraint_catalog
6667 .constraints()
6668 .last()
6669 .expect("targeted mutation rule should persist")
6670 .id();
6671 snapshot = snapshot.with_constraint_catalog(constraint_catalog);
6672
6673 let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
6674 let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
6675 let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
6676 let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
6677 let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
6678 let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
6679 let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
6680 let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
6681 let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
6682 BTreeMap::from([(entity_source, entity_tag)]),
6683 BTreeMap::from([
6684 ((entity_tag, id_source), FieldId::new(1)),
6685 ((entity_tag, profile_source), FieldId::new(2)),
6686 ((entity_tag, updated_at_source), FieldId::new(3)),
6687 ]),
6688 BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
6689 BTreeMap::new(),
6690 BTreeMap::new(),
6691 )
6692 .with_initial_named_types_for_tests(
6693 BTreeMap::from([
6694 (
6695 profile_type_source,
6696 AcceptedNamedTypeIdentity::Composite(profile_type),
6697 ),
6698 (
6699 degree_type_source,
6700 AcceptedNamedTypeIdentity::Composite(degree_type),
6701 ),
6702 ]),
6703 BTreeMap::new(),
6704 BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
6705 );
6706 let candidate = accepted_schema_candidate_with_catalogs_for_tests(
6707 STORE_PATH,
6708 AcceptedSchemaRevision::INITIAL,
6709 enum_catalog,
6710 composite_catalog,
6711 source_bindings,
6712 BTreeMap::from([(entity_tag, snapshot)]),
6713 );
6714
6715 let session = DbSession::<TestCanister>::new(
6716 &STORE_REGISTRY,
6717 &crate::db::RequestExecutionRoot::__new_runtime_root(),
6718 );
6719 session
6720 .db
6721 .drive_startup_recovery_page()
6722 .expect("targeted mutation test database should initialize");
6723 let store = session
6724 .db
6725 .store_handle(STORE_PATH)
6726 .expect("targeted mutation test store should resolve");
6727 crate::db::commit::publish_accepted_schema_candidate(
6728 STORE_PATH,
6729 store,
6730 AcceptedSchemaRevision::NONE,
6731 &candidate,
6732 )
6733 .expect("targeted mutation candidate should publish");
6734
6735 let dynamic_error = session
6736 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
6737 entity: "TargetedMutation".to_string(),
6738 patch: structural_patch(1, 12),
6739 })
6740 .expect_err("dynamic write must enforce the targeted rule");
6741 assert_eq!(
6742 targeted_constraint_id(&dynamic_error),
6743 targeted_rule_id.get()
6744 );
6745
6746 let binding = session
6747 .issue_typed_entity_binding(
6748 ENTITY_SOURCE,
6749 &[
6750 DynamicTypedFieldBindingRequest::new(
6751 ID_SOURCE.to_string(),
6752 DynamicTypedFieldType::Scalar(ScalarType::Nat64),
6753 false,
6754 ),
6755 DynamicTypedFieldBindingRequest::new(
6756 PROFILE_SOURCE.to_string(),
6757 DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
6758 false,
6759 ),
6760 DynamicTypedFieldBindingRequest::new(
6761 UPDATED_AT_SOURCE.to_string(),
6762 DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
6763 false,
6764 ),
6765 ],
6766 )
6767 .expect("targeted typed binding should issue");
6768 let typed_patch = binding
6769 .bind_write_fields(vec![
6770 (
6771 ID_SOURCE.to_string(),
6772 DynamicWriteCell::Value(InputValue::Nat64(2)),
6773 ),
6774 (
6775 PROFILE_SOURCE.to_string(),
6776 DynamicWriteCell::Value(profile_input(12)),
6777 ),
6778 ])
6779 .expect("targeted typed patch should bind");
6780 let typed_error = session
6781 .execute_trusted_typed_mutation(
6782 &binding,
6783 &DynamicTypedMutation::Insert { patch: typed_patch },
6784 )
6785 .expect_err("typed write must enforce the targeted rule");
6786 assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
6787
6788 #[cfg(feature = "sql")]
6789 {
6790 let sql_error = session
6791 .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
6792 .expect_err("SQL default resolution must enforce the targeted rule");
6793 let crate::db::QueryError::Execute(execute) = sql_error else {
6794 panic!("targeted SQL write should fail at shared execution admission");
6795 };
6796 assert_eq!(
6797 targeted_constraint_id(execute.as_internal()),
6798 targeted_rule_id.get()
6799 );
6800 }
6801
6802 session
6803 .execute_trusted_dynamic_mutation_batch(vec![
6804 DynamicMutation::Insert {
6805 entity: "TargetedMutation".to_string(),
6806 patch: structural_patch(4, 5),
6807 },
6808 DynamicMutation::Insert {
6809 entity: "TargetedMutation".to_string(),
6810 patch: structural_patch(5, 12),
6811 },
6812 ])
6813 .expect_err("one invalid targeted value must reject the whole batch");
6814 assert_eq!(
6815 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
6816 Some(0),
6817 "no frontend or earlier valid batch row may escape targeted admission",
6818 );
6819
6820 let admitted = session
6821 .execute_trusted_dynamic_mutation_batch(vec![
6822 DynamicMutation::Insert {
6823 entity: "TargetedMutation".to_string(),
6824 patch: structural_patch(6, 5),
6825 },
6826 DynamicMutation::Insert {
6827 entity: "TargetedMutation".to_string(),
6828 patch: structural_patch(7, 10),
6829 },
6830 ])
6831 .expect("compliant targeted values should share one accepted batch");
6832 let [first, second] = admitted.rows.as_slice() else {
6833 panic!("the mixed targeted batch should return two rows");
6834 };
6835 let first_timestamp = first
6836 .get(2)
6837 .expect("the first mixed row should contain its managed timestamp");
6838 assert!(matches!(
6839 first_timestamp,
6840 crate::value::OutputValue::Timestamp(_)
6841 ));
6842 assert_eq!(
6843 second.get(2),
6844 Some(first_timestamp),
6845 "one accepted mixed batch must materialize one managed timestamp",
6846 );
6847 assert_eq!(
6848 DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
6849 Some(2),
6850 );
6851 }
6852}