Skip to main content

icydb_core/db/session/
write.rs

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