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    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2557    use crate::db::executor::budget::{HardExecutionBudget, HardExecutionFailureHeadroom};
2558    use crate::{
2559        db::{
2560            commit::{database_incarnation_id, forget_recovered_domain_for_tests},
2561            data::DataStore,
2562            executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
2563            index::IndexStore,
2564            integrity::{
2565                PhysicalUnitCheckpoint, QuickIntegrityStatus, RowInspectionLimits,
2566                execute_quick_integrity, execute_row_integrity_page,
2567            },
2568            journal::JournalTailStore,
2569            registry::{
2570                StoreAllocationIdentities, StoreAllocationIdentity, StoreRegistry,
2571                StoreRuntimeStorageCapabilities,
2572            },
2573            schema::{
2574                AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
2575                FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2576                PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2577                PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
2578                SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
2579                accepted_schema_candidate_with_field_bindings_for_tests,
2580            },
2581            write_context::MutationMode,
2582        },
2583        error::{ErrorClass, ErrorOrigin, InternalError},
2584        testing::test_memory,
2585        traits::{CanisterKind, Path},
2586        types::{EntityTag, Timestamp},
2587        value::{InputValue, OutputValue, Value},
2588    };
2589    use icydb_schema::{FieldSourceKey, ScalarType};
2590    use std::{cell::RefCell, collections::BTreeMap, time::Instant};
2591
2592    const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
2593    const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
2594    const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
2595    const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
2596    const ENTITY_NAME: &str = "IdentityRow";
2597    const ENTITY_TAG: EntityTag = EntityTag::new(93);
2598    const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
2599
2600    struct TestCanister;
2601
2602    impl Path for TestCanister {
2603        const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
2604    }
2605
2606    impl CanisterKind for TestCanister {
2607        const COMMIT_MEMORY_ID: u8 = 45;
2608        const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
2609        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
2610        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2611            "icydb.identity_pre_key_tests.integrity.progress.v1";
2612    }
2613
2614    thread_local! {
2615        static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2616        static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2617        static SCHEMA_STORE: RefCell<SchemaStore> =
2618            const { RefCell::new(SchemaStore::init_heap()) };
2619        static STORE_REGISTRY: StoreRegistry = {
2620            let mut registry = StoreRegistry::new();
2621            registry.register_store(
2622                STORE_PATH,
2623                &DATA_STORE,
2624                &INDEX_STORE,
2625                &SCHEMA_STORE,
2626                StoreAllocationIdentities::absent(),
2627                StoreRuntimeStorageCapabilities::heap(),
2628            ).expect("identity pre-key test store should register");
2629            registry
2630        };
2631        static JOURNALED_DATA_STORE: RefCell<DataStore> =
2632            RefCell::new(DataStore::init_journaled(test_memory(186)));
2633        static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
2634            RefCell::new(IndexStore::init_journaled(test_memory(187)));
2635        static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
2636            RefCell::new(SchemaStore::init_journaled(test_memory(188)));
2637        static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
2638            RefCell::new(JournalTailStore::init(test_memory(189)));
2639        static JOURNALED_STORE_REGISTRY: StoreRegistry = {
2640            let mut registry = StoreRegistry::new();
2641            registry.register_journaled_store(
2642                JOURNALED_STORE_PATH,
2643                &JOURNALED_DATA_STORE,
2644                &JOURNALED_INDEX_STORE,
2645                &JOURNALED_SCHEMA_STORE,
2646                &JOURNALED_TAIL_STORE,
2647                StoreAllocationIdentities::new_journaled(
2648                    StoreAllocationIdentity::new(186, "icydb.test.identity-range.data.v1"),
2649                    StoreAllocationIdentity::new(187, "icydb.test.identity-range.index.v1"),
2650                    StoreAllocationIdentity::new(188, "icydb.test.identity-range.schema.v1"),
2651                    StoreAllocationIdentity::new(189, "icydb.test.identity-range.journal.v1"),
2652                ),
2653                StoreRuntimeStorageCapabilities::journaled(),
2654            ).expect("identity range journaled store should register");
2655            registry
2656        };
2657    }
2658
2659    struct JournaledTestCanister;
2660
2661    impl Path for JournaledTestCanister {
2662        const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
2663    }
2664
2665    impl CanisterKind for JournaledTestCanister {
2666        const COMMIT_MEMORY_ID: u8 = 190;
2667        const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
2668        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
2669        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2670            "icydb.identity_range_tests.integrity.progress.v1";
2671    }
2672
2673    fn source_key(source: &str) -> FieldSourceKey {
2674        FieldSourceKey::try_new(source).expect("identity test field source should admit")
2675    }
2676
2677    fn identity_snapshot(store_path: &str) -> PersistedSchemaSnapshot {
2678        let fields = vec![
2679            PersistedFieldSnapshot::new_initial_with_write_policy(
2680                FieldId::new(1),
2681                "id".to_string(),
2682                SchemaFieldSlot::new(0),
2683                AcceptedFieldKind::Nat64,
2684                Vec::new(),
2685                false,
2686                SchemaInsertDefault::None,
2687                SchemaFieldWritePolicy::from_model_policies(
2688                    Some(FieldInsertGeneration::Identity),
2689                    None,
2690                ),
2691                FieldStorageDecode::ByKind,
2692                LeafCodec::Scalar(ScalarCodec::Nat64),
2693            ),
2694            PersistedFieldSnapshot::new_initial(
2695                FieldId::new(2),
2696                "payload".to_string(),
2697                SchemaFieldSlot::new(1),
2698                AcceptedFieldKind::Nat64,
2699                Vec::new(),
2700                false,
2701                SchemaInsertDefault::None,
2702                FieldStorageDecode::ByKind,
2703                LeafCodec::Scalar(ScalarCodec::Nat64),
2704            ),
2705        ];
2706        PersistedSchemaSnapshot::new_with_indexes(
2707            SchemaVersion::initial(),
2708            ENTITY_SOURCE.to_string(),
2709            ENTITY_NAME.to_string(),
2710            FieldId::new(1),
2711            SchemaRowLayout::initial(
2712                fields
2713                    .iter()
2714                    .map(|field| (field.id(), field.slot()))
2715                    .collect(),
2716            ),
2717            fields,
2718            vec![PersistedIndexSnapshot::new(
2719                SchemaIndexId::new(1).expect("identity test index ID should admit"),
2720                1,
2721                "by_payload".to_string(),
2722                store_path.to_string(),
2723                false,
2724                PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2725                    FieldId::new(2),
2726                    SchemaFieldSlot::new(1),
2727                    vec!["payload".to_string()],
2728                    AcceptedFieldKind::Nat64,
2729                    false,
2730                )]),
2731                None,
2732            )],
2733        )
2734    }
2735
2736    fn initialize() -> DbSession<TestCanister> {
2737        DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2738        INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2739        SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2740        let session = DbSession::<TestCanister>::new(&STORE_REGISTRY);
2741        session
2742            .db
2743            .ensure_recovered_state()
2744            .expect("identity pre-key test database should initialize");
2745        let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2746            STORE_PATH,
2747            AcceptedSchemaRevision::INITIAL,
2748            BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH))]),
2749            BTreeMap::from([
2750                ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2751                ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2752            ]),
2753        );
2754        let store = session
2755            .db
2756            .store_handle(STORE_PATH)
2757            .expect("identity pre-key test store should resolve");
2758        crate::db::commit::publish_accepted_schema_candidate(
2759            STORE_PATH,
2760            store,
2761            AcceptedSchemaRevision::NONE,
2762            &candidate,
2763        )
2764        .expect("identity candidate should publish with explicit zero state");
2765        session
2766    }
2767
2768    fn initialize_journaled() -> DbSession<JournaledTestCanister> {
2769        let session = DbSession::<JournaledTestCanister>::new(&JOURNALED_STORE_REGISTRY);
2770        session
2771            .db
2772            .ensure_recovered_state()
2773            .expect("journaled identity database should initialize");
2774        let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2775            JOURNALED_STORE_PATH,
2776            AcceptedSchemaRevision::INITIAL,
2777            BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH))]),
2778            BTreeMap::from([
2779                ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2780                ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2781            ]),
2782        );
2783        let store = session
2784            .db
2785            .store_handle(JOURNALED_STORE_PATH)
2786            .expect("journaled identity store should resolve");
2787        crate::db::commit::publish_accepted_schema_candidate(
2788            JOURNALED_STORE_PATH,
2789            store,
2790            AcceptedSchemaRevision::NONE,
2791            &candidate,
2792        )
2793        .expect("journaled identity candidate should publish");
2794        session
2795    }
2796
2797    fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
2798        AcceptedMutationIntentPatch::new()
2799            .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
2800    }
2801
2802    fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
2803        DynamicStructuralPatch::new(vec![(
2804            "payload".to_string(),
2805            DynamicWriteCell::Value(InputValue::Nat64(value)),
2806        )])
2807    }
2808
2809    fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
2810        vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
2811    }
2812
2813    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2814    fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
2815        session
2816            .issue_typed_entity_binding(
2817                ENTITY_SOURCE,
2818                &[
2819                    DynamicTypedFieldBindingRequest::new(
2820                        ID_SOURCE.to_string(),
2821                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2822                        false,
2823                    ),
2824                    DynamicTypedFieldBindingRequest::new(
2825                        PAYLOAD_SOURCE.to_string(),
2826                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2827                        false,
2828                    ),
2829                ],
2830            )
2831            .expect("exact-key test binding should issue")
2832    }
2833
2834    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2835    fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
2836        let output = session
2837            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
2838                entity: ENTITY_NAME.to_string(),
2839                patch: dynamic_payload_patch(payload),
2840            })
2841            .expect("exact-key fixture insert should commit");
2842        match output.rows.as_slice() {
2843            [row] => match row.as_slice() {
2844                [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
2845                    if *actual_payload == payload =>
2846                {
2847                    *id
2848                }
2849                _ => panic!("exact-key fixture should return its identity and payload"),
2850            },
2851            _ => panic!("exact-key fixture insert should return one row"),
2852        }
2853    }
2854
2855    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2856    fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
2857        let first = insert_exact_key_fixture(session, 41);
2858        let second = insert_exact_key_fixture(session, 42);
2859        let missing = u64::MAX;
2860        let binding = exact_key_binding(session);
2861        let gets_before = DataStore::current_get_call_count();
2862        let result = session
2863            .execute_public_exact_key_batch_for_typed_binding(
2864                &binding,
2865                &[second, missing, first, second],
2866            )
2867            .expect("exact-key batch should execute")
2868            .expect("exact-key binding should remain current");
2869
2870        assert_eq!(result.positions, vec![0, 1, 2, 0]);
2871        assert_eq!(
2872            result.distinct_rows,
2873            vec![
2874                Some(expected_dynamic_row(second, 42)),
2875                None,
2876                Some(expected_dynamic_row(first, 41)),
2877            ],
2878        );
2879        assert_eq!(
2880            DataStore::current_get_call_count().saturating_sub(gets_before),
2881            3,
2882            "four input positions with one duplicate must perform three physical reads",
2883        );
2884    }
2885
2886    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2887    #[test]
2888    fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
2889        assert_exact_key_batch(&initialize());
2890        assert_exact_key_batch(&initialize_journaled());
2891    }
2892
2893    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2894    #[test]
2895    fn exact_key_batch_uses_typed_hard_execution_budget() {
2896        let session = initialize();
2897        let binding = exact_key_binding(&session);
2898        let budget =
2899            HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
2900        let error = session
2901            .execute_exact_key_batch_with_hard_budget_for_tests(&binding, &[u64::MAX], &budget)
2902            .expect_err("zero query budget should reject the exact-key route");
2903
2904        assert!(matches!(
2905            error.diagnostic().detail(),
2906            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
2907                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
2908            })
2909        ));
2910        let facts = error.diagnostic_facts();
2911        assert_eq!(
2912            &facts[..5],
2913            &[
2914                (
2915                    icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
2916                    icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions.raw(),
2917                ),
2918                (icydb_diagnostic_code::DiagnosticFactTag::Limit, 0),
2919                (icydb_diagnostic_code::DiagnosticFactTag::Actual, 1),
2920                (
2921                    icydb_diagnostic_code::DiagnosticFactTag::ExecutionBudgetScope,
2922                    icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution.raw(),
2923                ),
2924                (
2925                    icydb_diagnostic_code::DiagnosticFactTag::ExecutionLane,
2926                    icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead.raw(),
2927                ),
2928            ],
2929        );
2930        assert_eq!(
2931            facts[5].0,
2932            icydb_diagnostic_code::DiagnosticFactTag::QueryShapeFingerprintPrefix,
2933        );
2934        assert_ne!(facts[5].1, 0);
2935    }
2936
2937    fn assert_dynamic_payload(session: &DbSession<TestCanister>, key: u64, expected_payload: u64) {
2938        let unchanged = session
2939            .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
2940                entity: ENTITY_NAME.to_string(),
2941                key: InputValue::Nat64(key),
2942                patch: dynamic_payload_patch(expected_payload),
2943            })
2944            .expect("the expected row should remain readable through a no-op update");
2945        assert_eq!(unchanged.affected_rows, 0);
2946        assert_eq!(
2947            unchanged.rows,
2948            vec![expected_dynamic_row(key, expected_payload)],
2949        );
2950    }
2951
2952    fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
2953        values
2954            .iter()
2955            .map(|value| {
2956                AcceptedStructuralMutation::save(
2957                    MutationMode::Insert,
2958                    AcceptedStructuralMutationTarget::ResolveFromAfterImage,
2959                    payload_patch(*value),
2960                )
2961            })
2962            .collect()
2963    }
2964
2965    fn assert_identity_boundary(error: &InternalError) {
2966        assert_eq!(error.class(), ErrorClass::Unsupported);
2967        assert_eq!(error.origin(), ErrorOrigin::Identity);
2968    }
2969
2970    #[test]
2971    fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
2972        let generated = insert_key_exists_after_generation(true);
2973        assert_eq!(generated.class(), ErrorClass::Corruption);
2974        assert_eq!(generated.origin(), ErrorOrigin::Identity);
2975
2976        let ordinary = insert_key_exists_after_generation(false);
2977        assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
2978    }
2979
2980    #[cfg(target_pointer_width = "64")]
2981    #[test]
2982    fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
2983        let error = checked_pre_key_candidate_count(
2984            usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
2985        )
2986        .expect_err("candidate counts beyond u32 must reject");
2987        assert_identity_boundary(&error);
2988    }
2989
2990    #[test]
2991    #[expect(
2992        clippy::too_many_lines,
2993        reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
2994    )]
2995    fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
2996        let session = initialize();
2997        let seeded = session
2998            .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
2999            .expect("seed rows should commit");
3000        assert_eq!(seeded.affected_rows, 1);
3001
3002        let split = session
3003            .execute_trusted_dynamic_mutation_batch(vec![
3004                DynamicMutation::Update {
3005                    entity: ENTITY_NAME.to_string(),
3006                    key: InputValue::Nat64(1),
3007                    patch: dynamic_payload_patch(60),
3008                },
3009                DynamicMutation::Insert {
3010                    entity: ENTITY_NAME.to_string(),
3011                    patch: dynamic_payload_patch(40),
3012                },
3013            ])
3014            .expect("one holding should split atomically");
3015        assert_eq!(split.affected_rows, 2);
3016        assert_eq!(
3017            split.rows,
3018            vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
3019            "split after-images must retain input order and exact quantity",
3020        );
3021
3022        let rejected_split = session
3023            .execute_trusted_dynamic_mutation_batch(vec![
3024                DynamicMutation::Update {
3025                    entity: ENTITY_NAME.to_string(),
3026                    key: InputValue::Nat64(1),
3027                    patch: dynamic_payload_patch(50),
3028                },
3029                DynamicMutation::Insert {
3030                    entity: ENTITY_NAME.to_string(),
3031                    patch: DynamicStructuralPatch::new(Vec::new()),
3032                },
3033            ])
3034            .expect_err("an invalid split output must reject the staged source update");
3035        assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
3036        assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
3037        assert_eq!(
3038            rejected_split.diagnostic_facts(),
3039            vec![
3040                (
3041                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3042                    ENTITY_TAG.value(),
3043                ),
3044                (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
3045                (
3046                    icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3047                    icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3048                ),
3049                (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3050            ],
3051        );
3052        assert_dynamic_payload(&session, 1, 60);
3053        assert_dynamic_payload(&session, 2, 40);
3054
3055        let transfer = session
3056            .execute_trusted_dynamic_mutation_batch(vec![
3057                DynamicMutation::Update {
3058                    entity: ENTITY_NAME.to_string(),
3059                    key: InputValue::Nat64(1),
3060                    patch: dynamic_payload_patch(70),
3061                },
3062                DynamicMutation::Update {
3063                    entity: ENTITY_NAME.to_string(),
3064                    key: InputValue::Nat64(2),
3065                    patch: dynamic_payload_patch(30),
3066                },
3067            ])
3068            .expect("distinct transfer patches should share one atomic batch");
3069        assert_eq!(
3070            transfer.rows,
3071            vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
3072            "the transfer must preserve the exact total quantity",
3073        );
3074
3075        let merge = session
3076            .execute_trusted_dynamic_mutation_batch(vec![
3077                DynamicMutation::Delete {
3078                    entity: ENTITY_NAME.to_string(),
3079                    key: InputValue::Nat64(2),
3080                },
3081                DynamicMutation::Update {
3082                    entity: ENTITY_NAME.to_string(),
3083                    key: InputValue::Nat64(1),
3084                    patch: dynamic_payload_patch(100),
3085                },
3086            ])
3087            .expect("two holdings should merge atomically");
3088        assert_eq!(
3089            merge.rows,
3090            vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
3091            "delete before-images and update after-images must retain input order",
3092        );
3093
3094        let resplit = session
3095            .execute_trusted_dynamic_mutation_batch(vec![
3096                DynamicMutation::Update {
3097                    entity: ENTITY_NAME.to_string(),
3098                    key: InputValue::Nat64(1),
3099                    patch: dynamic_payload_patch(60),
3100                },
3101                DynamicMutation::Insert {
3102                    entity: ENTITY_NAME.to_string(),
3103                    patch: dynamic_payload_patch(40),
3104                },
3105            ])
3106            .expect("the merged holding should split again");
3107        assert_eq!(
3108            resplit.rows,
3109            vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
3110        );
3111
3112        let rejected_merge = session
3113            .execute_trusted_dynamic_mutation_batch(vec![
3114                DynamicMutation::Delete {
3115                    entity: ENTITY_NAME.to_string(),
3116                    key: InputValue::Nat64(3),
3117                },
3118                DynamicMutation::Update {
3119                    entity: ENTITY_NAME.to_string(),
3120                    key: InputValue::Nat64(99),
3121                    patch: dynamic_payload_patch(100),
3122                },
3123            ])
3124            .expect_err("a late missing merge target must preserve the earlier staged delete");
3125        assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
3126        assert_dynamic_payload(&session, 1, 60);
3127        assert_dynamic_payload(&session, 3, 40);
3128
3129        SCHEMA_STORE.with(|store| {
3130            let cursor = store
3131                .borrow()
3132                .identity_statement_cursor(
3133                    database_incarnation_id().expect("database incarnation should remain readable"),
3134                    ENTITY_TAG,
3135                    FieldId::new(1),
3136                    &AcceptedFieldKind::Nat64,
3137                )
3138                .expect("mixed Identity state should remain readable");
3139            assert_eq!(cursor.expected_high_water(), 3);
3140            assert!(!cursor.has_allocations());
3141        });
3142    }
3143
3144    #[test]
3145    fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
3146        let session = initialize();
3147        session
3148            .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3149            .expect("the holding fixture should initialize");
3150
3151        let duplicate = session
3152            .execute_trusted_dynamic_mutation_batch(vec![
3153                DynamicMutation::Update {
3154                    entity: ENTITY_NAME.to_string(),
3155                    key: InputValue::Nat64(1),
3156                    patch: dynamic_payload_patch(60),
3157                },
3158                DynamicMutation::Delete {
3159                    entity: ENTITY_NAME.to_string(),
3160                    key: InputValue::Nat64(1),
3161                },
3162            ])
3163            .expect_err("duplicate targets across operation kinds must reject");
3164        assert!(matches!(
3165            duplicate.diagnostic().detail(),
3166            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3167                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
3168            }),
3169        ));
3170        assert_eq!(
3171            duplicate.diagnostic_facts(),
3172            vec![
3173                (
3174                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3175                    ENTITY_TAG.value(),
3176                ),
3177                (
3178                    icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
3179                    0,
3180                ),
3181                (
3182                    icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
3183                    1,
3184                ),
3185            ],
3186        );
3187        assert_dynamic_payload(&session, 1, 100);
3188    }
3189
3190    #[test]
3191    fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
3192        let session = initialize();
3193        let empty = session
3194            .execute_trusted_dynamic_mutation_batch(Vec::new())
3195            .expect_err("an empty public batch must reject");
3196        assert!(matches!(
3197            empty.diagnostic().detail(),
3198            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3199                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
3200            }),
3201        ));
3202        assert_eq!(
3203            empty.diagnostic_facts(),
3204            vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
3205        );
3206
3207        let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
3208            .map(|_| DynamicMutation::Delete {
3209                entity: ENTITY_NAME.to_string(),
3210                key: InputValue::Nat64(1),
3211            })
3212            .collect();
3213        let over_bound = session
3214            .execute_trusted_dynamic_mutation_batch(requests)
3215            .expect_err("operation cap plus one must reject before row resolution");
3216        assert!(matches!(
3217            over_bound.diagnostic().detail(),
3218            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3219                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
3220            }),
3221        ));
3222        assert_eq!(
3223            over_bound.diagnostic_facts(),
3224            vec![
3225                (
3226                    icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
3227                    (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
3228                ),
3229                (
3230                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
3231                    MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
3232                ),
3233            ],
3234        );
3235    }
3236
3237    #[test]
3238    fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
3239        let mut exact = 0;
3240        add_structural_mutation_staged_bytes(
3241            &mut exact,
3242            [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
3243        )
3244        .expect("the exact staged-byte boundary should admit");
3245        assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
3246
3247        let error = add_structural_mutation_staged_bytes(&mut exact, [1])
3248            .expect_err("one byte above the staged-byte boundary must reject");
3249        assert!(matches!(
3250            error.diagnostic().detail(),
3251            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3252                boundary:
3253                    icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
3254            }),
3255        ));
3256        assert_eq!(
3257            error.diagnostic_facts(),
3258            vec![
3259                (
3260                    icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3261                    (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
3262                ),
3263                (
3264                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
3265                    MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
3266                ),
3267            ],
3268        );
3269
3270        validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
3271            .expect("the exact result-byte boundary should admit");
3272        let error = validate_structural_mutation_result_bytes(
3273            MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
3274        )
3275        .expect_err("one byte above the result-byte boundary must reject");
3276        assert!(matches!(
3277            error.diagnostic().detail(),
3278            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3279                boundary:
3280                    icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
3281            }),
3282        ));
3283        assert_eq!(
3284            error.diagnostic_facts(),
3285            vec![
3286                (
3287                    icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
3288                    (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
3289                ),
3290                (
3291                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
3292                    MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
3293                ),
3294            ],
3295        );
3296    }
3297
3298    #[expect(
3299        clippy::too_many_lines,
3300        reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
3301    )]
3302    #[test]
3303    fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
3304        let session = initialize();
3305        let catalog = session
3306            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3307            .expect("identity catalog should resolve");
3308        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3309            .expect("identity row layout should build");
3310        let initial_description = session
3311            .try_describe_entity_by_name(ENTITY_NAME)
3312            .expect("accepted Identity description should resolve");
3313        assert_eq!(
3314            initial_description.entity_tag(),
3315            catalog.identity().entity_tag().value()
3316        );
3317        assert_eq!(
3318            initial_description.accepted_schema_fingerprint_method(),
3319            catalog.fingerprint_method_version()
3320        );
3321        assert_eq!(
3322            initial_description.accepted_schema_fingerprint(),
3323            catalog.fingerprint()
3324        );
3325        let initial_identity = initial_description
3326            .identity()
3327            .expect("accepted Identity policy should be described");
3328        assert_eq!(initial_identity.field(), "id");
3329        assert_eq!(initial_identity.generator(), "Identity::next");
3330        assert_eq!(initial_identity.accepted_kind(), "nat64");
3331        assert_eq!(initial_identity.minimum(), 1);
3332        assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
3333        assert_eq!(initial_identity.high_water(), 0);
3334        assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
3335        assert!(!initial_identity.exhausted());
3336
3337        let rejected = session
3338            .execute_accepted_structural_save_batch(
3339                &catalog,
3340                &descriptor,
3341                batch(&[1_000, 2_000]),
3342                Timestamp::from_millis(6),
3343                |_| Err::<(), _>(InternalError::executor_unsupported()),
3344            )
3345            .expect_err("a rejected precommit result must not publish its tentative range");
3346        assert_eq!(rejected.class(), ErrorClass::Unsupported);
3347        assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
3348
3349        let rows = session
3350            .execute_accepted_structural_save_batch(
3351                &catalog,
3352                &descriptor,
3353                batch(&[10, 20, 30]),
3354                Timestamp::from_millis(7),
3355                Ok,
3356            )
3357            .expect("one accepted batch should commit rows and one identity range");
3358        assert_eq!(
3359            rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
3360            vec![
3361                vec![Value::Nat64(1), Value::Nat64(10)],
3362                vec![Value::Nat64(2), Value::Nat64(20)],
3363                vec![Value::Nat64(3), Value::Nat64(30)],
3364            ],
3365        );
3366
3367        let dynamic = session
3368            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3369                entity: ENTITY_NAME.to_string(),
3370                patch: DynamicStructuralPatch::new(vec![(
3371                    "payload".to_string(),
3372                    DynamicWriteCell::Value(InputValue::Nat64(40)),
3373                )]),
3374            })
3375            .expect("dynamic omission should commit through shared Identity generation");
3376        assert_eq!(dynamic.affected_rows, 1);
3377
3378        for (request, operation) in [
3379            (
3380                DynamicMutation::Insert {
3381                    entity: ENTITY_NAME.to_string(),
3382                    patch: DynamicStructuralPatch::new(vec![
3383                        (
3384                            "id".to_string(),
3385                            DynamicWriteCell::Value(InputValue::Nat64(41)),
3386                        ),
3387                        (
3388                            "payload".to_string(),
3389                            DynamicWriteCell::Value(InputValue::Nat64(42)),
3390                        ),
3391                    ]),
3392                },
3393                icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
3394            ),
3395            (
3396                DynamicMutation::Update {
3397                    entity: ENTITY_NAME.to_string(),
3398                    key: InputValue::Nat64(1),
3399                    patch: DynamicStructuralPatch::new(vec![(
3400                        "id".to_string(),
3401                        DynamicWriteCell::Default,
3402                    )]),
3403                },
3404                icydb_diagnostic_code::DiagnosticMutationOperation::Update,
3405            ),
3406        ] {
3407            let error = session
3408                .execute_trusted_dynamic_mutation(&request)
3409                .expect_err("structural Identity authorship and regeneration must reject");
3410            assert_eq!(error.class(), ErrorClass::Unsupported);
3411            assert_eq!(error.origin(), ErrorOrigin::Executor);
3412            assert_eq!(
3413                error.diagnostic_facts(),
3414                vec![
3415                    (
3416                        icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3417                        ENTITY_TAG.value(),
3418                    ),
3419                    (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3420                    (
3421                        icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3422                        operation.raw(),
3423                    ),
3424                    (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3425                ],
3426            );
3427        }
3428
3429        let binding = session
3430            .issue_typed_entity_binding(
3431                ENTITY_SOURCE,
3432                &[
3433                    DynamicTypedFieldBindingRequest::new(
3434                        ID_SOURCE.to_string(),
3435                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3436                        false,
3437                    ),
3438                    DynamicTypedFieldBindingRequest::new(
3439                        PAYLOAD_SOURCE.to_string(),
3440                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
3441                        false,
3442                    ),
3443                ],
3444            )
3445            .expect("typed output should bind the Identity field");
3446        let typed_patch = binding
3447            .bind_write_fields(vec![(
3448                PAYLOAD_SOURCE.to_string(),
3449                DynamicWriteCell::Value(InputValue::Nat64(50)),
3450            )])
3451            .expect("typed payload should lower");
3452        let typed = session
3453            .execute_trusted_typed_mutation(
3454                &binding,
3455                &DynamicTypedMutation::Insert { patch: typed_patch },
3456            )
3457            .expect("typed omission should commit through shared Identity generation");
3458        assert_eq!(
3459            typed
3460                .expect("typed insert should return one mutation result")
3461                .affected_rows,
3462            1,
3463        );
3464        let explicit_typed_patch = binding
3465            .bind_write_fields(vec![
3466                (
3467                    ID_SOURCE.to_string(),
3468                    DynamicWriteCell::Value(InputValue::Nat64(51)),
3469                ),
3470                (
3471                    PAYLOAD_SOURCE.to_string(),
3472                    DynamicWriteCell::Value(InputValue::Nat64(52)),
3473                ),
3474            ])
3475            .expect("the low-level binding should retain exact authored intent");
3476        let explicit_typed_error = session
3477            .execute_trusted_typed_mutation(
3478                &binding,
3479                &DynamicTypedMutation::Insert {
3480                    patch: explicit_typed_patch,
3481                },
3482            )
3483            .expect_err("typed Identity authorship must reject before allocation");
3484        assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
3485        assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
3486        assert_eq!(
3487            explicit_typed_error.diagnostic_facts(),
3488            vec![
3489                (
3490                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3491                    ENTITY_TAG.value(),
3492                ),
3493                (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
3494                (
3495                    icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3496                    icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3497                ),
3498                (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
3499            ],
3500        );
3501
3502        let replace_error = session
3503            .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
3504                entity: ENTITY_NAME.to_string(),
3505                key: InputValue::Nat64(99),
3506                patch: DynamicStructuralPatch::new(vec![(
3507                    "payload".to_string(),
3508                    DynamicWriteCell::Value(InputValue::Nat64(60)),
3509                )]),
3510            })
3511            .expect_err("save-as-insert with a chosen Identity must reject");
3512        assert_eq!(replace_error.class(), ErrorClass::Unsupported);
3513        assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
3514
3515        #[cfg(feature = "sql")]
3516        {
3517            for sql in [
3518                "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
3519                "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
3520            ] {
3521                let _result = session
3522                    .execute_trusted_sql_mutation(sql)
3523                    .expect("SQL omission and DEFAULT should commit Identity generation");
3524            }
3525
3526            let error = session
3527                .execute_trusted_sql_mutation(
3528                    "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
3529                )
3530                .expect_err("an explicit SQL Identity value must reject before allocation");
3531            let diagnostic = error.diagnostic();
3532            assert_eq!(
3533                diagnostic.code(),
3534                icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
3535            );
3536            assert!(matches!(
3537                diagnostic.detail(),
3538                Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
3539                    boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
3540                }),
3541            ));
3542        }
3543
3544        let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
3545        assert_eq!(
3546            DATA_STORE.with(|store| store.borrow().len()),
3547            expected_committed
3548        );
3549        SCHEMA_STORE.with(|store| {
3550            let cursor = store
3551                .borrow()
3552                .identity_statement_cursor(
3553                    database_incarnation_id().expect("database incarnation should remain readable"),
3554                    ENTITY_TAG,
3555                    FieldId::new(1),
3556                    &AcceptedFieldKind::Nat64,
3557                )
3558                .expect("committed writes must leave active state readable");
3559            assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
3560            assert!(!cursor.has_allocations());
3561        });
3562        let committed_description = session
3563            .try_describe_entity_by_name(ENTITY_NAME)
3564            .expect("committed Identity description should resolve");
3565        let committed_identity = committed_description
3566            .identity()
3567            .expect("accepted Identity policy should remain described");
3568        assert_eq!(
3569            committed_identity.high_water(),
3570            u128::from(expected_committed),
3571        );
3572        assert_eq!(
3573            committed_identity.remaining(),
3574            u128::from(u64::MAX - expected_committed),
3575        );
3576        assert!(!committed_identity.exhausted());
3577    }
3578
3579    #[test]
3580    #[expect(
3581        clippy::too_many_lines,
3582        reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
3583    )]
3584    fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
3585        let session = initialize_journaled();
3586        let catalog = session
3587            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3588            .expect("journaled identity catalog should resolve");
3589        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3590            .expect("journaled identity row layout should build");
3591
3592        for (ordinal, interruption) in [
3593            MutationCommitInterruption::MarkerPersisted,
3594            MutationCommitInterruption::JournalPublished,
3595            MutationCommitInterruption::RowsPublished,
3596            MutationCommitInterruption::StateMaterialized,
3597        ]
3598        .into_iter()
3599        .enumerate()
3600        {
3601            interrupt_next_mutation_commit_for_tests(interruption);
3602            let interrupted = session.execute_accepted_structural_save_batch(
3603                &catalog,
3604                &descriptor,
3605                batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
3606                Timestamp::from_millis(8),
3607                Ok,
3608            );
3609            assert!(
3610                interrupted.is_err(),
3611                "the selected durable boundary should interrupt",
3612            );
3613
3614            let committed = session
3615                .execute_accepted_structural_save_batch(
3616                    &catalog,
3617                    &descriptor,
3618                    batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
3619                    Timestamp::from_millis(9),
3620                    Ok,
3621                )
3622                .expect("the next mutation must recover before allocating");
3623            let expected_high_water =
3624                u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
3625            assert_eq!(
3626                committed
3627                    .into_iter()
3628                    .map(|row| row.values)
3629                    .collect::<Vec<_>>(),
3630                vec![vec![
3631                    Value::Nat64(expected_high_water),
3632                    Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
3633                ]],
3634            );
3635            assert_eq!(
3636                JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
3637                expected_high_water,
3638            );
3639            JOURNALED_SCHEMA_STORE.with(|store| {
3640                let cursor = store
3641                    .borrow()
3642                    .identity_statement_cursor(
3643                        database_incarnation_id()
3644                            .expect("database incarnation should remain readable"),
3645                        ENTITY_TAG,
3646                        FieldId::new(1),
3647                        &AcceptedFieldKind::Nat64,
3648                    )
3649                    .expect("guarded recovery must leave quiescent active state");
3650                assert_eq!(
3651                    cursor.expected_high_water(),
3652                    u128::from(expected_high_water),
3653                );
3654                assert!(!cursor.has_allocations());
3655            });
3656        }
3657
3658        for (ordinal, (interruption, deleted_key)) in [
3659            (MutationCommitInterruption::MarkerPersisted, 2),
3660            (MutationCommitInterruption::JournalPublished, 4),
3661            (MutationCommitInterruption::RowPrefixPublished, 6),
3662            (MutationCommitInterruption::RowsPublished, 8),
3663            (MutationCommitInterruption::StateMaterialized, 7),
3664        ]
3665        .into_iter()
3666        .enumerate()
3667        {
3668            let expected_payload =
3669                501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
3670            interrupt_next_mutation_commit_for_tests(interruption);
3671            let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
3672                DynamicMutation::Update {
3673                    entity: ENTITY_NAME.to_string(),
3674                    key: InputValue::Nat64(1),
3675                    patch: dynamic_payload_patch(expected_payload),
3676                },
3677                DynamicMutation::Delete {
3678                    entity: ENTITY_NAME.to_string(),
3679                    key: InputValue::Nat64(deleted_key),
3680                },
3681            ]);
3682            assert!(
3683                interrupted.is_err(),
3684                "the selected caller-key mixed publication boundary should interrupt",
3685            );
3686            let recovered_update = session
3687                .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3688                    entity: ENTITY_NAME.to_string(),
3689                    key: InputValue::Nat64(1),
3690                    patch: dynamic_payload_patch(expected_payload),
3691                })
3692                .expect("guarded reentry should complete the marker-authorized mixed batch");
3693            assert_eq!(
3694                recovered_update.affected_rows, 0,
3695                "the recovered update must already expose its admitted final image",
3696            );
3697            let recovered_delete = session
3698                .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
3699                    entity: ENTITY_NAME.to_string(),
3700                    key: InputValue::Nat64(deleted_key),
3701                })
3702                .expect_err("the recovered delete must already be materialized");
3703            assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
3704            JOURNALED_SCHEMA_STORE.with(|store| {
3705                let cursor = store
3706                    .borrow()
3707                    .identity_statement_cursor(
3708                        database_incarnation_id()
3709                            .expect("database incarnation should remain readable"),
3710                        ENTITY_TAG,
3711                        FieldId::new(1),
3712                        &AcceptedFieldKind::Nat64,
3713                    )
3714                    .expect("caller-key recovery must preserve active Identity state");
3715                assert_eq!(cursor.expected_high_water(), 8);
3716                assert!(!cursor.has_allocations());
3717            });
3718        }
3719
3720        forget_recovered_domain_for_tests(&session.db)
3721            .expect("the final journal tail should remain recoverable");
3722        session
3723            .db
3724            .ensure_recovered_state()
3725            .expect("derived rebuild must not allocate another identity");
3726
3727        let quick = execute_quick_integrity(&session.db, catalog.inspection_plan())
3728            .expect("quiescent Identity control inventory should be inspectable");
3729        assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
3730        let row_page = execute_row_integrity_page(
3731            &session.db,
3732            catalog.inspection_plan(),
3733            PhysicalUnitCheckpoint::BeforeFirst,
3734            RowInspectionLimits::standard(),
3735        )
3736        .expect("Identity rows should remain within committed high-water");
3737        assert!(row_page.exhausted());
3738        assert!(row_page.findings().is_empty());
3739
3740        assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
3741        assert!(
3742            JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
3743            "derived index rebuild should restore witnesses without allocating identities",
3744        );
3745        assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
3746        JOURNALED_SCHEMA_STORE.with(|store| {
3747            let cursor = store
3748                .borrow()
3749                .identity_statement_cursor(
3750                    database_incarnation_id().expect("database incarnation should remain readable"),
3751                    ENTITY_TAG,
3752                    FieldId::new(1),
3753                    &AcceptedFieldKind::Nat64,
3754                )
3755                .expect("folded identity state should reopen without allocating");
3756            assert_eq!(cursor.expected_high_water(), 8);
3757            assert!(!cursor.has_allocations());
3758        });
3759    }
3760
3761    #[test]
3762    #[ignore = "release-closeout native timing probe for one marker-authorized Identity recovery"]
3763    fn identity_recovery_closeout_reports_guarded_reentry_time() {
3764        let session = initialize_journaled();
3765        let catalog = session
3766            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
3767            .expect("journaled identity catalog should resolve");
3768        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
3769            .expect("journaled identity row layout should build");
3770
3771        interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
3772        let interrupted = session.execute_accepted_structural_save_batch(
3773            &catalog,
3774            &descriptor,
3775            batch(&[1]),
3776            Timestamp::from_millis(10),
3777            Ok,
3778        );
3779        assert!(
3780            interrupted.is_err(),
3781            "the selected publication boundary should interrupt",
3782        );
3783
3784        let start = Instant::now();
3785        let committed = session
3786            .execute_accepted_structural_save_batch(
3787                &catalog,
3788                &descriptor,
3789                batch(&[2]),
3790                Timestamp::from_millis(11),
3791                Ok,
3792            )
3793            .expect("guarded reentry should recover before allocation");
3794        let elapsed = start.elapsed();
3795        assert_eq!(
3796            committed
3797                .into_iter()
3798                .map(|row| row.values)
3799                .collect::<Vec<_>>(),
3800            vec![vec![Value::Nat64(2), Value::Nat64(2)]],
3801        );
3802
3803        println!(
3804            "identity recovery closeout: guarded_reentry_nanos={}",
3805            elapsed.as_nanos(),
3806        );
3807    }
3808}
3809
3810#[cfg(test)]
3811mod targeted_rule_mutation_tests {
3812    use super::{
3813        DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
3814        DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
3815    };
3816    use crate::{
3817        db::{
3818            data::{DataStore, encode_input_value_for_candidate_field_contract},
3819            index::IndexStore,
3820            registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
3821            schema::{
3822                AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
3823                AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
3824                AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
3825                ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
3826                PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
3827                ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
3828                SchemaRowLayout, SchemaStore, SchemaVersion,
3829                accepted_schema_candidate_with_catalogs_for_tests,
3830                build_record_newtype_composite_catalog_for_tests,
3831                empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
3832            },
3833        },
3834        error::InternalError,
3835        traits::{CanisterKind, Path},
3836        types::EntityTag,
3837        value::InputValue,
3838    };
3839    use icydb_schema::{
3840        ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
3841    };
3842    use std::{cell::RefCell, collections::BTreeMap};
3843
3844    const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
3845    const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
3846    const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
3847    const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
3848    const UPDATED_AT_SOURCE: &str =
3849        "session::write::targeted_rule_mutation_tests::Entity::updated_at";
3850    const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
3851    const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
3852    const DEGREE_MEMBER_SOURCE: &str =
3853        "session::write::targeted_rule_mutation_tests::Profile::degree";
3854    const DEGREE_RULE_SOURCE: &str =
3855        "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
3856
3857    struct TestCanister;
3858
3859    impl Path for TestCanister {
3860        const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
3861    }
3862
3863    impl CanisterKind for TestCanister {
3864        const COMMIT_MEMORY_ID: u8 = 43;
3865        const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
3866        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
3867        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3868            "icydb.targeted_mutation_tests.integrity.progress.v1";
3869    }
3870
3871    thread_local! {
3872        static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
3873        static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
3874        static SCHEMA_STORE: RefCell<SchemaStore> =
3875            const { RefCell::new(SchemaStore::init_heap()) };
3876        static STORE_REGISTRY: StoreRegistry = {
3877            let mut registry = StoreRegistry::new();
3878            registry.register_store(
3879                STORE_PATH,
3880                &DATA_STORE,
3881                &INDEX_STORE,
3882                &SCHEMA_STORE,
3883                StoreAllocationIdentities::absent(),
3884                StoreRuntimeStorageCapabilities::heap(),
3885            ).expect("targeted mutation test store should register");
3886            registry
3887        };
3888    }
3889
3890    fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
3891        parse(raw.to_string()).expect("test source identity should admit")
3892    }
3893
3894    fn profile_input(degree: u64) -> InputValue {
3895        InputValue::Map(vec![(
3896            InputValue::Text("degree".to_string()),
3897            InputValue::Nat64(degree),
3898        )])
3899    }
3900
3901    fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
3902        DynamicStructuralPatch::new(vec![
3903            (
3904                "id".to_string(),
3905                DynamicWriteCell::Value(InputValue::Nat64(id)),
3906            ),
3907            (
3908                "profile".to_string(),
3909                DynamicWriteCell::Value(profile_input(degree)),
3910            ),
3911        ])
3912    }
3913
3914    fn encoded_value(
3915        enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
3916        composite_catalog: &AcceptedCompositeCatalog,
3917        name: &str,
3918        kind: &AcceptedFieldKind,
3919        storage_decode: FieldStorageDecode,
3920        leaf_codec: LeafCodec,
3921        value: InputValue,
3922    ) -> Vec<u8> {
3923        let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
3924        encode_input_value_for_candidate_field_contract(
3925            enum_catalog,
3926            composite_catalog,
3927            field,
3928            value,
3929            &mut ValueAdmissionBudget::standard(),
3930        )
3931        .expect("test accepted value should encode")
3932    }
3933
3934    fn nat64_literal(
3935        enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
3936        composite_catalog: &AcceptedCompositeCatalog,
3937        value: u64,
3938    ) -> AcceptedCheckLiteralV1 {
3939        let kind = AcceptedFieldKind::Nat64;
3940        AcceptedCheckLiteralV1::from_accepted_parts(
3941            kind.clone(),
3942            FieldStorageDecode::ByKind,
3943            LeafCodec::Scalar(ScalarCodec::Nat64),
3944            encoded_value(
3945                enum_catalog,
3946                composite_catalog,
3947                "degree_bound",
3948                &kind,
3949                FieldStorageDecode::ByKind,
3950                LeafCodec::Scalar(ScalarCodec::Nat64),
3951                InputValue::Nat64(value),
3952            ),
3953        )
3954    }
3955
3956    fn targeted_constraint_id(error: &InternalError) -> u32 {
3957        let facts = error.diagnostic_facts();
3958        assert!(facts.contains(&(
3959            icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3960            icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3961        )));
3962        assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
3963        assert!(facts.contains(&(
3964            icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
3965            icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
3966        )));
3967        assert_eq!(
3968            facts
3969                .iter()
3970                .filter(|(tag, _)| matches!(
3971                    tag,
3972                    icydb_diagnostic_code::DiagnosticFactTag::RootField
3973                        | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
3974                ))
3975                .copied()
3976                .collect::<Vec<_>>(),
3977            vec![
3978                (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
3979                (
3980                    icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
3981                    icydb_diagnostic_code::pack_u32_pair(1, 1),
3982                ),
3983            ]
3984        );
3985        let value = facts
3986            .iter()
3987            .find_map(|(tag, value)| {
3988                (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
3989            })
3990            .expect("targeted mutation should retain its accepted constraint ID");
3991        u32::try_from(value).expect("accepted constraint ID fits u32")
3992    }
3993
3994    #[expect(
3995        clippy::too_many_lines,
3996        reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
3997    )]
3998    #[test]
3999    fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
4000        DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
4001        INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
4002        SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
4003
4004        let entity_tag = EntityTag::new(93);
4005        let enum_catalog = empty_accepted_enum_catalog_for_tests();
4006        let (composite_catalog, profile_type, degree_type, degree_member) =
4007            build_record_newtype_composite_catalog_for_tests(
4008                "tests::TargetedProfile".to_string(),
4009                "degree".to_string(),
4010                "tests::TargetedDegree".to_string(),
4011                AcceptedFieldKind::Nat64,
4012                &enum_catalog,
4013            )
4014            .expect("targeted mutation composites should close");
4015        let profile_kind = AcceptedFieldKind::Composite {
4016            type_id: profile_type,
4017        };
4018        let profile_default = encoded_value(
4019            &enum_catalog,
4020            &composite_catalog,
4021            "profile",
4022            &profile_kind,
4023            FieldStorageDecode::CatalogValue,
4024            LeafCodec::Structural,
4025            profile_input(12),
4026        );
4027        let fields = vec![
4028            PersistedFieldSnapshot::new_initial(
4029                FieldId::new(1),
4030                "id".to_string(),
4031                SchemaFieldSlot::new(0),
4032                AcceptedFieldKind::Nat64,
4033                Vec::new(),
4034                false,
4035                SchemaInsertDefault::None,
4036                FieldStorageDecode::ByKind,
4037                LeafCodec::Scalar(ScalarCodec::Nat64),
4038            ),
4039            PersistedFieldSnapshot::new_initial(
4040                FieldId::new(2),
4041                "profile".to_string(),
4042                SchemaFieldSlot::new(1),
4043                profile_kind,
4044                vec![PersistedNestedLeafSnapshot::new(
4045                    vec!["degree".to_string()],
4046                    AcceptedFieldKind::Composite {
4047                        type_id: degree_type,
4048                    },
4049                    false,
4050                )],
4051                false,
4052                SchemaInsertDefault::SlotPayload(profile_default),
4053                FieldStorageDecode::CatalogValue,
4054                LeafCodec::Structural,
4055            ),
4056            PersistedFieldSnapshot::new_initial_with_write_policy(
4057                FieldId::new(3),
4058                "updated_at".to_string(),
4059                SchemaFieldSlot::new(2),
4060                AcceptedFieldKind::Timestamp,
4061                Vec::new(),
4062                false,
4063                SchemaInsertDefault::None,
4064                SchemaFieldWritePolicy::from_model_policies(
4065                    None,
4066                    Some(FieldWriteManagement::UpdatedAt),
4067                ),
4068                FieldStorageDecode::ByKind,
4069                LeafCodec::Scalar(ScalarCodec::Timestamp),
4070            ),
4071        ];
4072        let mut snapshot = PersistedSchemaSnapshot::new(
4073            SchemaVersion::initial(),
4074            ENTITY_SOURCE.to_string(),
4075            "TargetedMutation".to_string(),
4076            FieldId::new(1),
4077            SchemaRowLayout::initial(
4078                fields
4079                    .iter()
4080                    .map(|field| (field.id(), field.slot()))
4081                    .collect(),
4082            ),
4083            fields,
4084        );
4085        let constraint_catalog = snapshot
4086            .constraint_catalog()
4087            .clone()
4088            .with_added_targeted_rule(
4089                "profile_degree_multiple".to_string(),
4090                ConstraintOrigin::Generated,
4091                AcceptedRuleTarget::new(
4092                    FieldId::new(2),
4093                    AcceptedNamedTypeIdentity::Composite(degree_type),
4094                ),
4095                AcceptedRuleOperation::MultipleOf {
4096                    divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
4097                },
4098            )
4099            .expect("targeted mutation rule should allocate");
4100        let targeted_rule_id = constraint_catalog
4101            .constraints()
4102            .last()
4103            .expect("targeted mutation rule should persist")
4104            .id();
4105        snapshot = snapshot.with_constraint_catalog(constraint_catalog);
4106
4107        let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
4108        let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
4109        let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
4110        let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
4111        let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
4112        let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
4113        let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
4114        let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
4115        let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
4116            BTreeMap::from([(entity_source, entity_tag)]),
4117            BTreeMap::from([
4118                ((entity_tag, id_source), FieldId::new(1)),
4119                ((entity_tag, profile_source), FieldId::new(2)),
4120                ((entity_tag, updated_at_source), FieldId::new(3)),
4121            ]),
4122            BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
4123            BTreeMap::new(),
4124            BTreeMap::new(),
4125        )
4126        .with_initial_named_types_for_tests(
4127            BTreeMap::from([
4128                (
4129                    profile_type_source,
4130                    AcceptedNamedTypeIdentity::Composite(profile_type),
4131                ),
4132                (
4133                    degree_type_source,
4134                    AcceptedNamedTypeIdentity::Composite(degree_type),
4135                ),
4136            ]),
4137            BTreeMap::new(),
4138            BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
4139        );
4140        let candidate = accepted_schema_candidate_with_catalogs_for_tests(
4141            STORE_PATH,
4142            AcceptedSchemaRevision::INITIAL,
4143            enum_catalog,
4144            composite_catalog,
4145            source_bindings,
4146            BTreeMap::from([(entity_tag, snapshot)]),
4147        );
4148
4149        let session = DbSession::<TestCanister>::new(&STORE_REGISTRY);
4150        session
4151            .db
4152            .ensure_recovered_state()
4153            .expect("targeted mutation test database should initialize");
4154        let store = session
4155            .db
4156            .store_handle(STORE_PATH)
4157            .expect("targeted mutation test store should resolve");
4158        crate::db::commit::publish_accepted_schema_candidate(
4159            STORE_PATH,
4160            store,
4161            AcceptedSchemaRevision::NONE,
4162            &candidate,
4163        )
4164        .expect("targeted mutation candidate should publish");
4165
4166        let dynamic_error = session
4167            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4168                entity: "TargetedMutation".to_string(),
4169                patch: structural_patch(1, 12),
4170            })
4171            .expect_err("dynamic write must enforce the targeted rule");
4172        assert_eq!(
4173            targeted_constraint_id(&dynamic_error),
4174            targeted_rule_id.get()
4175        );
4176
4177        let binding = session
4178            .issue_typed_entity_binding(
4179                ENTITY_SOURCE,
4180                &[
4181                    DynamicTypedFieldBindingRequest::new(
4182                        ID_SOURCE.to_string(),
4183                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4184                        false,
4185                    ),
4186                    DynamicTypedFieldBindingRequest::new(
4187                        PROFILE_SOURCE.to_string(),
4188                        DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
4189                        false,
4190                    ),
4191                    DynamicTypedFieldBindingRequest::new(
4192                        UPDATED_AT_SOURCE.to_string(),
4193                        DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
4194                        false,
4195                    ),
4196                ],
4197            )
4198            .expect("targeted typed binding should issue");
4199        let typed_patch = binding
4200            .bind_write_fields(vec![
4201                (
4202                    ID_SOURCE.to_string(),
4203                    DynamicWriteCell::Value(InputValue::Nat64(2)),
4204                ),
4205                (
4206                    PROFILE_SOURCE.to_string(),
4207                    DynamicWriteCell::Value(profile_input(12)),
4208                ),
4209            ])
4210            .expect("targeted typed patch should bind");
4211        let typed_error = session
4212            .execute_trusted_typed_mutation(
4213                &binding,
4214                &DynamicTypedMutation::Insert { patch: typed_patch },
4215            )
4216            .expect_err("typed write must enforce the targeted rule");
4217        assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
4218
4219        #[cfg(feature = "sql")]
4220        {
4221            let sql_error = session
4222                .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
4223                .expect_err("SQL default resolution must enforce the targeted rule");
4224            let crate::db::QueryError::Execute(execute) = sql_error else {
4225                panic!("targeted SQL write should fail at shared execution admission");
4226            };
4227            assert_eq!(
4228                targeted_constraint_id(execute.as_internal()),
4229                targeted_rule_id.get()
4230            );
4231        }
4232
4233        session
4234            .execute_trusted_dynamic_mutation_batch(vec![
4235                DynamicMutation::Insert {
4236                    entity: "TargetedMutation".to_string(),
4237                    patch: structural_patch(4, 5),
4238                },
4239                DynamicMutation::Insert {
4240                    entity: "TargetedMutation".to_string(),
4241                    patch: structural_patch(5, 12),
4242                },
4243            ])
4244            .expect_err("one invalid targeted value must reject the whole batch");
4245        assert_eq!(
4246            DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4247            Some(0),
4248            "no frontend or earlier valid batch row may escape targeted admission",
4249        );
4250
4251        let admitted = session
4252            .execute_trusted_dynamic_mutation_batch(vec![
4253                DynamicMutation::Insert {
4254                    entity: "TargetedMutation".to_string(),
4255                    patch: structural_patch(6, 5),
4256                },
4257                DynamicMutation::Insert {
4258                    entity: "TargetedMutation".to_string(),
4259                    patch: structural_patch(7, 10),
4260                },
4261            ])
4262            .expect("compliant targeted values should share one accepted batch");
4263        let [first, second] = admitted.rows.as_slice() else {
4264            panic!("the mixed targeted batch should return two rows");
4265        };
4266        let first_timestamp = first
4267            .get(2)
4268            .expect("the first mixed row should contain its managed timestamp");
4269        assert!(matches!(
4270            first_timestamp,
4271            crate::value::OutputValue::Timestamp(_)
4272        ));
4273        assert_eq!(
4274            second.get(2),
4275            Some(first_timestamp),
4276            "one accepted mixed batch must materialize one managed timestamp",
4277        );
4278        assert_eq!(
4279            DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
4280            Some(2),
4281        );
4282    }
4283}