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