Skip to main content

icydb_core/db/session/
write.rs

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