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::data::DecodedDataStoreKey;
2687    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2688    use crate::db::executor::budget::{
2689        HardExecutionBudget, HardExecutionContext, HardExecutionFailureHeadroom,
2690        with_query_execution_budget_for_tests,
2691    };
2692    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2693    use crate::db::{
2694        CompareProofAndAdvanceError, DynamicQuery, ExhaustiveReadError, RawDataStoreKey,
2695        ReadSetRevisionError, ResumableJobAdvance, ResumableJobAdvanceRequest,
2696        ResumableJobAdvanceStatus, ResumableJobError, ResumableJobId, ResumableJobIdempotencyKey,
2697        ResumableJobStatus, asc,
2698    };
2699    use crate::{
2700        db::{
2701            commit::{database_incarnation_id, forget_recovered_domain_for_tests},
2702            data::DataStore,
2703            executor::{MutationCommitInterruption, interrupt_next_mutation_commit_for_tests},
2704            index::IndexStore,
2705            integrity::{
2706                PhysicalUnitCheckpoint, QuickIntegrityStatus, RowInspectionLimits,
2707                execute_quick_integrity, execute_row_integrity_page,
2708            },
2709            journal::JournalTailStore,
2710            registry::{
2711                StoreAllocationIdentities, StoreAllocationIdentity, StoreRegistry,
2712                StoreRuntimeStorageCapabilities,
2713            },
2714            schema::{
2715                AcceptedFieldKind, AcceptedSchemaRevision, FieldId, FieldInsertGeneration,
2716                FieldStorageDecode, LeafCodec, PersistedFieldSnapshot,
2717                PersistedIndexFieldPathSnapshot, PersistedIndexKeySnapshot, PersistedIndexSnapshot,
2718                PersistedSchemaSnapshot, ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy,
2719                SchemaIndexId, SchemaInsertDefault, SchemaRowLayout, SchemaStore, SchemaVersion,
2720                accepted_schema_candidate_with_field_bindings_for_tests,
2721            },
2722            write_context::MutationMode,
2723        },
2724        error::{ErrorClass, ErrorOrigin, InternalError},
2725        testing::test_memory,
2726        traits::{CanisterKind, Path},
2727        types::{EntityTag, Timestamp},
2728        value::{InputValue, OutputValue, Value},
2729    };
2730    use icydb_schema::{FieldSourceKey, ScalarType};
2731    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2732    use std::cell::Cell;
2733    use std::{cell::RefCell, collections::BTreeMap, time::Instant};
2734
2735    const STORE_PATH: &str = "session::write::identity_pre_key_tests::Store";
2736    const ENTITY_SOURCE: &str = "session::write::identity_pre_key_tests::Entity";
2737    const ID_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::id";
2738    const PAYLOAD_SOURCE: &str = "session::write::identity_pre_key_tests::Entity::payload";
2739    const ENTITY_NAME: &str = "IdentityRow";
2740    const ENTITY_TAG: EntityTag = EntityTag::new(93);
2741    const JOURNALED_STORE_PATH: &str = "session::write::identity_pre_key_tests::JournaledStore";
2742    const UNRELATED_STORE_PATH: &str = "session::write::identity_pre_key_tests::UnrelatedStore";
2743
2744    struct TestCanister;
2745
2746    impl Path for TestCanister {
2747        const PATH: &'static str = "session::write::identity_pre_key_tests::Canister";
2748    }
2749
2750    impl CanisterKind for TestCanister {
2751        const COMMIT_MEMORY_ID: u8 = 45;
2752        const COMMIT_STABLE_KEY: &'static str = "icydb.identity_pre_key_tests.commit.v1";
2753        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 46;
2754        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2755            "icydb.identity_pre_key_tests.integrity.progress.v1";
2756    }
2757
2758    thread_local! {
2759        static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
2760        static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
2761        static SCHEMA_STORE: RefCell<SchemaStore> =
2762            const { RefCell::new(SchemaStore::init_heap()) };
2763        static UNRELATED_DATA_STORE: RefCell<DataStore> =
2764            const { RefCell::new(DataStore::init_heap()) };
2765        static UNRELATED_INDEX_STORE: RefCell<IndexStore> =
2766            const { RefCell::new(IndexStore::init_heap()) };
2767        static UNRELATED_SCHEMA_STORE: RefCell<SchemaStore> =
2768            const { RefCell::new(SchemaStore::init_heap()) };
2769        static STORE_REGISTRY: StoreRegistry = {
2770            let mut registry = StoreRegistry::new();
2771            registry.register_store(
2772                STORE_PATH,
2773                &DATA_STORE,
2774                &INDEX_STORE,
2775                &SCHEMA_STORE,
2776                StoreAllocationIdentities::absent(),
2777                StoreRuntimeStorageCapabilities::heap(),
2778            ).expect("identity pre-key test store should register");
2779            registry.register_store(
2780                UNRELATED_STORE_PATH,
2781                &UNRELATED_DATA_STORE,
2782                &UNRELATED_INDEX_STORE,
2783                &UNRELATED_SCHEMA_STORE,
2784                StoreAllocationIdentities::absent(),
2785                StoreRuntimeStorageCapabilities::heap(),
2786            ).expect("unrelated identity test store should register");
2787            registry
2788        };
2789        static JOURNALED_DATA_STORE: RefCell<DataStore> =
2790            RefCell::new(DataStore::init_journaled(test_memory(186)));
2791        static JOURNALED_INDEX_STORE: RefCell<IndexStore> =
2792            RefCell::new(IndexStore::init_journaled(test_memory(187)));
2793        static JOURNALED_SCHEMA_STORE: RefCell<SchemaStore> =
2794            RefCell::new(SchemaStore::init_journaled(test_memory(188)));
2795        static JOURNALED_TAIL_STORE: RefCell<JournalTailStore> =
2796            RefCell::new(JournalTailStore::init(test_memory(189)));
2797        static JOURNALED_STORE_REGISTRY: StoreRegistry = {
2798            let mut registry = StoreRegistry::new();
2799            registry.register_journaled_store(
2800                JOURNALED_STORE_PATH,
2801                &JOURNALED_DATA_STORE,
2802                &JOURNALED_INDEX_STORE,
2803                &JOURNALED_SCHEMA_STORE,
2804                &JOURNALED_TAIL_STORE,
2805                StoreAllocationIdentities::new_journaled(
2806                    StoreAllocationIdentity::new(186, "icydb.test.identity-range.data.v1"),
2807                    StoreAllocationIdentity::new(187, "icydb.test.identity-range.index.v1"),
2808                    StoreAllocationIdentity::new(188, "icydb.test.identity-range.schema.v1"),
2809                    StoreAllocationIdentity::new(189, "icydb.test.identity-range.journal.v1"),
2810                ),
2811                StoreRuntimeStorageCapabilities::journaled(),
2812            ).expect("identity range journaled store should register");
2813            registry
2814        };
2815    }
2816
2817    struct JournaledTestCanister;
2818
2819    impl Path for JournaledTestCanister {
2820        const PATH: &'static str = "session::write::identity_pre_key_tests::JournaledCanister";
2821    }
2822
2823    impl CanisterKind for JournaledTestCanister {
2824        const COMMIT_MEMORY_ID: u8 = 190;
2825        const COMMIT_STABLE_KEY: &'static str = "icydb.identity_range_tests.commit.v1";
2826        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 191;
2827        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
2828            "icydb.identity_range_tests.integrity.progress.v1";
2829    }
2830
2831    fn source_key(source: &str) -> FieldSourceKey {
2832        FieldSourceKey::try_new(source).expect("identity test field source should admit")
2833    }
2834
2835    fn identity_snapshot(store_path: &str) -> PersistedSchemaSnapshot {
2836        let fields = vec![
2837            PersistedFieldSnapshot::new_initial_with_write_policy(
2838                FieldId::new(1),
2839                "id".to_string(),
2840                SchemaFieldSlot::new(0),
2841                AcceptedFieldKind::Nat64,
2842                Vec::new(),
2843                false,
2844                SchemaInsertDefault::None,
2845                SchemaFieldWritePolicy::from_model_policies(
2846                    Some(FieldInsertGeneration::Identity),
2847                    None,
2848                ),
2849                FieldStorageDecode::ByKind,
2850                LeafCodec::Scalar(ScalarCodec::Nat64),
2851            ),
2852            PersistedFieldSnapshot::new_initial(
2853                FieldId::new(2),
2854                "payload".to_string(),
2855                SchemaFieldSlot::new(1),
2856                AcceptedFieldKind::Nat64,
2857                Vec::new(),
2858                false,
2859                SchemaInsertDefault::None,
2860                FieldStorageDecode::ByKind,
2861                LeafCodec::Scalar(ScalarCodec::Nat64),
2862            ),
2863        ];
2864        PersistedSchemaSnapshot::new_with_indexes(
2865            SchemaVersion::initial(),
2866            ENTITY_SOURCE.to_string(),
2867            ENTITY_NAME.to_string(),
2868            FieldId::new(1),
2869            SchemaRowLayout::initial(
2870                fields
2871                    .iter()
2872                    .map(|field| (field.id(), field.slot()))
2873                    .collect(),
2874            ),
2875            fields,
2876            vec![PersistedIndexSnapshot::new(
2877                SchemaIndexId::new(1).expect("identity test index ID should admit"),
2878                1,
2879                "by_payload".to_string(),
2880                store_path.to_string(),
2881                false,
2882                PersistedIndexKeySnapshot::FieldPath(vec![PersistedIndexFieldPathSnapshot::new(
2883                    FieldId::new(2),
2884                    SchemaFieldSlot::new(1),
2885                    vec!["payload".to_string()],
2886                    AcceptedFieldKind::Nat64,
2887                    false,
2888                )]),
2889                None,
2890            )],
2891        )
2892    }
2893
2894    fn initialize() -> DbSession<TestCanister> {
2895        DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2896        INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2897        SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2898        UNRELATED_DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
2899        UNRELATED_INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
2900        UNRELATED_SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
2901        let session = DbSession::<TestCanister>::new(
2902            &STORE_REGISTRY,
2903            &crate::db::RequestExecutionRoot::__new_runtime_root(),
2904        );
2905        session
2906            .db
2907            .ensure_recovered_state()
2908            .expect("identity pre-key test database should initialize");
2909        let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2910            STORE_PATH,
2911            AcceptedSchemaRevision::INITIAL,
2912            BTreeMap::from([(ENTITY_TAG, identity_snapshot(STORE_PATH))]),
2913            BTreeMap::from([
2914                ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2915                ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2916            ]),
2917        );
2918        let store = session
2919            .db
2920            .store_handle(STORE_PATH)
2921            .expect("identity pre-key test store should resolve");
2922        crate::db::commit::publish_accepted_schema_candidate(
2923            STORE_PATH,
2924            store,
2925            AcceptedSchemaRevision::NONE,
2926            &candidate,
2927        )
2928        .expect("identity candidate should publish with explicit zero state");
2929        session
2930    }
2931
2932    fn initialize_journaled_with_root() -> (
2933        DbSession<JournaledTestCanister>,
2934        crate::db::RequestExecutionRoot,
2935    ) {
2936        let root = crate::db::RequestExecutionRoot::__new_runtime_root();
2937        let session = DbSession::<JournaledTestCanister>::new(&JOURNALED_STORE_REGISTRY, &root);
2938        session
2939            .db
2940            .ensure_recovered_state()
2941            .expect("journaled identity database should initialize");
2942        let candidate = accepted_schema_candidate_with_field_bindings_for_tests(
2943            JOURNALED_STORE_PATH,
2944            AcceptedSchemaRevision::INITIAL,
2945            BTreeMap::from([(ENTITY_TAG, identity_snapshot(JOURNALED_STORE_PATH))]),
2946            BTreeMap::from([
2947                ((ENTITY_TAG, source_key(ID_SOURCE)), FieldId::new(1)),
2948                ((ENTITY_TAG, source_key(PAYLOAD_SOURCE)), FieldId::new(2)),
2949            ]),
2950        );
2951        let store = session
2952            .db
2953            .store_handle(JOURNALED_STORE_PATH)
2954            .expect("journaled identity store should resolve");
2955        crate::db::commit::publish_accepted_schema_candidate(
2956            JOURNALED_STORE_PATH,
2957            store,
2958            AcceptedSchemaRevision::NONE,
2959            &candidate,
2960        )
2961        .expect("journaled identity candidate should publish");
2962        (session, root)
2963    }
2964
2965    fn initialize_journaled() -> DbSession<JournaledTestCanister> {
2966        initialize_journaled_with_root().0
2967    }
2968
2969    fn payload_patch(value: u64) -> AcceptedMutationIntentPatch {
2970        AcceptedMutationIntentPatch::new()
2971            .set_authored(FieldSlot::from_validated_index(1), InputValue::Nat64(value))
2972    }
2973
2974    fn dynamic_payload_patch(value: u64) -> DynamicStructuralPatch {
2975        DynamicStructuralPatch::new(vec![(
2976            "payload".to_string(),
2977            DynamicWriteCell::Value(InputValue::Nat64(value)),
2978        )])
2979    }
2980
2981    fn expected_dynamic_row(id: u64, payload: u64) -> Vec<OutputValue> {
2982        vec![OutputValue::Nat64(id), OutputValue::Nat64(payload)]
2983    }
2984
2985    #[cfg(all(feature = "sql", feature = "diagnostics"))]
2986    fn exact_key_binding<C: CanisterKind>(session: &DbSession<C>) -> DynamicTypedEntityBinding {
2987        session
2988            .issue_typed_entity_binding(
2989                ENTITY_SOURCE,
2990                &[
2991                    DynamicTypedFieldBindingRequest::new(
2992                        ID_SOURCE.to_string(),
2993                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2994                        false,
2995                    ),
2996                    DynamicTypedFieldBindingRequest::new(
2997                        PAYLOAD_SOURCE.to_string(),
2998                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
2999                        false,
3000                    ),
3001                ],
3002            )
3003            .expect("exact-key test binding should issue")
3004    }
3005
3006    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3007    fn insert_exact_key_fixture<C: CanisterKind>(session: &DbSession<C>, payload: u64) -> u64 {
3008        let output = session
3009            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
3010                entity: ENTITY_NAME.to_string(),
3011                patch: dynamic_payload_patch(payload),
3012            })
3013            .expect("exact-key fixture insert should commit");
3014        match output.rows.as_slice() {
3015            [row] => match row.as_slice() {
3016                [OutputValue::Nat64(id), OutputValue::Nat64(actual_payload)]
3017                    if *actual_payload == payload =>
3018                {
3019                    *id
3020                }
3021                _ => panic!("exact-key fixture should return its identity and payload"),
3022            },
3023            _ => panic!("exact-key fixture insert should return one row"),
3024        }
3025    }
3026
3027    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3028    fn identity_row_stored_bytes<C: CanisterKind>(
3029        session: &DbSession<C>,
3030        store_path: &'static str,
3031        key: u64,
3032    ) -> u64 {
3033        let data_key = DecodedDataStoreKey::try_from_structural_key(ENTITY_TAG, &Value::Nat64(key))
3034            .expect("identity row key should encode");
3035        let raw_key = data_key.to_raw().expect("identity raw key should encode");
3036        let store = session
3037            .db
3038            .recovered_store(store_path)
3039            .expect("identity store should resolve");
3040        store.with_data(|data_store| {
3041            u64::try_from(
3042                data_store
3043                    .get(&raw_key)
3044                    .expect("inserted identity row should exist")
3045                    .len(),
3046            )
3047            .expect("bounded row length should fit u64")
3048        })
3049    }
3050
3051    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3052    fn with_stored_bytes_limit<T>(
3053        limit: u64,
3054        shape_fingerprint_prefix: u64,
3055        operation: impl FnOnce() -> Result<T, crate::db::query::intent::QueryError>,
3056    ) -> Result<T, crate::db::query::intent::QueryError> {
3057        let budget = HardExecutionBudget::uniform_for_tests(
3058            u64::MAX,
3059            HardExecutionFailureHeadroom::new(500, 256),
3060        )
3061        .with_limit_for_tests(
3062            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::StoredBytesRead,
3063            limit,
3064        );
3065        let context = HardExecutionContext::new(
3066            icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3067            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3068            shape_fingerprint_prefix,
3069        );
3070
3071        with_query_execution_budget_for_tests(budget, context, operation)
3072    }
3073
3074    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3075    fn assert_exact_key_batch<C: CanisterKind>(session: &DbSession<C>) {
3076        let first = insert_exact_key_fixture(session, 41);
3077        let second = insert_exact_key_fixture(session, 42);
3078        let missing = u64::MAX;
3079        let binding = exact_key_binding(session);
3080        let gets_before = DataStore::current_get_call_count();
3081        let result = session
3082            .execute_public_exact_key_batch_for_typed_binding(
3083                &binding,
3084                &[second, missing, first, second],
3085            )
3086            .expect("exact-key batch should execute")
3087            .expect("exact-key binding should remain current");
3088
3089        assert_eq!(result.positions, vec![0, 1, 2, 0]);
3090        assert_eq!(
3091            result.distinct_rows,
3092            vec![
3093                Some(expected_dynamic_row(second, 42)),
3094                None,
3095                Some(expected_dynamic_row(first, 41)),
3096            ],
3097        );
3098        assert_eq!(
3099            DataStore::current_get_call_count().saturating_sub(gets_before),
3100            3,
3101            "four input positions with one duplicate must perform three physical reads",
3102        );
3103    }
3104
3105    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3106    #[test]
3107    fn exact_key_batches_preserve_semantics_across_heap_and_journaled_stores() {
3108        assert_exact_key_batch(&initialize());
3109        assert_exact_key_batch(&initialize_journaled());
3110    }
3111
3112    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3113    fn assert_primary_range_materialization_fetches_once<C: CanisterKind>(
3114        session: &DbSession<C>,
3115        store_path: &'static str,
3116    ) {
3117        let key = insert_exact_key_fixture(session, 41);
3118        let stored_bytes = identity_row_stored_bytes(session, store_path, key);
3119
3120        let scalar = DynamicQuery::new(ENTITY_NAME)
3121            .select(["id", "payload"])
3122            .order_by(asc("id"))
3123            .limit(1);
3124        let gets_before = DataStore::current_get_call_count();
3125        let scalar_page = with_stored_bytes_limit(stored_bytes, 0x7072_696d_6172_792d, || {
3126            session.execute_trusted_live_page(&scalar, None)
3127        })
3128        .expect("one scalar primary-range row should fit one payload-read allowance");
3129        assert_eq!(scalar_page.row_count, 1);
3130        assert_eq!(
3131            DataStore::current_get_call_count().saturating_sub(gets_before),
3132            1,
3133            "scalar primary traversal should fetch its emitted row exactly once",
3134        );
3135
3136        let grouped = DynamicQuery::new(ENTITY_NAME)
3137            .group_by("payload")
3138            .aggregate(crate::db::count())
3139            .grouped_limits(10, 16 * 1_024)
3140            .limit(1);
3141        let gets_before = DataStore::current_get_call_count();
3142        let grouped_page = with_stored_bytes_limit(stored_bytes, 0x6772_6f75_7065_642d, || {
3143            session.execute_trusted_dynamic_grouped_query(&grouped)
3144        })
3145        .expect("one grouped primary-range row should fit one payload-read allowance");
3146        assert_eq!(grouped_page.row_count, 1);
3147        assert_eq!(
3148            DataStore::current_get_call_count().saturating_sub(gets_before),
3149            1,
3150            "grouped primary traversal should fetch its source row exactly once",
3151        );
3152    }
3153
3154    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3155    #[test]
3156    fn row_materialization_fetches_each_required_payload_at_most_once() {
3157        assert_primary_range_materialization_fetches_once(&initialize(), STORE_PATH);
3158        assert_primary_range_materialization_fetches_once(
3159            &initialize_journaled(),
3160            JOURNALED_STORE_PATH,
3161        );
3162    }
3163
3164    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3165    #[test]
3166    fn exhaustive_pages_require_and_recompare_the_complete_source_proof() {
3167        let session = initialize();
3168        let first = insert_exact_key_fixture(&session, 41);
3169        let second = insert_exact_key_fixture(&session, 42);
3170        let third = insert_exact_key_fixture(&session, 43);
3171        let query = DynamicQuery::new(ENTITY_NAME)
3172            .select(["id", "payload"])
3173            .order_by(asc("id"));
3174
3175        let page = session
3176            .execute_trusted_exhaustive_page(&query, None, None)
3177            .expect("initial exhaustive page should capture its source proof");
3178        assert_eq!(
3179            page.rows,
3180            vec![
3181                expected_dynamic_row(first, 41),
3182                expected_dynamic_row(second, 42),
3183            ],
3184        );
3185        let continuation = page
3186            .continuation
3187            .as_deref()
3188            .expect("unreturned row should retain exhaustive continuation");
3189        assert!(matches!(
3190            session.execute_trusted_exhaustive_page(&query, Some(continuation), None),
3191            Err(ExhaustiveReadError::Revision(
3192                ReadSetRevisionError::ResumeProofRequired
3193            )),
3194        ));
3195        let resumed = session
3196            .execute_trusted_exhaustive_page(&query, Some(continuation), Some(&page.proof))
3197            .expect("unchanged proof should resume exhaustive traversal");
3198        assert_eq!(resumed.rows, vec![expected_dynamic_row(third, 43)]);
3199        assert_eq!(resumed.continuation, None);
3200
3201        let stale_page = session
3202            .execute_trusted_exhaustive_page(&query, None, None)
3203            .expect("fresh exhaustive page should capture current revision");
3204        let stale_continuation = stale_page
3205            .continuation
3206            .as_deref()
3207            .expect("fresh three-row traversal should retain continuation");
3208        let _ = insert_exact_key_fixture(&session, 44);
3209        assert!(matches!(
3210            session.execute_trusted_exhaustive_page(
3211                &query,
3212                Some(stale_continuation),
3213                Some(&stale_page.proof),
3214            ),
3215            Err(ExhaustiveReadError::Revision(
3216                ReadSetRevisionError::StoreDataChanged { .. }
3217            )),
3218        ));
3219    }
3220
3221    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3222    #[test]
3223    fn heap_sources_cannot_back_durable_resumable_jobs() {
3224        let session = initialize();
3225        let proof = session
3226            .capture_read_set_revision_proof(&[ENTITY_NAME])
3227            .expect("heap source proof should capture for one-call exhaustive reads");
3228        let job_id = ResumableJobId::try_from_bytes([70; 32])
3229            .expect("nonzero heap test job identity should admit");
3230
3231        assert!(matches!(
3232            session.start_resumable_job(job_id, proof, Vec::new()),
3233            Err(ResumableJobError::SourceProof(
3234                ReadSetRevisionError::DurableStoreRequired { .. }
3235            )),
3236        ));
3237    }
3238
3239    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3240    #[test]
3241    fn proof_and_progress_controls_charge_one_shared_request_scope() {
3242        let (session, root) = initialize_journaled_with_root();
3243        let resource = icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions;
3244        let before = root.observed(resource);
3245        let proof = session
3246            .capture_read_set_revision_proof(&[ENTITY_NAME])
3247            .expect("proof capture should use the retained request scope");
3248        let job_id = ResumableJobId::try_from_bytes([75; 32])
3249            .expect("nonzero accounting job identity should admit");
3250        session
3251            .start_resumable_job(job_id, proof, Vec::new())
3252            .expect("job start should use the same retained request scope");
3253        let _ = session
3254            .resumable_job_state(job_id)
3255            .expect("job load should use the same retained request scope");
3256
3257        assert_eq!(root.observed(resource).saturating_sub(before), 3);
3258    }
3259
3260    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3261    #[test]
3262    fn source_proofs_ignore_unrelated_stores_but_bind_access_state_changes() {
3263        let session = initialize();
3264        let proof = session
3265            .capture_read_set_revision_proof(&[ENTITY_NAME])
3266            .expect("source proof should cover only the entity's physical store");
3267        let shared_store_proof = session
3268            .capture_read_set_revision_proof(&[ENTITY_NAME, ENTITY_NAME])
3269            .expect("entities sharing one physical source should deduplicate");
3270        assert_eq!(shared_store_proof, proof);
3271        assert_eq!(shared_store_proof.stores().len(), 1);
3272        let unrelated = session
3273            .db
3274            .store_handle(UNRELATED_STORE_PATH)
3275            .expect("unrelated registered store should resolve");
3276        unrelated.with_data_mut(|store| {
3277            let _ = store.remove(&RawDataStoreKey::from_persisted_bytes(vec![1]));
3278        });
3279        session
3280            .verify_read_set_revision_proof(&proof)
3281            .expect("a nonparticipating store mutation must not invalidate the proof");
3282
3283        let source = session
3284            .db
3285            .store_handle(STORE_PATH)
3286            .expect("participating source store should resolve");
3287        source
3288            .mark_index_building()
3289            .expect("source access-state transition should advance its revision");
3290        assert!(matches!(
3291            session.verify_read_set_revision_proof(&proof),
3292            Err(ExhaustiveReadError::Revision(
3293                ReadSetRevisionError::StoreAccessChanged { .. }
3294            )),
3295        ));
3296    }
3297
3298    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3299    #[expect(
3300        clippy::too_many_lines,
3301        reason = "one lifecycle test proves successful replay plus pre-page and post-page source invalidation without sharing progress state across tests"
3302    )]
3303    #[test]
3304    fn journaled_job_advance_is_idempotent_and_revision_checked_on_both_sides() {
3305        let session = initialize_journaled();
3306        let proof = session
3307            .capture_read_set_revision_proof(&[ENTITY_NAME])
3308            .expect("journaled source proof should capture");
3309        let job_id =
3310            ResumableJobId::try_from_bytes([71; 32]).expect("nonzero job identity should admit");
3311        session
3312            .start_resumable_job(job_id, proof, vec![0])
3313            .expect("journaled job should start outside its protected source revision");
3314        let request = ResumableJobAdvanceRequest::new(
3315            job_id,
3316            0,
3317            ResumableJobIdempotencyKey::new("page-0")
3318                .expect("bounded idempotency key should admit"),
3319        );
3320        let calls = Cell::new(0_u8);
3321        let receipt = session
3322            .compare_proof_and_advance(&request, |state| {
3323                calls.set(calls.get() + 1);
3324                assert_eq!(state.application_state, vec![0]);
3325                Ok::<_, ()>(
3326                    ResumableJobAdvance::new(Some("cursor-1".to_string()), vec![1], vec![9])
3327                        .expect("bounded application advance should admit"),
3328                )
3329            })
3330            .expect("unchanged source should advance exactly once");
3331        assert_eq!(calls.get(), 1);
3332        assert_eq!(receipt.status, ResumableJobAdvanceStatus::Advanced);
3333        assert_eq!(receipt.committed_sequence, 1);
3334
3335        let replay = session
3336            .compare_proof_and_advance::<()>(&request, |_| {
3337                panic!("lost-response replay must not execute application work")
3338            })
3339            .expect("same request identity should return its persisted receipt");
3340        assert_eq!(replay, receipt);
3341        let retained = session
3342            .resumable_job_state(job_id)
3343            .expect("advanced state should remain durable");
3344        assert_eq!(retained.sequence, 1);
3345        assert_eq!(retained.application_state, vec![1]);
3346
3347        let _ = insert_exact_key_fixture(&session, 51);
3348        let pre_change_request = ResumableJobAdvanceRequest::new(
3349            job_id,
3350            1,
3351            ResumableJobIdempotencyKey::new("page-1")
3352                .expect("bounded idempotency key should admit"),
3353        );
3354        let pre_change_calls = Cell::new(0_u8);
3355        let invalidated = session
3356            .compare_proof_and_advance::<()>(&pre_change_request, |_| {
3357                pre_change_calls.set(pre_change_calls.get() + 1);
3358                unreachable!("pre-page proof failure must reject before application work")
3359            })
3360            .expect("source drift should persist one replayable invalidation receipt");
3361        assert_eq!(pre_change_calls.get(), 0);
3362        assert_eq!(invalidated.status, ResumableJobAdvanceStatus::Invalidated);
3363        let invalidated_state = session
3364            .resumable_job_state(job_id)
3365            .expect("invalidated job should remain inspectable");
3366        assert_eq!(invalidated_state.status, ResumableJobStatus::Invalidated);
3367        assert_eq!(invalidated_state.continuation, None);
3368        assert_eq!(invalidated_state.application_state, vec![1]);
3369        assert_eq!(
3370            session
3371                .compare_proof_and_advance::<()>(&pre_change_request, |_| {
3372                    panic!("invalidation replay must not execute application work")
3373                })
3374                .expect("lost invalidation reply should replay exactly"),
3375            invalidated,
3376        );
3377
3378        let post_proof = session
3379            .capture_read_set_revision_proof(&[ENTITY_NAME])
3380            .expect("post-change journaled proof should capture");
3381        let post_job_id = ResumableJobId::try_from_bytes([72; 32])
3382            .expect("nonzero post-change job identity should admit");
3383        session
3384            .start_resumable_job(post_job_id, post_proof, vec![7])
3385            .expect("post-change journaled job should start");
3386        let post_request = ResumableJobAdvanceRequest::new(
3387            post_job_id,
3388            0,
3389            ResumableJobIdempotencyKey::new("post-page-0")
3390                .expect("bounded idempotency key should admit"),
3391        );
3392        let post_receipt = session
3393            .compare_proof_and_advance::<()>(&post_request, |_| {
3394                let _ = insert_exact_key_fixture(&session, 52);
3395                Ok(ResumableJobAdvance::new(None, vec![8], vec![10])
3396                    .expect("bounded post-change candidate should admit"))
3397            })
3398            .expect("post-page drift should discard the candidate and persist invalidation");
3399        assert_eq!(post_receipt.status, ResumableJobAdvanceStatus::Invalidated);
3400        let post_state = session
3401            .resumable_job_state(post_job_id)
3402            .expect("post-page invalidation should remain inspectable");
3403        assert_eq!(post_state.status, ResumableJobStatus::Invalidated);
3404        assert_eq!(post_state.application_state, vec![7]);
3405        session
3406            .acknowledge_resumable_job(post_job_id, post_state.sequence)
3407            .expect("terminal job acknowledgement should remove retained progress");
3408        session
3409            .acknowledge_resumable_job(post_job_id, post_state.sequence)
3410            .expect("lost acknowledgement reply should be safely replayable");
3411        assert_eq!(
3412            session.resumable_job_state(post_job_id),
3413            Err(ResumableJobError::NotFound),
3414        );
3415
3416        let completed_job_id = ResumableJobId::try_from_bytes([74; 32])
3417            .expect("nonzero completed job identity should admit");
3418        let completed_proof = session
3419            .capture_read_set_revision_proof(&[ENTITY_NAME])
3420            .expect("completed-job source proof should capture");
3421        session
3422            .start_resumable_job(completed_job_id, completed_proof, Vec::new())
3423            .expect("completed-job fixture should start");
3424        let completed_request = ResumableJobAdvanceRequest::new(
3425            completed_job_id,
3426            0,
3427            ResumableJobIdempotencyKey::new("complete")
3428                .expect("bounded completion key should admit"),
3429        );
3430        let completed_receipt = session
3431            .compare_proof_and_advance::<()>(&completed_request, |_| {
3432                Ok(ResumableJobAdvance::new(None, vec![99], vec![100])
3433                    .expect("bounded terminal advance should admit"))
3434            })
3435            .expect("null continuation should commit terminal completion");
3436        let completed_state = session
3437            .resumable_job_state(completed_job_id)
3438            .expect("completed state should remain replayable before acknowledgement");
3439        assert_eq!(completed_state.status, ResumableJobStatus::Completed);
3440        assert_eq!(
3441            session
3442                .compare_proof_and_advance::<()>(&completed_request, |_| {
3443                    panic!("completed request replay must not execute application work")
3444                })
3445                .expect("completed request should replay until acknowledgement"),
3446            completed_receipt,
3447        );
3448        let after_completion = ResumableJobAdvanceRequest::new(
3449            completed_job_id,
3450            1,
3451            ResumableJobIdempotencyKey::new("after-complete")
3452                .expect("bounded post-completion key should admit"),
3453        );
3454        assert!(matches!(
3455            session.compare_proof_and_advance::<()>(&after_completion, |_| {
3456                panic!("completed jobs cannot execute another page")
3457            }),
3458            Err(CompareProofAndAdvanceError::Protocol(
3459                ResumableJobError::Completed
3460            )),
3461        ));
3462        session
3463            .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
3464            .expect("completed job should acknowledge and free capacity");
3465        session
3466            .acknowledge_resumable_job(completed_job_id, completed_state.sequence)
3467            .expect("completion acknowledgement should be idempotent");
3468
3469        let stale_job_id = ResumableJobId::try_from_bytes([73; 32])
3470            .expect("nonzero stale-sequence job identity should admit");
3471        let stale_proof = session
3472            .capture_read_set_revision_proof(&[ENTITY_NAME])
3473            .expect("stale-sequence source proof should capture");
3474        session
3475            .start_resumable_job(stale_job_id, stale_proof, Vec::new())
3476            .expect("stale-sequence job should start");
3477        let stale_request = ResumableJobAdvanceRequest::new(
3478            stale_job_id,
3479            4,
3480            ResumableJobIdempotencyKey::new("stale").expect("bounded idempotency key should admit"),
3481        );
3482        assert!(matches!(
3483            session.compare_proof_and_advance::<()>(&stale_request, |_| {
3484                panic!("stale sequence must reject before application work")
3485            }),
3486            Err(CompareProofAndAdvanceError::Protocol(
3487                ResumableJobError::StaleSequence {
3488                    expected: 4,
3489                    actual: 0,
3490                }
3491            )),
3492        ));
3493        assert_eq!(
3494            session.acknowledge_resumable_job(stale_job_id, 0),
3495            Err(ResumableJobError::NotTerminal),
3496        );
3497    }
3498
3499    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3500    #[test]
3501    fn exact_key_batch_uses_typed_hard_execution_budget() {
3502        let session = initialize();
3503        let binding = exact_key_binding(&session);
3504        let budget =
3505            HardExecutionBudget::uniform_for_tests(0, HardExecutionFailureHeadroom::new(500, 256));
3506        let error = session
3507            .execute_exact_key_batch_with_hard_budget_for_tests(&binding, &[u64::MAX], &budget)
3508            .expect_err("zero query budget should reject the exact-key route");
3509
3510        assert!(matches!(
3511            error.diagnostic().detail(),
3512            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3513                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3514            })
3515        ));
3516        let facts = error.diagnostic_facts();
3517        assert_eq!(
3518            &facts[..5],
3519            &[
3520                (
3521                    icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3522                    icydb_diagnostic_code::DiagnosticExecutionBudgetResource::QueryExecutions.raw(),
3523                ),
3524                (icydb_diagnostic_code::DiagnosticFactTag::Limit, 0),
3525                (icydb_diagnostic_code::DiagnosticFactTag::Actual, 1),
3526                (
3527                    icydb_diagnostic_code::DiagnosticFactTag::ExecutionBudgetScope,
3528                    icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution.raw(),
3529                ),
3530                (
3531                    icydb_diagnostic_code::DiagnosticFactTag::ExecutionLane,
3532                    icydb_diagnostic_code::DiagnosticExecutionLane::PublicRead.raw(),
3533                ),
3534            ],
3535        );
3536        assert_eq!(
3537            facts[5].0,
3538            icydb_diagnostic_code::DiagnosticFactTag::QueryShapeFingerprintPrefix,
3539        );
3540        assert_ne!(facts[5].1, 0);
3541    }
3542
3543    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3544    fn assert_planned_query_exhausts(
3545        session: &DbSession<TestCanister>,
3546        query: &crate::db::DynamicQuery,
3547        resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3548    ) {
3549        let budget = HardExecutionBudget::uniform_for_tests(
3550            u64::MAX,
3551            HardExecutionFailureHeadroom::new(500, 256),
3552        )
3553        .with_limit_for_tests(resource, 0);
3554        let context = HardExecutionContext::new(
3555            icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3556            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3557            0x7068_7973_6963_616c,
3558        );
3559        let error = with_query_execution_budget_for_tests(budget, context, || {
3560            session.execute_trusted_live_page(query, None)
3561        })
3562        .expect_err("the injected zero resource allowance should reject planned execution");
3563
3564        assert!(matches!(
3565            error.diagnostic().detail(),
3566            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3567                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3568            })
3569        ));
3570        assert_eq!(
3571            error.diagnostic_facts()[0],
3572            (
3573                icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3574                resource.raw(),
3575            ),
3576        );
3577    }
3578
3579    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3580    fn assert_grouped_query_exhausts(
3581        session: &DbSession<TestCanister>,
3582        query: &crate::db::DynamicQuery,
3583        resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3584    ) {
3585        let budget = HardExecutionBudget::uniform_for_tests(
3586            u64::MAX,
3587            HardExecutionFailureHeadroom::new(500, 256),
3588        )
3589        .with_limit_for_tests(resource, 0);
3590        let context = HardExecutionContext::new(
3591            icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3592            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3593            0x6772_6f75_7065_642d,
3594        );
3595        let error = with_query_execution_budget_for_tests(budget, context, || {
3596            session.execute_trusted_dynamic_grouped_query(query)
3597        })
3598        .expect_err("the injected zero resource allowance should reject grouped execution");
3599
3600        assert!(matches!(
3601            error.diagnostic().detail(),
3602            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3603                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3604            })
3605        ));
3606        assert_eq!(
3607            error.diagnostic_facts()[0],
3608            (
3609                icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3610                resource.raw(),
3611            ),
3612        );
3613    }
3614
3615    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3616    fn assert_sql_query_exhausts(
3617        session: &DbSession<TestCanister>,
3618        sql: &str,
3619        resource: icydb_diagnostic_code::DiagnosticExecutionBudgetResource,
3620    ) {
3621        let budget = HardExecutionBudget::uniform_for_tests(
3622            u64::MAX,
3623            HardExecutionFailureHeadroom::new(500, 256),
3624        )
3625        .with_limit_for_tests(resource, 0);
3626        let context = HardExecutionContext::new(
3627            icydb_diagnostic_code::DiagnosticExecutionBudgetScope::Execution,
3628            icydb_diagnostic_code::DiagnosticExecutionLane::TrustedRead,
3629            0x7371_6c2d_736f_7274,
3630        );
3631        let error = with_query_execution_budget_for_tests(budget, context, || {
3632            session.execute_trusted_sql_query(sql)
3633        })
3634        .expect_err("the injected zero resource allowance should reject SQL execution");
3635
3636        assert!(matches!(
3637            error.diagnostic().detail(),
3638            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3639                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::ExecutionBudgetExceeded,
3640            })
3641        ));
3642        assert_eq!(
3643            error.diagnostic_facts()[0],
3644            (
3645                icydb_diagnostic_code::DiagnosticFactTag::BudgetResource,
3646                resource.raw(),
3647            ),
3648        );
3649    }
3650
3651    #[cfg(all(feature = "sql", feature = "diagnostics"))]
3652    #[test]
3653    fn planned_read_routes_share_physical_resource_accounting() {
3654        let session = initialize();
3655        let first = insert_exact_key_fixture(&session, 41);
3656        insert_exact_key_fixture(&session, 42);
3657
3658        let fallback = crate::db::DynamicQuery::new(ENTITY_NAME)
3659            .filter(crate::db::FieldRef::new("id").eq(first))
3660            .select(["id", "payload"])
3661            .order_by(crate::db::asc("id"))
3662            .limit(1);
3663        assert_eq!(
3664            session
3665                .execute_trusted_live_page(&fallback, None)
3666                .expect("bounded fallback execution should preserve its result")
3667                .row_count,
3668            1,
3669        );
3670        assert_planned_query_exhausts(
3671            &session,
3672            &fallback,
3673            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::RowsVisited,
3674        );
3675
3676        let covering = crate::db::DynamicQuery::new(ENTITY_NAME)
3677            .filter(crate::db::FieldRef::new("payload").eq(41_u64))
3678            .select(["payload"])
3679            .order_by(crate::db::asc("payload"))
3680            .limit(1);
3681        assert_eq!(
3682            session
3683                .execute_trusted_live_page(&covering, None)
3684                .expect("bounded covering execution should preserve its result")
3685                .row_count,
3686            1,
3687        );
3688        assert_planned_query_exhausts(
3689            &session,
3690            &covering,
3691            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
3692        );
3693
3694        let residual = crate::db::DynamicQuery::new(ENTITY_NAME)
3695            .filter(crate::db::FieldRef::new("payload").eq_field("id"))
3696            .select(["id"])
3697            .order_by(crate::db::asc("id"))
3698            .limit(1);
3699        assert_eq!(
3700            session
3701                .execute_trusted_live_page(&residual, None)
3702                .expect("bounded residual execution should preserve its result")
3703                .row_count,
3704            0,
3705        );
3706        assert_planned_query_exhausts(
3707            &session,
3708            &residual,
3709            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::PredicateExpressionSteps,
3710        );
3711
3712        assert_planned_query_exhausts(
3713            &session,
3714            &fallback,
3715            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::ResultBytes,
3716        );
3717
3718        let grouped = crate::db::DynamicQuery::new(ENTITY_NAME)
3719            .group_by("payload")
3720            .aggregate(crate::db::count())
3721            .order_by(crate::db::asc("payload"))
3722            .grouped_limits(10, 16 * 1_024)
3723            .limit(1);
3724        let grouped_result = session
3725            .execute_trusted_dynamic_grouped_query(&grouped)
3726            .expect("bounded grouped execution should preserve its result");
3727        assert_eq!(grouped_result.row_count, 1);
3728        assert!(grouped_result.next_cursor.is_some());
3729        assert_grouped_query_exhausts(
3730            &session,
3731            &grouped,
3732            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::GroupDistinctEntries,
3733        );
3734        assert_grouped_query_exhausts(
3735            &session,
3736            &grouped,
3737            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::CursorSteps,
3738        );
3739
3740        assert_sql_query_exhausts(
3741            &session,
3742            "SELECT payload, COUNT(*) AS row_count FROM IdentityRow \
3743             GROUP BY payload ORDER BY row_count DESC, payload ASC LIMIT 1",
3744            icydb_diagnostic_code::DiagnosticExecutionBudgetResource::SortEntries,
3745        );
3746    }
3747
3748    fn assert_dynamic_payload(session: &DbSession<TestCanister>, key: u64, expected_payload: u64) {
3749        let unchanged = session
3750            .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
3751                entity: ENTITY_NAME.to_string(),
3752                key: InputValue::Nat64(key),
3753                patch: dynamic_payload_patch(expected_payload),
3754            })
3755            .expect("the expected row should remain readable through a no-op update");
3756        assert_eq!(unchanged.affected_rows, 0);
3757        assert_eq!(
3758            unchanged.rows,
3759            vec![expected_dynamic_row(key, expected_payload)],
3760        );
3761    }
3762
3763    fn batch(values: &[u64]) -> Vec<AcceptedStructuralMutation> {
3764        values
3765            .iter()
3766            .map(|value| {
3767                AcceptedStructuralMutation::save(
3768                    MutationMode::Insert,
3769                    AcceptedStructuralMutationTarget::ResolveFromAfterImage,
3770                    payload_patch(*value),
3771                )
3772            })
3773            .collect()
3774    }
3775
3776    fn assert_identity_boundary(error: &InternalError) {
3777        assert_eq!(error.class(), ErrorClass::Unsupported);
3778        assert_eq!(error.origin(), ErrorOrigin::Identity);
3779    }
3780
3781    #[test]
3782    fn generated_candidate_collision_is_identity_corruption_before_generic_uniqueness() {
3783        let generated = insert_key_exists_after_generation(true);
3784        assert_eq!(generated.class(), ErrorClass::Corruption);
3785        assert_eq!(generated.origin(), ErrorOrigin::Identity);
3786
3787        let ordinary = insert_key_exists_after_generation(false);
3788        assert_ne!(ordinary.origin(), ErrorOrigin::Identity);
3789    }
3790
3791    #[cfg(target_pointer_width = "64")]
3792    #[test]
3793    fn pre_key_candidate_count_rejects_values_beyond_the_persisted_u32_bound() {
3794        let error = checked_pre_key_candidate_count(
3795            usize::try_from(u64::from(u32::MAX) + 1).expect("64-bit usize should hold u32 + 1"),
3796        )
3797        .expect_err("candidate counts beyond u32 must reject");
3798        assert_identity_boundary(&error);
3799    }
3800
3801    #[test]
3802    #[expect(
3803        clippy::too_many_lines,
3804        reason = "one holding lifecycle proves split, merge, transfer, late-failure neutrality, result order, and Identity state"
3805    )]
3806    fn mixed_structural_batch_preserves_holding_conservation_and_failure_atomicity() {
3807        let session = initialize();
3808        let seeded = session
3809            .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3810            .expect("seed rows should commit");
3811        assert_eq!(seeded.affected_rows, 1);
3812
3813        let split = session
3814            .execute_trusted_dynamic_mutation_batch(vec![
3815                DynamicMutation::Update {
3816                    entity: ENTITY_NAME.to_string(),
3817                    key: InputValue::Nat64(1),
3818                    patch: dynamic_payload_patch(60),
3819                },
3820                DynamicMutation::Insert {
3821                    entity: ENTITY_NAME.to_string(),
3822                    patch: dynamic_payload_patch(40),
3823                },
3824            ])
3825            .expect("one holding should split atomically");
3826        assert_eq!(split.affected_rows, 2);
3827        assert_eq!(
3828            split.rows,
3829            vec![expected_dynamic_row(1, 60), expected_dynamic_row(2, 40),],
3830            "split after-images must retain input order and exact quantity",
3831        );
3832
3833        let rejected_split = session
3834            .execute_trusted_dynamic_mutation_batch(vec![
3835                DynamicMutation::Update {
3836                    entity: ENTITY_NAME.to_string(),
3837                    key: InputValue::Nat64(1),
3838                    patch: dynamic_payload_patch(50),
3839                },
3840                DynamicMutation::Insert {
3841                    entity: ENTITY_NAME.to_string(),
3842                    patch: DynamicStructuralPatch::new(Vec::new()),
3843                },
3844            ])
3845            .expect_err("an invalid split output must reject the staged source update");
3846        assert_eq!(rejected_split.class(), ErrorClass::Unsupported);
3847        assert_eq!(rejected_split.origin(), ErrorOrigin::Executor);
3848        assert_eq!(
3849            rejected_split.diagnostic_facts(),
3850            vec![
3851                (
3852                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3853                    ENTITY_TAG.value(),
3854                ),
3855                (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 2),
3856                (
3857                    icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
3858                    icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
3859                ),
3860                (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 1,),
3861            ],
3862        );
3863        assert_dynamic_payload(&session, 1, 60);
3864        assert_dynamic_payload(&session, 2, 40);
3865
3866        let transfer = session
3867            .execute_trusted_dynamic_mutation_batch(vec![
3868                DynamicMutation::Update {
3869                    entity: ENTITY_NAME.to_string(),
3870                    key: InputValue::Nat64(1),
3871                    patch: dynamic_payload_patch(70),
3872                },
3873                DynamicMutation::Update {
3874                    entity: ENTITY_NAME.to_string(),
3875                    key: InputValue::Nat64(2),
3876                    patch: dynamic_payload_patch(30),
3877                },
3878            ])
3879            .expect("distinct transfer patches should share one atomic batch");
3880        assert_eq!(
3881            transfer.rows,
3882            vec![expected_dynamic_row(1, 70), expected_dynamic_row(2, 30),],
3883            "the transfer must preserve the exact total quantity",
3884        );
3885
3886        let merge = session
3887            .execute_trusted_dynamic_mutation_batch(vec![
3888                DynamicMutation::Delete {
3889                    entity: ENTITY_NAME.to_string(),
3890                    key: InputValue::Nat64(2),
3891                },
3892                DynamicMutation::Update {
3893                    entity: ENTITY_NAME.to_string(),
3894                    key: InputValue::Nat64(1),
3895                    patch: dynamic_payload_patch(100),
3896                },
3897            ])
3898            .expect("two holdings should merge atomically");
3899        assert_eq!(
3900            merge.rows,
3901            vec![expected_dynamic_row(2, 30), expected_dynamic_row(1, 100),],
3902            "delete before-images and update after-images must retain input order",
3903        );
3904
3905        let resplit = session
3906            .execute_trusted_dynamic_mutation_batch(vec![
3907                DynamicMutation::Update {
3908                    entity: ENTITY_NAME.to_string(),
3909                    key: InputValue::Nat64(1),
3910                    patch: dynamic_payload_patch(60),
3911                },
3912                DynamicMutation::Insert {
3913                    entity: ENTITY_NAME.to_string(),
3914                    patch: dynamic_payload_patch(40),
3915                },
3916            ])
3917            .expect("the merged holding should split again");
3918        assert_eq!(
3919            resplit.rows,
3920            vec![expected_dynamic_row(1, 60), expected_dynamic_row(3, 40),],
3921        );
3922
3923        let rejected_merge = session
3924            .execute_trusted_dynamic_mutation_batch(vec![
3925                DynamicMutation::Delete {
3926                    entity: ENTITY_NAME.to_string(),
3927                    key: InputValue::Nat64(3),
3928                },
3929                DynamicMutation::Update {
3930                    entity: ENTITY_NAME.to_string(),
3931                    key: InputValue::Nat64(99),
3932                    patch: dynamic_payload_patch(100),
3933                },
3934            ])
3935            .expect_err("a late missing merge target must preserve the earlier staged delete");
3936        assert_eq!(rejected_merge.class(), ErrorClass::NotFound);
3937        assert_dynamic_payload(&session, 1, 60);
3938        assert_dynamic_payload(&session, 3, 40);
3939
3940        SCHEMA_STORE.with(|store| {
3941            let cursor = store
3942                .borrow()
3943                .identity_statement_cursor(
3944                    database_incarnation_id().expect("database incarnation should remain readable"),
3945                    ENTITY_TAG,
3946                    FieldId::new(1),
3947                    &AcceptedFieldKind::Nat64,
3948                )
3949                .expect("mixed Identity state should remain readable");
3950            assert_eq!(cursor.expected_high_water(), 3);
3951            assert!(!cursor.has_allocations());
3952        });
3953    }
3954
3955    #[test]
3956    fn mixed_structural_batch_rejects_duplicate_holding_targets_without_mutation() {
3957        let session = initialize();
3958        session
3959            .execute_trusted_dynamic_insert_batch(ENTITY_NAME, vec![dynamic_payload_patch(100)])
3960            .expect("the holding fixture should initialize");
3961
3962        let duplicate = session
3963            .execute_trusted_dynamic_mutation_batch(vec![
3964                DynamicMutation::Update {
3965                    entity: ENTITY_NAME.to_string(),
3966                    key: InputValue::Nat64(1),
3967                    patch: dynamic_payload_patch(60),
3968                },
3969                DynamicMutation::Delete {
3970                    entity: ENTITY_NAME.to_string(),
3971                    key: InputValue::Nat64(1),
3972                },
3973            ])
3974            .expect_err("duplicate targets across operation kinds must reject");
3975        assert!(matches!(
3976            duplicate.diagnostic().detail(),
3977            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
3978                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchDuplicateKey,
3979            }),
3980        ));
3981        assert_eq!(
3982            duplicate.diagnostic_facts(),
3983            vec![
3984                (
3985                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
3986                    ENTITY_TAG.value(),
3987                ),
3988                (
3989                    icydb_diagnostic_code::DiagnosticFactTag::FirstBatchPosition,
3990                    0,
3991                ),
3992                (
3993                    icydb_diagnostic_code::DiagnosticFactTag::DuplicateBatchPosition,
3994                    1,
3995                ),
3996            ],
3997        );
3998        assert_dynamic_payload(&session, 1, 100);
3999    }
4000
4001    #[test]
4002    fn mixed_structural_batch_rejects_empty_and_over_bound_before_resolution() {
4003        let session = initialize();
4004        let empty = session
4005            .execute_trusted_dynamic_mutation_batch(Vec::new())
4006            .expect_err("an empty public batch must reject");
4007        assert!(matches!(
4008            empty.diagnostic().detail(),
4009            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4010                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchEmpty,
4011            }),
4012        ));
4013        assert_eq!(
4014            empty.diagnostic_facts(),
4015            vec![(icydb_diagnostic_code::DiagnosticFactTag::ActualCount, 0,)],
4016        );
4017
4018        let requests = (0..=MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS)
4019            .map(|_| DynamicMutation::Delete {
4020                entity: ENTITY_NAME.to_string(),
4021                key: InputValue::Nat64(1),
4022            })
4023            .collect();
4024        let over_bound = session
4025            .execute_trusted_dynamic_mutation_batch(requests)
4026            .expect_err("operation cap plus one must reject before row resolution");
4027        assert!(matches!(
4028            over_bound.diagnostic().detail(),
4029            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4030                boundary: icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchTooManyItems,
4031            }),
4032        ));
4033        assert_eq!(
4034            over_bound.diagnostic_facts(),
4035            vec![
4036                (
4037                    icydb_diagnostic_code::DiagnosticFactTag::ActualCount,
4038                    (MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS + 1) as u64,
4039                ),
4040                (
4041                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
4042                    MAX_STRUCTURAL_MUTATION_BATCH_OPERATIONS as u64,
4043                ),
4044            ],
4045        );
4046    }
4047
4048    #[test]
4049    fn mixed_structural_batch_staged_byte_bound_uses_checked_exact_boundary() {
4050        let mut exact = 0;
4051        add_structural_mutation_staged_bytes(
4052            &mut exact,
4053            [MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES],
4054        )
4055        .expect("the exact staged-byte boundary should admit");
4056        assert_eq!(exact, MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES);
4057
4058        let error = add_structural_mutation_staged_bytes(&mut exact, [1])
4059            .expect_err("one byte above the staged-byte boundary must reject");
4060        assert!(matches!(
4061            error.diagnostic().detail(),
4062            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4063                boundary:
4064                    icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchStagedBytesExceeded,
4065            }),
4066        ));
4067        assert_eq!(
4068            error.diagnostic_facts(),
4069            vec![
4070                (
4071                    icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4072                    (MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES + 1) as u64,
4073                ),
4074                (
4075                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
4076                    MAX_STRUCTURAL_MUTATION_BATCH_STAGED_BYTES as u64,
4077                ),
4078            ],
4079        );
4080
4081        validate_structural_mutation_result_bytes(MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES)
4082            .expect("the exact result-byte boundary should admit");
4083        let error = validate_structural_mutation_result_bytes(
4084            MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1,
4085        )
4086        .expect_err("one byte above the result-byte boundary must reject");
4087        assert!(matches!(
4088            error.diagnostic().detail(),
4089            Some(icydb_diagnostic_code::DiagnosticDetail::RuntimeBoundary {
4090                boundary:
4091                    icydb_diagnostic_code::RuntimeBoundaryCode::MutationBatchResultBytesExceeded,
4092            }),
4093        ));
4094        assert_eq!(
4095            error.diagnostic_facts(),
4096            vec![
4097                (
4098                    icydb_diagnostic_code::DiagnosticFactTag::ActualLength,
4099                    (MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES + 1) as u64,
4100                ),
4101                (
4102                    icydb_diagnostic_code::DiagnosticFactTag::Limit,
4103                    MAX_STRUCTURAL_MUTATION_BATCH_RESULT_BYTES as u64,
4104                ),
4105            ],
4106        );
4107    }
4108
4109    #[expect(
4110        clippy::too_many_lines,
4111        reason = "one lifecycle proves shared materialization and every maintained frontend against the same zero-state owner"
4112    )]
4113    #[test]
4114    fn identity_insert_frontends_share_one_committed_range_without_rejected_consumption() {
4115        let session = initialize();
4116        let catalog = session
4117            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4118            .expect("identity catalog should resolve");
4119        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4120            .expect("identity row layout should build");
4121        let initial_description = session
4122            .try_describe_entity_by_name(ENTITY_NAME)
4123            .expect("accepted Identity description should resolve");
4124        assert_eq!(
4125            initial_description.entity_tag(),
4126            catalog.identity().entity_tag().value()
4127        );
4128        assert_eq!(
4129            initial_description.accepted_schema_fingerprint_method(),
4130            catalog.fingerprint_method_version()
4131        );
4132        assert_eq!(
4133            initial_description.accepted_schema_fingerprint(),
4134            catalog.fingerprint()
4135        );
4136        let initial_identity = initial_description
4137            .identity()
4138            .expect("accepted Identity policy should be described");
4139        assert_eq!(initial_identity.field(), "id");
4140        assert_eq!(initial_identity.generator(), "Identity::next");
4141        assert_eq!(initial_identity.accepted_kind(), "nat64");
4142        assert_eq!(initial_identity.minimum(), 1);
4143        assert_eq!(initial_identity.maximum(), u128::from(u64::MAX));
4144        assert_eq!(initial_identity.high_water(), 0);
4145        assert_eq!(initial_identity.remaining(), u128::from(u64::MAX));
4146        assert!(!initial_identity.exhausted());
4147
4148        let rejected = session
4149            .execute_accepted_structural_save_batch(
4150                &catalog,
4151                &descriptor,
4152                batch(&[1_000, 2_000]),
4153                Timestamp::from_millis(6),
4154                |_| Err::<(), _>(InternalError::executor_unsupported()),
4155            )
4156            .expect_err("a rejected precommit result must not publish its tentative range");
4157        assert_eq!(rejected.class(), ErrorClass::Unsupported);
4158        assert_eq!(DATA_STORE.with(|store| store.borrow().len()), 0);
4159
4160        let rows = session
4161            .execute_accepted_structural_save_batch(
4162                &catalog,
4163                &descriptor,
4164                batch(&[10, 20, 30]),
4165                Timestamp::from_millis(7),
4166                Ok,
4167            )
4168            .expect("one accepted batch should commit rows and one identity range");
4169        assert_eq!(
4170            rows.into_iter().map(|row| row.values).collect::<Vec<_>>(),
4171            vec![
4172                vec![Value::Nat64(1), Value::Nat64(10)],
4173                vec![Value::Nat64(2), Value::Nat64(20)],
4174                vec![Value::Nat64(3), Value::Nat64(30)],
4175            ],
4176        );
4177
4178        let dynamic = session
4179            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4180                entity: ENTITY_NAME.to_string(),
4181                patch: DynamicStructuralPatch::new(vec![(
4182                    "payload".to_string(),
4183                    DynamicWriteCell::Value(InputValue::Nat64(40)),
4184                )]),
4185            })
4186            .expect("dynamic omission should commit through shared Identity generation");
4187        assert_eq!(dynamic.affected_rows, 1);
4188
4189        for (request, operation) in [
4190            (
4191                DynamicMutation::Insert {
4192                    entity: ENTITY_NAME.to_string(),
4193                    patch: DynamicStructuralPatch::new(vec![
4194                        (
4195                            "id".to_string(),
4196                            DynamicWriteCell::Value(InputValue::Nat64(41)),
4197                        ),
4198                        (
4199                            "payload".to_string(),
4200                            DynamicWriteCell::Value(InputValue::Nat64(42)),
4201                        ),
4202                    ]),
4203                },
4204                icydb_diagnostic_code::DiagnosticMutationOperation::Insert,
4205            ),
4206            (
4207                DynamicMutation::Update {
4208                    entity: ENTITY_NAME.to_string(),
4209                    key: InputValue::Nat64(1),
4210                    patch: DynamicStructuralPatch::new(vec![(
4211                        "id".to_string(),
4212                        DynamicWriteCell::Default,
4213                    )]),
4214                },
4215                icydb_diagnostic_code::DiagnosticMutationOperation::Update,
4216            ),
4217        ] {
4218            let error = session
4219                .execute_trusted_dynamic_mutation(&request)
4220                .expect_err("structural Identity authorship and regeneration must reject");
4221            assert_eq!(error.class(), ErrorClass::Unsupported);
4222            assert_eq!(error.origin(), ErrorOrigin::Executor);
4223            assert_eq!(
4224                error.diagnostic_facts(),
4225                vec![
4226                    (
4227                        icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4228                        ENTITY_TAG.value(),
4229                    ),
4230                    (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
4231                    (
4232                        icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4233                        operation.raw(),
4234                    ),
4235                    (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
4236                ],
4237            );
4238        }
4239
4240        let binding = session
4241            .issue_typed_entity_binding(
4242                ENTITY_SOURCE,
4243                &[
4244                    DynamicTypedFieldBindingRequest::new(
4245                        ID_SOURCE.to_string(),
4246                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4247                        false,
4248                    ),
4249                    DynamicTypedFieldBindingRequest::new(
4250                        PAYLOAD_SOURCE.to_string(),
4251                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4252                        false,
4253                    ),
4254                ],
4255            )
4256            .expect("typed output should bind the Identity field");
4257        let typed_patch = binding
4258            .bind_write_fields(vec![(
4259                PAYLOAD_SOURCE.to_string(),
4260                DynamicWriteCell::Value(InputValue::Nat64(50)),
4261            )])
4262            .expect("typed payload should lower");
4263        let typed = session
4264            .execute_trusted_typed_mutation(
4265                &binding,
4266                &DynamicTypedMutation::Insert { patch: typed_patch },
4267            )
4268            .expect("typed omission should commit through shared Identity generation");
4269        assert_eq!(
4270            typed
4271                .expect("typed insert should return one mutation result")
4272                .affected_rows,
4273            1,
4274        );
4275        let explicit_typed_patch = binding
4276            .bind_write_fields(vec![
4277                (
4278                    ID_SOURCE.to_string(),
4279                    DynamicWriteCell::Value(InputValue::Nat64(51)),
4280                ),
4281                (
4282                    PAYLOAD_SOURCE.to_string(),
4283                    DynamicWriteCell::Value(InputValue::Nat64(52)),
4284                ),
4285            ])
4286            .expect("the low-level binding should retain exact authored intent");
4287        let explicit_typed_error = session
4288            .execute_trusted_typed_mutation(
4289                &binding,
4290                &DynamicTypedMutation::Insert {
4291                    patch: explicit_typed_patch,
4292                },
4293            )
4294            .expect_err("typed Identity authorship must reject before allocation");
4295        assert_eq!(explicit_typed_error.class(), ErrorClass::Unsupported);
4296        assert_eq!(explicit_typed_error.origin(), ErrorOrigin::Executor);
4297        assert_eq!(
4298            explicit_typed_error.diagnostic_facts(),
4299            vec![
4300                (
4301                    icydb_diagnostic_code::DiagnosticFactTag::EntityTag,
4302                    ENTITY_TAG.value(),
4303                ),
4304                (icydb_diagnostic_code::DiagnosticFactTag::FieldId, 1),
4305                (
4306                    icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4307                    icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4308                ),
4309                (icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,),
4310            ],
4311        );
4312
4313        let replace_error = session
4314            .execute_trusted_dynamic_mutation(&DynamicMutation::Replace {
4315                entity: ENTITY_NAME.to_string(),
4316                key: InputValue::Nat64(99),
4317                patch: DynamicStructuralPatch::new(vec![(
4318                    "payload".to_string(),
4319                    DynamicWriteCell::Value(InputValue::Nat64(60)),
4320                )]),
4321            })
4322            .expect_err("save-as-insert with a chosen Identity must reject");
4323        assert_eq!(replace_error.class(), ErrorClass::Unsupported);
4324        assert_eq!(replace_error.origin(), ErrorOrigin::Executor);
4325
4326        #[cfg(feature = "sql")]
4327        {
4328            for sql in [
4329                "INSERT INTO IdentityRow (payload) VALUES (70) RETURNING id, payload",
4330                "INSERT INTO IdentityRow (id, payload) VALUES (DEFAULT, 80) RETURNING id",
4331            ] {
4332                let _result = session
4333                    .execute_trusted_sql_mutation(sql)
4334                    .expect("SQL omission and DEFAULT should commit Identity generation");
4335            }
4336
4337            let error = session
4338                .execute_trusted_sql_mutation(
4339                    "INSERT INTO IdentityRow (id, payload) VALUES (42, 90)",
4340                )
4341                .expect_err("an explicit SQL Identity value must reject before allocation");
4342            let diagnostic = error.diagnostic();
4343            assert_eq!(
4344                diagnostic.code(),
4345                icydb_diagnostic_code::DiagnosticCode::QuerySqlWriteBoundary,
4346            );
4347            assert!(matches!(
4348                diagnostic.detail(),
4349                Some(icydb_diagnostic_code::DiagnosticDetail::SqlWriteBoundary {
4350                    boundary: icydb_diagnostic_code::SqlWriteBoundaryCode::ExplicitGeneratedField,
4351                }),
4352            ));
4353        }
4354
4355        let expected_committed = if cfg!(feature = "sql") { 7 } else { 5 };
4356        assert_eq!(
4357            DATA_STORE.with(|store| store.borrow().len()),
4358            expected_committed
4359        );
4360        SCHEMA_STORE.with(|store| {
4361            let cursor = store
4362                .borrow()
4363                .identity_statement_cursor(
4364                    database_incarnation_id().expect("database incarnation should remain readable"),
4365                    ENTITY_TAG,
4366                    FieldId::new(1),
4367                    &AcceptedFieldKind::Nat64,
4368                )
4369                .expect("committed writes must leave active state readable");
4370            assert_eq!(cursor.expected_high_water(), u128::from(expected_committed),);
4371            assert!(!cursor.has_allocations());
4372        });
4373        let committed_description = session
4374            .try_describe_entity_by_name(ENTITY_NAME)
4375            .expect("committed Identity description should resolve");
4376        let committed_identity = committed_description
4377            .identity()
4378            .expect("accepted Identity policy should remain described");
4379        assert_eq!(
4380            committed_identity.high_water(),
4381            u128::from(expected_committed),
4382        );
4383        assert_eq!(
4384            committed_identity.remaining(),
4385            u128::from(u64::MAX - expected_committed),
4386        );
4387        assert!(!committed_identity.exhausted());
4388    }
4389
4390    #[test]
4391    #[expect(
4392        clippy::too_many_lines,
4393        reason = "one ordered scenario exercises every durable interruption boundary, guarded recovery, derived rebuild, and both integrity tiers"
4394    )]
4395    fn journaled_identity_recovery_quiesces_every_publication_interruption_before_reallocation() {
4396        let session = initialize_journaled();
4397        let catalog = session
4398            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4399            .expect("journaled identity catalog should resolve");
4400        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4401            .expect("journaled identity row layout should build");
4402
4403        for (ordinal, interruption) in [
4404            MutationCommitInterruption::MarkerPersisted,
4405            MutationCommitInterruption::JournalPublished,
4406            MutationCommitInterruption::RowsPublished,
4407            MutationCommitInterruption::StateMaterialized,
4408        ]
4409        .into_iter()
4410        .enumerate()
4411        {
4412            interrupt_next_mutation_commit_for_tests(interruption);
4413            let interrupted = session.execute_accepted_structural_save_batch(
4414                &catalog,
4415                &descriptor,
4416                batch(&[u64::try_from(ordinal).expect("ordinal should fit")]),
4417                Timestamp::from_millis(8),
4418                Ok,
4419            );
4420            assert!(
4421                interrupted.is_err(),
4422                "the selected durable boundary should interrupt",
4423            );
4424
4425            let committed = session
4426                .execute_accepted_structural_save_batch(
4427                    &catalog,
4428                    &descriptor,
4429                    batch(&[100 + u64::try_from(ordinal).expect("ordinal should fit")]),
4430                    Timestamp::from_millis(9),
4431                    Ok,
4432                )
4433                .expect("the next mutation must recover before allocating");
4434            let expected_high_water =
4435                u64::try_from((ordinal + 1) * 2).expect("small test high-water should fit");
4436            assert_eq!(
4437                committed
4438                    .into_iter()
4439                    .map(|row| row.values)
4440                    .collect::<Vec<_>>(),
4441                vec![vec![
4442                    Value::Nat64(expected_high_water),
4443                    Value::Nat64(100 + u64::try_from(ordinal).expect("ordinal should fit")),
4444                ]],
4445            );
4446            assert_eq!(
4447                JOURNALED_DATA_STORE.with(|store| store.borrow().len()),
4448                expected_high_water,
4449            );
4450            JOURNALED_SCHEMA_STORE.with(|store| {
4451                let cursor = store
4452                    .borrow()
4453                    .identity_statement_cursor(
4454                        database_incarnation_id()
4455                            .expect("database incarnation should remain readable"),
4456                        ENTITY_TAG,
4457                        FieldId::new(1),
4458                        &AcceptedFieldKind::Nat64,
4459                    )
4460                    .expect("guarded recovery must leave quiescent active state");
4461                assert_eq!(
4462                    cursor.expected_high_water(),
4463                    u128::from(expected_high_water),
4464                );
4465                assert!(!cursor.has_allocations());
4466            });
4467        }
4468
4469        for (ordinal, (interruption, deleted_key)) in [
4470            (MutationCommitInterruption::MarkerPersisted, 2),
4471            (MutationCommitInterruption::JournalPublished, 4),
4472            (MutationCommitInterruption::RowPrefixPublished, 6),
4473            (MutationCommitInterruption::RowsPublished, 8),
4474            (MutationCommitInterruption::StateMaterialized, 7),
4475        ]
4476        .into_iter()
4477        .enumerate()
4478        {
4479            let expected_payload =
4480                501 + u64::try_from(ordinal).expect("small interruption ordinal should fit");
4481            interrupt_next_mutation_commit_for_tests(interruption);
4482            let interrupted = session.execute_trusted_dynamic_mutation_batch(vec![
4483                DynamicMutation::Update {
4484                    entity: ENTITY_NAME.to_string(),
4485                    key: InputValue::Nat64(1),
4486                    patch: dynamic_payload_patch(expected_payload),
4487                },
4488                DynamicMutation::Delete {
4489                    entity: ENTITY_NAME.to_string(),
4490                    key: InputValue::Nat64(deleted_key),
4491                },
4492            ]);
4493            assert!(
4494                interrupted.is_err(),
4495                "the selected caller-key mixed publication boundary should interrupt",
4496            );
4497            let recovered_update = session
4498                .execute_trusted_dynamic_mutation(&DynamicMutation::Update {
4499                    entity: ENTITY_NAME.to_string(),
4500                    key: InputValue::Nat64(1),
4501                    patch: dynamic_payload_patch(expected_payload),
4502                })
4503                .expect("guarded reentry should complete the marker-authorized mixed batch");
4504            assert_eq!(
4505                recovered_update.affected_rows, 0,
4506                "the recovered update must already expose its admitted final image",
4507            );
4508            let recovered_delete = session
4509                .execute_trusted_dynamic_mutation(&DynamicMutation::Delete {
4510                    entity: ENTITY_NAME.to_string(),
4511                    key: InputValue::Nat64(deleted_key),
4512                })
4513                .expect_err("the recovered delete must already be materialized");
4514            assert_eq!(recovered_delete.class(), ErrorClass::NotFound);
4515            JOURNALED_SCHEMA_STORE.with(|store| {
4516                let cursor = store
4517                    .borrow()
4518                    .identity_statement_cursor(
4519                        database_incarnation_id()
4520                            .expect("database incarnation should remain readable"),
4521                        ENTITY_TAG,
4522                        FieldId::new(1),
4523                        &AcceptedFieldKind::Nat64,
4524                    )
4525                    .expect("caller-key recovery must preserve active Identity state");
4526                assert_eq!(cursor.expected_high_water(), 8);
4527                assert!(!cursor.has_allocations());
4528            });
4529        }
4530
4531        forget_recovered_domain_for_tests(&session.db)
4532            .expect("the final journal tail should remain recoverable");
4533        session
4534            .db
4535            .ensure_recovered_state()
4536            .expect("derived rebuild must not allocate another identity");
4537
4538        let quick = execute_quick_integrity(&session.db, catalog.inspection_plan())
4539            .expect("quiescent Identity control inventory should be inspectable");
4540        assert_eq!(quick.status(), &QuickIntegrityStatus::CompleteClean);
4541        let row_page = execute_row_integrity_page(
4542            &session.db,
4543            catalog.inspection_plan(),
4544            PhysicalUnitCheckpoint::BeforeFirst,
4545            RowInspectionLimits::standard(),
4546        )
4547        .expect("Identity rows should remain within committed high-water");
4548        assert!(row_page.exhausted());
4549        assert!(row_page.findings().is_empty());
4550
4551        assert_eq!(JOURNALED_DATA_STORE.with(|store| store.borrow().len()), 3);
4552        assert!(
4553            JOURNALED_INDEX_STORE.with(|store| !store.borrow().is_empty()),
4554            "derived index rebuild should restore witnesses without allocating identities",
4555        );
4556        assert!(!JOURNALED_TAIL_STORE.with(|tail| tail.borrow().has_stored_batch()));
4557        JOURNALED_SCHEMA_STORE.with(|store| {
4558            let cursor = store
4559                .borrow()
4560                .identity_statement_cursor(
4561                    database_incarnation_id().expect("database incarnation should remain readable"),
4562                    ENTITY_TAG,
4563                    FieldId::new(1),
4564                    &AcceptedFieldKind::Nat64,
4565                )
4566                .expect("folded identity state should reopen without allocating");
4567            assert_eq!(cursor.expected_high_water(), 8);
4568            assert!(!cursor.has_allocations());
4569        });
4570    }
4571
4572    #[test]
4573    #[ignore = "release-closeout native timing probe for one marker-authorized Identity recovery"]
4574    fn identity_recovery_closeout_reports_guarded_reentry_time() {
4575        let session = initialize_journaled();
4576        let catalog = session
4577            .accepted_schema_catalog_context_for_entity_name(Some(ENTITY_NAME))
4578            .expect("journaled identity catalog should resolve");
4579        let descriptor = AcceptedRowLayoutRuntimeContract::from_accepted_schema(catalog.snapshot())
4580            .expect("journaled identity row layout should build");
4581
4582        interrupt_next_mutation_commit_for_tests(MutationCommitInterruption::RowsPublished);
4583        let interrupted = session.execute_accepted_structural_save_batch(
4584            &catalog,
4585            &descriptor,
4586            batch(&[1]),
4587            Timestamp::from_millis(10),
4588            Ok,
4589        );
4590        assert!(
4591            interrupted.is_err(),
4592            "the selected publication boundary should interrupt",
4593        );
4594
4595        let start = Instant::now();
4596        let committed = session
4597            .execute_accepted_structural_save_batch(
4598                &catalog,
4599                &descriptor,
4600                batch(&[2]),
4601                Timestamp::from_millis(11),
4602                Ok,
4603            )
4604            .expect("guarded reentry should recover before allocation");
4605        let elapsed = start.elapsed();
4606        assert_eq!(
4607            committed
4608                .into_iter()
4609                .map(|row| row.values)
4610                .collect::<Vec<_>>(),
4611            vec![vec![Value::Nat64(2), Value::Nat64(2)]],
4612        );
4613
4614        println!(
4615            "identity recovery closeout: guarded_reentry_nanos={}",
4616            elapsed.as_nanos(),
4617        );
4618    }
4619}
4620
4621#[cfg(test)]
4622mod targeted_rule_mutation_tests {
4623    use super::{
4624        DbSession, DynamicMutation, DynamicStructuralPatch, DynamicTypedFieldBindingRequest,
4625        DynamicTypedFieldType, DynamicTypedMutation, DynamicWriteCell,
4626    };
4627    use crate::{
4628        db::{
4629            data::{DataStore, encode_input_value_for_candidate_field_contract},
4630            index::IndexStore,
4631            registry::{StoreAllocationIdentities, StoreRegistry, StoreRuntimeStorageCapabilities},
4632            schema::{
4633                AcceptedCheckLiteralV1, AcceptedCompositeCatalog, AcceptedFieldDecodeContract,
4634                AcceptedFieldKind, AcceptedNamedTypeIdentity, AcceptedRuleOperation,
4635                AcceptedRuleTarget, AcceptedSchemaRevision, AcceptedSourceBindingCatalog,
4636                ConstraintOrigin, FieldId, FieldStorageDecode, FieldWriteManagement, LeafCodec,
4637                PersistedFieldSnapshot, PersistedNestedLeafSnapshot, PersistedSchemaSnapshot,
4638                ScalarCodec, SchemaFieldSlot, SchemaFieldWritePolicy, SchemaInsertDefault,
4639                SchemaRowLayout, SchemaStore, SchemaVersion,
4640                accepted_schema_candidate_with_catalogs_for_tests,
4641                build_record_newtype_composite_catalog_for_tests,
4642                empty_accepted_enum_catalog_for_tests, enum_catalog::ValueAdmissionBudget,
4643            },
4644        },
4645        error::InternalError,
4646        traits::{CanisterKind, Path},
4647        types::EntityTag,
4648        value::InputValue,
4649    };
4650    use icydb_schema::{
4651        ConstraintSourceKey, EntitySourceKey, FieldSourceKey, ScalarType, TypeSourceKey,
4652    };
4653    use std::{cell::RefCell, collections::BTreeMap};
4654
4655    const STORE_PATH: &str = "session::write::targeted_rule_mutation_tests::Store";
4656    const ENTITY_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity";
4657    const ID_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::id";
4658    const PROFILE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Entity::profile";
4659    const UPDATED_AT_SOURCE: &str =
4660        "session::write::targeted_rule_mutation_tests::Entity::updated_at";
4661    const PROFILE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Profile";
4662    const DEGREE_TYPE_SOURCE: &str = "session::write::targeted_rule_mutation_tests::Degree";
4663    const DEGREE_MEMBER_SOURCE: &str =
4664        "session::write::targeted_rule_mutation_tests::Profile::degree";
4665    const DEGREE_RULE_SOURCE: &str =
4666        "session::write::targeted_rule_mutation_tests::Profile::degree_multiple";
4667
4668    struct TestCanister;
4669
4670    impl Path for TestCanister {
4671        const PATH: &'static str = "session::write::targeted_rule_mutation_tests::Canister";
4672    }
4673
4674    impl CanisterKind for TestCanister {
4675        const COMMIT_MEMORY_ID: u8 = 43;
4676        const COMMIT_STABLE_KEY: &'static str = "icydb.targeted_mutation_tests.commit.v1";
4677        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 44;
4678        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
4679            "icydb.targeted_mutation_tests.integrity.progress.v1";
4680    }
4681
4682    thread_local! {
4683        static DATA_STORE: RefCell<DataStore> = const { RefCell::new(DataStore::init_heap()) };
4684        static INDEX_STORE: RefCell<IndexStore> = const { RefCell::new(IndexStore::init_heap()) };
4685        static SCHEMA_STORE: RefCell<SchemaStore> =
4686            const { RefCell::new(SchemaStore::init_heap()) };
4687        static STORE_REGISTRY: StoreRegistry = {
4688            let mut registry = StoreRegistry::new();
4689            registry.register_store(
4690                STORE_PATH,
4691                &DATA_STORE,
4692                &INDEX_STORE,
4693                &SCHEMA_STORE,
4694                StoreAllocationIdentities::absent(),
4695                StoreRuntimeStorageCapabilities::heap(),
4696            ).expect("targeted mutation test store should register");
4697            registry
4698        };
4699    }
4700
4701    fn source<T, E: std::fmt::Debug>(raw: &str, parse: impl FnOnce(String) -> Result<T, E>) -> T {
4702        parse(raw.to_string()).expect("test source identity should admit")
4703    }
4704
4705    fn profile_input(degree: u64) -> InputValue {
4706        InputValue::Map(vec![(
4707            InputValue::Text("degree".to_string()),
4708            InputValue::Nat64(degree),
4709        )])
4710    }
4711
4712    fn structural_patch(id: u64, degree: u64) -> DynamicStructuralPatch {
4713        DynamicStructuralPatch::new(vec![
4714            (
4715                "id".to_string(),
4716                DynamicWriteCell::Value(InputValue::Nat64(id)),
4717            ),
4718            (
4719                "profile".to_string(),
4720                DynamicWriteCell::Value(profile_input(degree)),
4721            ),
4722        ])
4723    }
4724
4725    fn encoded_value(
4726        enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4727        composite_catalog: &AcceptedCompositeCatalog,
4728        name: &str,
4729        kind: &AcceptedFieldKind,
4730        storage_decode: FieldStorageDecode,
4731        leaf_codec: LeafCodec,
4732        value: InputValue,
4733    ) -> Vec<u8> {
4734        let field = AcceptedFieldDecodeContract::new(name, kind, false, storage_decode, leaf_codec);
4735        encode_input_value_for_candidate_field_contract(
4736            enum_catalog,
4737            composite_catalog,
4738            field,
4739            value,
4740            &mut ValueAdmissionBudget::standard(),
4741        )
4742        .expect("test accepted value should encode")
4743    }
4744
4745    fn nat64_literal(
4746        enum_catalog: &crate::db::schema::AcceptedEnumCatalog,
4747        composite_catalog: &AcceptedCompositeCatalog,
4748        value: u64,
4749    ) -> AcceptedCheckLiteralV1 {
4750        let kind = AcceptedFieldKind::Nat64;
4751        AcceptedCheckLiteralV1::from_accepted_parts(
4752            kind.clone(),
4753            FieldStorageDecode::ByKind,
4754            LeafCodec::Scalar(ScalarCodec::Nat64),
4755            encoded_value(
4756                enum_catalog,
4757                composite_catalog,
4758                "degree_bound",
4759                &kind,
4760                FieldStorageDecode::ByKind,
4761                LeafCodec::Scalar(ScalarCodec::Nat64),
4762                InputValue::Nat64(value),
4763            ),
4764        )
4765    }
4766
4767    fn targeted_constraint_id(error: &InternalError) -> u32 {
4768        let facts = error.diagnostic_facts();
4769        assert!(facts.contains(&(
4770            icydb_diagnostic_code::DiagnosticFactTag::MutationOperation,
4771            icydb_diagnostic_code::DiagnosticMutationOperation::Insert.raw(),
4772        )));
4773        assert!(facts.contains(&(icydb_diagnostic_code::DiagnosticFactTag::BatchPosition, 0,)));
4774        assert!(facts.contains(&(
4775            icydb_diagnostic_code::DiagnosticFactTag::ConstraintKind,
4776            icydb_diagnostic_code::DiagnosticConstraintKind::TargetedRule.raw(),
4777        )));
4778        assert_eq!(
4779            facts
4780                .iter()
4781                .filter(|(tag, _)| matches!(
4782                    tag,
4783                    icydb_diagnostic_code::DiagnosticFactTag::RootField
4784                        | icydb_diagnostic_code::DiagnosticFactTag::RecordMember
4785                ))
4786                .copied()
4787                .collect::<Vec<_>>(),
4788            vec![
4789                (icydb_diagnostic_code::DiagnosticFactTag::RootField, 2),
4790                (
4791                    icydb_diagnostic_code::DiagnosticFactTag::RecordMember,
4792                    icydb_diagnostic_code::pack_u32_pair(1, 1),
4793                ),
4794            ]
4795        );
4796        let value = facts
4797            .iter()
4798            .find_map(|(tag, value)| {
4799                (*tag == icydb_diagnostic_code::DiagnosticFactTag::ConstraintId).then_some(*value)
4800            })
4801            .expect("targeted mutation should retain its accepted constraint ID");
4802        u32::try_from(value).expect("accepted constraint ID fits u32")
4803    }
4804
4805    #[expect(
4806        clippy::too_many_lines,
4807        reason = "one end-to-end fixture proves every maintained write frontend converges on the same accepted targeted-rule schedule"
4808    )]
4809    #[test]
4810    fn targeted_rules_converge_across_dynamic_typed_sql_default_timestamp_and_batch_writes() {
4811        DATA_STORE.with(|store| *store.borrow_mut() = DataStore::init_heap());
4812        INDEX_STORE.with(|store| *store.borrow_mut() = IndexStore::init_heap());
4813        SCHEMA_STORE.with(|store| *store.borrow_mut() = SchemaStore::init_heap());
4814
4815        let entity_tag = EntityTag::new(93);
4816        let enum_catalog = empty_accepted_enum_catalog_for_tests();
4817        let (composite_catalog, profile_type, degree_type, degree_member) =
4818            build_record_newtype_composite_catalog_for_tests(
4819                "tests::TargetedProfile".to_string(),
4820                "degree".to_string(),
4821                "tests::TargetedDegree".to_string(),
4822                AcceptedFieldKind::Nat64,
4823                &enum_catalog,
4824            )
4825            .expect("targeted mutation composites should close");
4826        let profile_kind = AcceptedFieldKind::Composite {
4827            type_id: profile_type,
4828        };
4829        let profile_default = encoded_value(
4830            &enum_catalog,
4831            &composite_catalog,
4832            "profile",
4833            &profile_kind,
4834            FieldStorageDecode::CatalogValue,
4835            LeafCodec::Structural,
4836            profile_input(12),
4837        );
4838        let fields = vec![
4839            PersistedFieldSnapshot::new_initial(
4840                FieldId::new(1),
4841                "id".to_string(),
4842                SchemaFieldSlot::new(0),
4843                AcceptedFieldKind::Nat64,
4844                Vec::new(),
4845                false,
4846                SchemaInsertDefault::None,
4847                FieldStorageDecode::ByKind,
4848                LeafCodec::Scalar(ScalarCodec::Nat64),
4849            ),
4850            PersistedFieldSnapshot::new_initial(
4851                FieldId::new(2),
4852                "profile".to_string(),
4853                SchemaFieldSlot::new(1),
4854                profile_kind,
4855                vec![PersistedNestedLeafSnapshot::new(
4856                    vec!["degree".to_string()],
4857                    AcceptedFieldKind::Composite {
4858                        type_id: degree_type,
4859                    },
4860                    false,
4861                )],
4862                false,
4863                SchemaInsertDefault::SlotPayload(profile_default),
4864                FieldStorageDecode::CatalogValue,
4865                LeafCodec::Structural,
4866            ),
4867            PersistedFieldSnapshot::new_initial_with_write_policy(
4868                FieldId::new(3),
4869                "updated_at".to_string(),
4870                SchemaFieldSlot::new(2),
4871                AcceptedFieldKind::Timestamp,
4872                Vec::new(),
4873                false,
4874                SchemaInsertDefault::None,
4875                SchemaFieldWritePolicy::from_model_policies(
4876                    None,
4877                    Some(FieldWriteManagement::UpdatedAt),
4878                ),
4879                FieldStorageDecode::ByKind,
4880                LeafCodec::Scalar(ScalarCodec::Timestamp),
4881            ),
4882        ];
4883        let mut snapshot = PersistedSchemaSnapshot::new(
4884            SchemaVersion::initial(),
4885            ENTITY_SOURCE.to_string(),
4886            "TargetedMutation".to_string(),
4887            FieldId::new(1),
4888            SchemaRowLayout::initial(
4889                fields
4890                    .iter()
4891                    .map(|field| (field.id(), field.slot()))
4892                    .collect(),
4893            ),
4894            fields,
4895        );
4896        let constraint_catalog = snapshot
4897            .constraint_catalog()
4898            .clone()
4899            .with_added_targeted_rule(
4900                "profile_degree_multiple".to_string(),
4901                ConstraintOrigin::Generated,
4902                AcceptedRuleTarget::new(
4903                    FieldId::new(2),
4904                    AcceptedNamedTypeIdentity::Composite(degree_type),
4905                ),
4906                AcceptedRuleOperation::MultipleOf {
4907                    divisor: nat64_literal(&enum_catalog, &composite_catalog, 5),
4908                },
4909            )
4910            .expect("targeted mutation rule should allocate");
4911        let targeted_rule_id = constraint_catalog
4912            .constraints()
4913            .last()
4914            .expect("targeted mutation rule should persist")
4915            .id();
4916        snapshot = snapshot.with_constraint_catalog(constraint_catalog);
4917
4918        let entity_source = source(ENTITY_SOURCE, EntitySourceKey::try_new);
4919        let id_source = source(ID_SOURCE, FieldSourceKey::try_new);
4920        let profile_source = source(PROFILE_SOURCE, FieldSourceKey::try_new);
4921        let updated_at_source = source(UPDATED_AT_SOURCE, FieldSourceKey::try_new);
4922        let profile_type_source = source(PROFILE_TYPE_SOURCE, TypeSourceKey::try_new);
4923        let degree_type_source = source(DEGREE_TYPE_SOURCE, TypeSourceKey::try_new);
4924        let degree_member_source = source(DEGREE_MEMBER_SOURCE, FieldSourceKey::try_new);
4925        let degree_rule_source = source(DEGREE_RULE_SOURCE, ConstraintSourceKey::try_new);
4926        let source_bindings = AcceptedSourceBindingCatalog::initial_for_tests(
4927            BTreeMap::from([(entity_source, entity_tag)]),
4928            BTreeMap::from([
4929                ((entity_tag, id_source), FieldId::new(1)),
4930                ((entity_tag, profile_source), FieldId::new(2)),
4931                ((entity_tag, updated_at_source), FieldId::new(3)),
4932            ]),
4933            BTreeMap::from([((entity_tag, degree_rule_source), targeted_rule_id)]),
4934            BTreeMap::new(),
4935            BTreeMap::new(),
4936        )
4937        .with_initial_named_types_for_tests(
4938            BTreeMap::from([
4939                (
4940                    profile_type_source,
4941                    AcceptedNamedTypeIdentity::Composite(profile_type),
4942                ),
4943                (
4944                    degree_type_source,
4945                    AcceptedNamedTypeIdentity::Composite(degree_type),
4946                ),
4947            ]),
4948            BTreeMap::new(),
4949            BTreeMap::from([((profile_type, degree_member_source), degree_member)]),
4950        );
4951        let candidate = accepted_schema_candidate_with_catalogs_for_tests(
4952            STORE_PATH,
4953            AcceptedSchemaRevision::INITIAL,
4954            enum_catalog,
4955            composite_catalog,
4956            source_bindings,
4957            BTreeMap::from([(entity_tag, snapshot)]),
4958        );
4959
4960        let session = DbSession::<TestCanister>::new(
4961            &STORE_REGISTRY,
4962            &crate::db::RequestExecutionRoot::__new_runtime_root(),
4963        );
4964        session
4965            .db
4966            .ensure_recovered_state()
4967            .expect("targeted mutation test database should initialize");
4968        let store = session
4969            .db
4970            .store_handle(STORE_PATH)
4971            .expect("targeted mutation test store should resolve");
4972        crate::db::commit::publish_accepted_schema_candidate(
4973            STORE_PATH,
4974            store,
4975            AcceptedSchemaRevision::NONE,
4976            &candidate,
4977        )
4978        .expect("targeted mutation candidate should publish");
4979
4980        let dynamic_error = session
4981            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4982                entity: "TargetedMutation".to_string(),
4983                patch: structural_patch(1, 12),
4984            })
4985            .expect_err("dynamic write must enforce the targeted rule");
4986        assert_eq!(
4987            targeted_constraint_id(&dynamic_error),
4988            targeted_rule_id.get()
4989        );
4990
4991        let binding = session
4992            .issue_typed_entity_binding(
4993                ENTITY_SOURCE,
4994                &[
4995                    DynamicTypedFieldBindingRequest::new(
4996                        ID_SOURCE.to_string(),
4997                        DynamicTypedFieldType::Scalar(ScalarType::Nat64),
4998                        false,
4999                    ),
5000                    DynamicTypedFieldBindingRequest::new(
5001                        PROFILE_SOURCE.to_string(),
5002                        DynamicTypedFieldType::Named(PROFILE_TYPE_SOURCE.to_string()),
5003                        false,
5004                    ),
5005                    DynamicTypedFieldBindingRequest::new(
5006                        UPDATED_AT_SOURCE.to_string(),
5007                        DynamicTypedFieldType::Scalar(ScalarType::Timestamp),
5008                        false,
5009                    ),
5010                ],
5011            )
5012            .expect("targeted typed binding should issue");
5013        let typed_patch = binding
5014            .bind_write_fields(vec![
5015                (
5016                    ID_SOURCE.to_string(),
5017                    DynamicWriteCell::Value(InputValue::Nat64(2)),
5018                ),
5019                (
5020                    PROFILE_SOURCE.to_string(),
5021                    DynamicWriteCell::Value(profile_input(12)),
5022                ),
5023            ])
5024            .expect("targeted typed patch should bind");
5025        let typed_error = session
5026            .execute_trusted_typed_mutation(
5027                &binding,
5028                &DynamicTypedMutation::Insert { patch: typed_patch },
5029            )
5030            .expect_err("typed write must enforce the targeted rule");
5031        assert_eq!(targeted_constraint_id(&typed_error), targeted_rule_id.get());
5032
5033        #[cfg(feature = "sql")]
5034        {
5035            let sql_error = session
5036                .execute_trusted_sql_mutation("INSERT INTO TargetedMutation (id) VALUES (3)")
5037                .expect_err("SQL default resolution must enforce the targeted rule");
5038            let crate::db::QueryError::Execute(execute) = sql_error else {
5039                panic!("targeted SQL write should fail at shared execution admission");
5040            };
5041            assert_eq!(
5042                targeted_constraint_id(execute.as_internal()),
5043                targeted_rule_id.get()
5044            );
5045        }
5046
5047        session
5048            .execute_trusted_dynamic_mutation_batch(vec![
5049                DynamicMutation::Insert {
5050                    entity: "TargetedMutation".to_string(),
5051                    patch: structural_patch(4, 5),
5052                },
5053                DynamicMutation::Insert {
5054                    entity: "TargetedMutation".to_string(),
5055                    patch: structural_patch(5, 12),
5056                },
5057            ])
5058            .expect_err("one invalid targeted value must reject the whole batch");
5059        assert_eq!(
5060            DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
5061            Some(0),
5062            "no frontend or earlier valid batch row may escape targeted admission",
5063        );
5064
5065        let admitted = session
5066            .execute_trusted_dynamic_mutation_batch(vec![
5067                DynamicMutation::Insert {
5068                    entity: "TargetedMutation".to_string(),
5069                    patch: structural_patch(6, 5),
5070                },
5071                DynamicMutation::Insert {
5072                    entity: "TargetedMutation".to_string(),
5073                    patch: structural_patch(7, 10),
5074                },
5075            ])
5076            .expect("compliant targeted values should share one accepted batch");
5077        let [first, second] = admitted.rows.as_slice() else {
5078            panic!("the mixed targeted batch should return two rows");
5079        };
5080        let first_timestamp = first
5081            .get(2)
5082            .expect("the first mixed row should contain its managed timestamp");
5083        assert!(matches!(
5084            first_timestamp,
5085            crate::value::OutputValue::Timestamp(_)
5086        ));
5087        assert_eq!(
5088            second.get(2),
5089            Some(first_timestamp),
5090            "one accepted mixed batch must materialize one managed timestamp",
5091        );
5092        assert_eq!(
5093            DATA_STORE.with(|store| store.borrow().exact_entity_count(entity_tag)),
5094            Some(2),
5095        );
5096    }
5097}