Skip to main content

icydb_core/db/schema/
store.rs

1//! Module: db::schema::store
2//! Responsibility: stable BTreeMap-backed schema metadata persistence.
3//! Does not own: reconciliation policy, typed snapshot encoding, or generated proposal construction.
4//! Boundary: provides the third per-store stable memory alongside row and index stores.
5
6use crate::db::schema::identity_state::{
7    IdentityAdvanceId, IdentityRangeAdvance, IdentityRangeCommitState, IdentityState,
8    IdentityStateInventory, IdentityStateLifecycle, IdentityStateTransition,
9    IdentityStatementCursor, MAX_IDENTITY_STATE_RECORDS_PER_DATABASE, decode_identity_state,
10    encode_identity_state, prepare_identity_state_transition, validate_identity_state_closure,
11};
12use crate::db::schema::{
13    cardinality_build::{
14        CardinalityAcceptedDomain, CardinalityReadyCandidate, EmptyCardinalityReadyCandidate,
15    },
16    cardinality_generation::{
17        CardinalityAcceptedRootIdentity, CardinalityBuildCursor, CardinalityCountDigest,
18        CardinalityCountRecord, CardinalityCountSlot, CardinalityGenerationHeader,
19        CardinalityGenerationId, CardinalityGenerationState, CardinalitySourceIdentity,
20    },
21};
22use crate::{
23    db::{
24        codec::{
25            finalize_hash_sha256, new_hash_sha256, write_hash_len_u32, write_hash_str_u32,
26            write_hash_tag_u8, write_hash_u32, write_hash_u64,
27        },
28        commit::CommitSchemaFingerprint,
29        direction::Direction,
30        integrity::DatabaseIncarnationId,
31        journal::{JournalBatch, JournalRecord},
32        ordered_overlay::{OrderedOverlayEntry, ordered_overlay_entries},
33        positioned_overlay::{
34            JournalOverlayPosition, PositionedOverlayMetadata, PositionedOverlayRetirement,
35        },
36        runtime_entity_catalog::AcceptedRuntimeEntity,
37        schema::{
38            AcceptedFieldKind, AcceptedSchemaSnapshot, ConstraintActivationKind,
39            ConstraintActivationState, ConstraintId, ConstraintOrigin, ConstraintValidationJob,
40            FieldId, PersistedIndexKeyItemSnapshot, PersistedIndexKeySnapshot,
41            PersistedSchemaSnapshot, SchemaVersion, accepted_schema_cache_fingerprint,
42            accepted_schema_cache_fingerprint_for_persisted_snapshot,
43            accepted_schema_cache_fingerprint_method_version, decode_constraint_validation_job,
44            decode_persisted_schema_snapshot, encode_constraint_validation_job,
45            encode_persisted_schema_snapshot,
46            enum_catalog::{
47                AcceptedSchemaAuthority, AcceptedSchemaPublicationError, AcceptedSchemaRevision,
48                AcceptedSchemaRevisionBundle, AcceptedSchemaRootSelection,
49                AcceptedStoreCatalogScope, AcceptedValueCatalogHandle, CandidateSchemaRevision,
50                decode_verified_accepted_schema_revision_bundle,
51                prepare_accepted_schema_root_publication, select_current_accepted_schema_root,
52            },
53            schema_snapshot_integrity_detail,
54        },
55    },
56    error::InternalError,
57    types::EntityTag,
58};
59use ic_memory::RuntimeMemory;
60use ic_memory::ic_stable_structures::{
61    BTreeMap as StableBTreeMap, DefaultMemoryImpl, Storable, storable::Bound as StorableBound,
62};
63use sha2::Digest;
64use std::borrow::Cow;
65#[cfg(test)]
66use std::cell::Cell;
67use std::cell::{OnceCell, Ref, RefCell};
68use std::collections::{BTreeMap as StdBTreeMap, BTreeSet};
69#[cfg(test)]
70use std::convert::Infallible;
71use std::ops::Bound as RangeBound;
72use std::rc::Rc;
73
74const SCHEMA_KEY_BYTES_USIZE: usize = 16;
75const SCHEMA_KEY_BYTES: u32 = 16;
76const SCHEMA_KEY_NAMESPACE_ENTITY_SNAPSHOT: u8 = 0;
77const SCHEMA_KEY_NAMESPACE_ACCEPTED_BUNDLE: u8 = 1;
78const SCHEMA_KEY_NAMESPACE_ACCEPTED_ROOT: u8 = 2;
79const SCHEMA_KEY_NAMESPACE_CONSTRAINT_VALIDATION_JOB: u8 = 3;
80const SCHEMA_KEY_NAMESPACE_IDENTITY_STATE: u8 = 4;
81// Every role exposes the sole current method version while its separate domain
82// tag keeps data, index, and full-catalog fingerprint inputs disjoint.
83const SCHEMA_STORE_FINGERPRINT_METHOD_VERSION: u8 = 1;
84const SCHEMA_STORE_CATALOG_FINGERPRINT_DOMAIN: u8 = 1;
85const SCHEMA_STORE_DATA_ALLOCATION_FINGERPRINT_DOMAIN: u8 = 2;
86const SCHEMA_STORE_INDEX_ALLOCATION_FINGERPRINT_DOMAIN: u8 = 3;
87const ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_BOOL: u8 = 3;
88const ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_LIST: u8 = 29;
89const ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_SET: u8 = 30;
90const ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_MAP: u8 = 31;
91const ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_COMPOSITE: u8 = 32;
92const RAW_SCHEMA_SNAPSHOT_MAGIC: &[u8; 8] = b"ICYDBCAT";
93const RAW_SCHEMA_SNAPSHOT_VALUE_VERSION: u8 = 1;
94const RAW_SCHEMA_SNAPSHOT_HEADER_BYTES: usize = 25;
95
96#[cfg(test)]
97thread_local! {
98    static ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES: Cell<u64> = const { Cell::new(0) };
99}
100
101#[cfg(test)]
102fn reset_accepted_schema_bundle_cache_miss_count_for_tests() {
103    ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES.with(|misses| misses.set(0));
104}
105
106#[cfg(test)]
107fn accepted_schema_bundle_cache_miss_count_for_tests() -> u64 {
108    ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES.with(Cell::get)
109}
110
111///
112/// RawSchemaKey
113///
114/// Stable key for one persisted schema snapshot entry.
115/// It combines the entity tag and schema version so reconciliation can load
116/// concrete versions without depending on generated entity names.
117///
118
119#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
120struct RawSchemaKey([u8; SCHEMA_KEY_BYTES_USIZE]);
121
122impl RawSchemaKey {
123    /// Build the raw persisted key for one entity schema version.
124    #[must_use]
125    fn from_entity_version(entity: EntityTag, version: SchemaVersion) -> Self {
126        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
127        out[0] = SCHEMA_KEY_NAMESPACE_ENTITY_SNAPSHOT;
128        out[4..12].copy_from_slice(&entity.value().to_be_bytes());
129        out[12..].copy_from_slice(&version.get().to_be_bytes());
130
131        Self(out)
132    }
133
134    fn from_accepted_bundle(bundle_key: super::enum_catalog::AcceptedSchemaBundleKey) -> Self {
135        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
136        out[0] = SCHEMA_KEY_NAMESPACE_ACCEPTED_BUNDLE;
137        out[4..12].copy_from_slice(&bundle_key.get().to_be_bytes());
138        Self(out)
139    }
140
141    fn from_accepted_root_slot(slot: usize) -> Result<Self, InternalError> {
142        let slot = u32::try_from(slot).map_err(|_| InternalError::store_invariant())?;
143        if slot > 1 {
144            return Err(InternalError::store_invariant());
145        }
146        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
147        out[0] = SCHEMA_KEY_NAMESPACE_ACCEPTED_ROOT;
148        out[12..].copy_from_slice(&slot.to_be_bytes());
149        Ok(Self(out))
150    }
151
152    fn from_constraint_validation_job(entity: EntityTag, constraint_id: ConstraintId) -> Self {
153        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
154        out[0] = SCHEMA_KEY_NAMESPACE_CONSTRAINT_VALIDATION_JOB;
155        out[4..12].copy_from_slice(&entity.value().to_be_bytes());
156        out[12..].copy_from_slice(&constraint_id.get().to_be_bytes());
157        Self(out)
158    }
159
160    fn from_identity_state(entity: EntityTag, field_id: FieldId) -> Self {
161        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
162        out[0] = SCHEMA_KEY_NAMESPACE_IDENTITY_STATE;
163        out[4..12].copy_from_slice(&entity.value().to_be_bytes());
164        out[12..].copy_from_slice(&field_id.get().to_be_bytes());
165        Self(out)
166    }
167
168    /// Return the entity tag encoded in this schema key.
169    #[must_use]
170    fn entity_tag(self) -> EntityTag {
171        let mut bytes = [0u8; size_of::<u64>()];
172        bytes.copy_from_slice(&self.0[4..12]);
173
174        EntityTag::new(u64::from_be_bytes(bytes))
175    }
176
177    /// Return the schema version encoded in this schema key.
178    #[must_use]
179    fn version(self) -> u32 {
180        let mut bytes = [0u8; size_of::<u32>()];
181        bytes.copy_from_slice(&self.0[12..]);
182
183        u32::from_be_bytes(bytes)
184    }
185
186    const fn all_entity_range_bounds() -> (RangeBound<Self>, RangeBound<Self>) {
187        let mut end = [u8::MAX; SCHEMA_KEY_BYTES_USIZE];
188        end[0] = SCHEMA_KEY_NAMESPACE_ENTITY_SNAPSHOT;
189        (
190            RangeBound::Included(Self([0; SCHEMA_KEY_BYTES_USIZE])),
191            RangeBound::Included(Self(end)),
192        )
193    }
194
195    #[cfg(test)]
196    fn entity_range_bounds(entity: EntityTag) -> (RangeBound<Self>, RangeBound<Self>) {
197        (
198            RangeBound::Included(Self::from_entity_version(entity, SchemaVersion::initial())),
199            RangeBound::Included(Self::from_entity_version(
200                entity,
201                SchemaVersion::new(u32::MAX),
202            )),
203        )
204    }
205
206    const fn all_constraint_validation_job_range_bounds() -> (RangeBound<Self>, RangeBound<Self>) {
207        let mut start = [0u8; SCHEMA_KEY_BYTES_USIZE];
208        start[0] = SCHEMA_KEY_NAMESPACE_CONSTRAINT_VALIDATION_JOB;
209        let mut end = [u8::MAX; SCHEMA_KEY_BYTES_USIZE];
210        end[0] = SCHEMA_KEY_NAMESPACE_CONSTRAINT_VALIDATION_JOB;
211        (
212            RangeBound::Included(Self(start)),
213            RangeBound::Included(Self(end)),
214        )
215    }
216
217    const fn all_identity_state_range_bounds() -> (RangeBound<Self>, RangeBound<Self>) {
218        let mut start = [0u8; SCHEMA_KEY_BYTES_USIZE];
219        start[0] = SCHEMA_KEY_NAMESPACE_IDENTITY_STATE;
220        let mut end = [u8::MAX; SCHEMA_KEY_BYTES_USIZE];
221        end[0] = SCHEMA_KEY_NAMESPACE_IDENTITY_STATE;
222        (
223            RangeBound::Included(Self(start)),
224            RangeBound::Included(Self(end)),
225        )
226    }
227
228    #[cfg(test)]
229    const fn is_entity_snapshot(self) -> bool {
230        self.0[0] == SCHEMA_KEY_NAMESPACE_ENTITY_SNAPSHOT
231    }
232
233    const fn is_accepted_root(self) -> bool {
234        self.0[0] == SCHEMA_KEY_NAMESPACE_ACCEPTED_ROOT
235    }
236
237    const fn is_constraint_validation_job(self) -> bool {
238        self.0[0] == SCHEMA_KEY_NAMESPACE_CONSTRAINT_VALIDATION_JOB
239    }
240
241    const fn is_identity_state(self) -> bool {
242        self.0[0] == SCHEMA_KEY_NAMESPACE_IDENTITY_STATE
243    }
244
245    fn constraint_id(self) -> Option<ConstraintId> {
246        self.is_constraint_validation_job()
247            .then(|| ConstraintId::new(self.version()))
248            .flatten()
249    }
250}
251
252impl RawSchemaKey {
253    const NAMESPACE_CARDINALITY_CONTROL: u8 = 5;
254    const NAMESPACE_CARDINALITY_COUNT_A: u8 = 6;
255    const NAMESPACE_CARDINALITY_COUNT_B: u8 = 7;
256    const CARDINALITY_HEADER_DISCRIMINATOR: u8 = 0;
257    const CARDINALITY_BUILD_CURSOR_DISCRIMINATOR: u8 = 1;
258
259    const fn from_cardinality_generation_header() -> Self {
260        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
261        out[0] = Self::NAMESPACE_CARDINALITY_CONTROL;
262        out[SCHEMA_KEY_BYTES_USIZE - 1] = Self::CARDINALITY_HEADER_DISCRIMINATOR;
263        Self(out)
264    }
265
266    const fn from_cardinality_build_cursor() -> Self {
267        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
268        out[0] = Self::NAMESPACE_CARDINALITY_CONTROL;
269        out[SCHEMA_KEY_BYTES_USIZE - 1] = Self::CARDINALITY_BUILD_CURSOR_DISCRIMINATOR;
270        Self(out)
271    }
272
273    fn from_cardinality_count(slot: CardinalityCountSlot, digest: CardinalityCountDigest) -> Self {
274        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
275        out[0] = match slot {
276            CardinalityCountSlot::A => Self::NAMESPACE_CARDINALITY_COUNT_A,
277            CardinalityCountSlot::B => Self::NAMESPACE_CARDINALITY_COUNT_B,
278        };
279        out[1..].copy_from_slice(&digest.as_bytes()[..SCHEMA_KEY_BYTES_USIZE - 1]);
280        Self(out)
281    }
282
283    const fn cardinality_count_range_bounds(
284        slot: CardinalityCountSlot,
285    ) -> (RangeBound<Self>, RangeBound<Self>) {
286        let namespace = match slot {
287            CardinalityCountSlot::A => Self::NAMESPACE_CARDINALITY_COUNT_A,
288            CardinalityCountSlot::B => Self::NAMESPACE_CARDINALITY_COUNT_B,
289        };
290        let mut start = [0_u8; SCHEMA_KEY_BYTES_USIZE];
291        start[0] = namespace;
292        let mut end = [u8::MAX; SCHEMA_KEY_BYTES_USIZE];
293        end[0] = namespace;
294        (
295            RangeBound::Included(Self(start)),
296            RangeBound::Included(Self(end)),
297        )
298    }
299}
300
301impl Storable for RawSchemaKey {
302    fn to_bytes(&self) -> Cow<'_, [u8]> {
303        Cow::Borrowed(&self.0)
304    }
305
306    fn from_bytes(bytes: Cow<'_, [u8]>) -> Self {
307        debug_assert_eq!(
308            bytes.len(),
309            SCHEMA_KEY_BYTES_USIZE,
310            "RawSchemaKey::from_bytes received unexpected byte length",
311        );
312
313        if bytes.len() != SCHEMA_KEY_BYTES_USIZE {
314            return Self([0u8; SCHEMA_KEY_BYTES_USIZE]);
315        }
316
317        let mut out = [0u8; SCHEMA_KEY_BYTES_USIZE];
318        out.copy_from_slice(bytes.as_ref());
319        Self(out)
320    }
321
322    fn into_bytes(self) -> Vec<u8> {
323        self.0.to_vec()
324    }
325
326    const BOUND: StorableBound = StorableBound::Bounded {
327        max_size: SCHEMA_KEY_BYTES,
328        is_fixed_size: true,
329    };
330}
331
332///
333/// RawSchemaSnapshot
334///
335/// Raw persisted value in the schema metadata store.
336///
337/// Entity snapshots carry this wrapper's identity header. Accepted catalog
338/// bundles and root slots are already-versioned control records and remain
339/// opaque here. Key-specific readers decide which representation is required.
340///
341
342#[derive(Clone, Debug, Eq, PartialEq)]
343struct RawSchemaSnapshot {
344    payload: Vec<u8>,
345    accepted_schema_fingerprint: Option<CommitSchemaFingerprint>,
346}
347
348impl RawSchemaSnapshot {
349    /// Encode one typed persisted-schema snapshot into a raw store payload.
350    fn from_persisted_snapshot(snapshot: &PersistedSchemaSnapshot) -> Result<Self, InternalError> {
351        validate_typed_schema_snapshot_for_store(snapshot)?;
352
353        let accepted_schema_fingerprint =
354            accepted_schema_cache_fingerprint_for_persisted_snapshot(snapshot)?;
355        let payload = encode_persisted_schema_snapshot(snapshot)?;
356
357        Ok(Self {
358            payload,
359            accepted_schema_fingerprint: Some(accepted_schema_fingerprint),
360        })
361    }
362
363    /// Store one already-versioned accepted-catalog control record.
364    #[must_use]
365    const fn from_encoded_control_record(payload: Vec<u8>) -> Self {
366        Self {
367            payload,
368            accepted_schema_fingerprint: None,
369        }
370    }
371
372    /// Build a framed entity snapshot around deliberately untrusted payload
373    /// bytes so decode-boundary tests can exercise current-format corruption.
374    #[cfg(test)]
375    #[must_use]
376    const fn from_unchecked_persisted_snapshot_payload(payload: Vec<u8>) -> Self {
377        Self {
378            payload,
379            accepted_schema_fingerprint: Some([0; size_of::<CommitSchemaFingerprint>()]),
380        }
381    }
382
383    /// Borrow the encoded schema snapshot payload.
384    #[must_use]
385    const fn as_bytes(&self) -> &[u8] {
386        self.payload.as_slice()
387    }
388
389    /// Return the accepted schema identity fingerprint stored beside the raw
390    /// payload, without decoding the persisted snapshot.
391    fn accepted_schema_fingerprint(&self) -> Result<CommitSchemaFingerprint, InternalError> {
392        self.accepted_schema_fingerprint
393            .ok_or_else(InternalError::store_corruption)
394    }
395
396    /// Decode this raw store payload into a typed persisted-schema snapshot.
397    fn decode_persisted_snapshot(&self) -> Result<PersistedSchemaSnapshot, InternalError> {
398        // The identity header is the outer format gate. Do not pass a
399        // headerless value or a control record into the schema payload codec.
400        let _fingerprint = self.accepted_schema_fingerprint()?;
401        decode_persisted_schema_snapshot(self.as_bytes())
402    }
403}
404
405#[cfg(test)]
406pub(in crate::db::schema) fn validate_raw_schema_snapshot_bytes_for_tests(
407    bytes: Vec<u8>,
408) -> Result<(), InternalError> {
409    let raw = <RawSchemaSnapshot as Storable>::from_bytes(Cow::Owned(bytes));
410    raw.decode_persisted_snapshot().map(drop)
411}
412
413#[derive(Clone, Debug, Eq, PartialEq)]
414pub(in crate::db) struct AcceptedCatalogIdentity {
415    entity_tag: EntityTag,
416    entity_path: Rc<str>,
417    store_path: &'static str,
418    accepted_schema_revision: AcceptedSchemaRevision,
419    accepted_schema_version: SchemaVersion,
420    fingerprint_method_version: u8,
421    accepted_schema_fingerprint: CommitSchemaFingerprint,
422}
423
424impl AcceptedCatalogIdentity {
425    #[must_use]
426    pub(in crate::db) fn new(
427        entity_tag: EntityTag,
428        entity_path: impl Into<Rc<str>>,
429        store_path: &'static str,
430        accepted_schema_revision: AcceptedSchemaRevision,
431        accepted_schema_version: SchemaVersion,
432        accepted_schema_fingerprint: CommitSchemaFingerprint,
433    ) -> Self {
434        Self {
435            entity_tag,
436            entity_path: entity_path.into(),
437            store_path,
438            accepted_schema_revision,
439            accepted_schema_version,
440            fingerprint_method_version: accepted_schema_cache_fingerprint_method_version(),
441            accepted_schema_fingerprint,
442        }
443    }
444
445    #[must_use]
446    pub(in crate::db) const fn entity_tag(&self) -> EntityTag {
447        self.entity_tag
448    }
449
450    #[must_use]
451    pub(in crate::db) fn entity_path(&self) -> &str {
452        self.entity_path.as_ref()
453    }
454
455    #[must_use]
456    pub(in crate::db) fn entity_path_handle(&self) -> Rc<str> {
457        self.entity_path.clone()
458    }
459
460    #[must_use]
461    pub(in crate::db) const fn store_path(&self) -> &'static str {
462        self.store_path
463    }
464
465    #[must_use]
466    pub(in crate::db) const fn accepted_schema_revision(&self) -> AcceptedSchemaRevision {
467        self.accepted_schema_revision
468    }
469
470    #[must_use]
471    pub(in crate::db) const fn accepted_schema_version(&self) -> SchemaVersion {
472        self.accepted_schema_version
473    }
474
475    #[must_use]
476    pub(in crate::db) const fn fingerprint_method_version(&self) -> u8 {
477        self.fingerprint_method_version
478    }
479
480    #[must_use]
481    pub(in crate::db) const fn accepted_schema_fingerprint(&self) -> CommitSchemaFingerprint {
482        self.accepted_schema_fingerprint
483    }
484}
485
486#[derive(Clone, Debug, Eq, PartialEq)]
487pub(in crate::db) struct AcceptedCatalogSnapshotSelection {
488    identity: AcceptedCatalogIdentity,
489    value_catalog: AcceptedValueCatalogHandle,
490    snapshot: Rc<AcceptedSchemaSnapshot>,
491}
492
493impl AcceptedCatalogSnapshotSelection {
494    // Derive identity and shared authority together from the verified bundle.
495    // Consumers never need to serialize or revalidate this immutable selection.
496    fn from_verified_snapshot(
497        entity: EntityTag,
498        store_path: &'static str,
499        snapshot: &PersistedSchemaSnapshot,
500        value_catalog: AcceptedValueCatalogHandle,
501    ) -> Result<Self, InternalError> {
502        let snapshot = Rc::new(AcceptedSchemaSnapshot::try_new(snapshot.clone())?);
503        let fingerprint = accepted_schema_cache_fingerprint(&snapshot)?;
504        let identity = AcceptedCatalogIdentity::new(
505            entity,
506            snapshot.entity_path(),
507            store_path,
508            value_catalog.revision(),
509            snapshot.persisted_snapshot().version(),
510            fingerprint,
511        );
512
513        Ok(Self {
514            identity,
515            value_catalog,
516            snapshot,
517        })
518    }
519
520    #[must_use]
521    pub(in crate::db) fn identity(&self) -> AcceptedCatalogIdentity {
522        self.identity.clone()
523    }
524
525    #[must_use]
526    pub(in crate::db) const fn value_catalog_handle(&self) -> &AcceptedValueCatalogHandle {
527        &self.value_catalog
528    }
529
530    /// Select one entity snapshot and catalog directly from a verified schema
531    /// candidate during migration preparation.
532    #[cfg(any(test, feature = "migration"))]
533    pub(in crate::db) fn from_candidate(
534        candidate: &CandidateSchemaRevision,
535        entity_tag: EntityTag,
536        entity_path: &str,
537        store_path: &'static str,
538    ) -> Result<Option<Self>, InternalError> {
539        if candidate.store_path() != store_path {
540            return Err(InternalError::store_corruption());
541        }
542        let Some(snapshot) = candidate.bundle().entity_snapshots().get(&entity_tag) else {
543            return Ok(None);
544        };
545        if snapshot.entity_path() != entity_path {
546            return Err(InternalError::store_corruption());
547        }
548
549        Self::from_verified_snapshot(
550            entity_tag,
551            store_path,
552            snapshot,
553            AcceptedValueCatalogHandle::new(
554                candidate.bundle().enum_catalog().clone(),
555                candidate.bundle().composite_catalog().clone(),
556                AcceptedStoreCatalogScope::new(),
557                candidate.revision(),
558                candidate.root().fingerprint(),
559            ),
560        )
561        .map(Some)
562    }
563
564    /// Share the immutable snapshot verified when this selection was built.
565    #[must_use]
566    pub(in crate::db) fn snapshot(&self) -> Rc<AcceptedSchemaSnapshot> {
567        self.snapshot.clone()
568    }
569}
570
571impl Storable for RawSchemaSnapshot {
572    fn to_bytes(&self) -> Cow<'_, [u8]> {
573        let Some(fingerprint) = self.accepted_schema_fingerprint else {
574            return Cow::Borrowed(self.as_bytes());
575        };
576
577        let mut bytes = Vec::with_capacity(RAW_SCHEMA_SNAPSHOT_HEADER_BYTES + self.payload.len());
578        bytes.extend_from_slice(RAW_SCHEMA_SNAPSHOT_MAGIC);
579        bytes.push(RAW_SCHEMA_SNAPSHOT_VALUE_VERSION);
580        bytes.extend_from_slice(&fingerprint);
581        bytes.extend_from_slice(self.as_bytes());
582
583        Cow::Owned(bytes)
584    }
585
586    fn from_bytes(bytes: Cow<'_, [u8]>) -> Self {
587        let bytes = bytes.into_owned();
588        if bytes.len() >= RAW_SCHEMA_SNAPSHOT_HEADER_BYTES
589            && &bytes[..RAW_SCHEMA_SNAPSHOT_MAGIC.len()] == RAW_SCHEMA_SNAPSHOT_MAGIC
590            && bytes[RAW_SCHEMA_SNAPSHOT_MAGIC.len()] == RAW_SCHEMA_SNAPSHOT_VALUE_VERSION
591        {
592            let fingerprint_start = RAW_SCHEMA_SNAPSHOT_MAGIC.len() + size_of::<u8>();
593            let fingerprint_end = fingerprint_start + size_of::<CommitSchemaFingerprint>();
594            let mut fingerprint = [0_u8; size_of::<CommitSchemaFingerprint>()];
595            fingerprint.copy_from_slice(&bytes[fingerprint_start..fingerprint_end]);
596
597            return Self {
598                payload: bytes[fingerprint_end..].to_vec(),
599                accepted_schema_fingerprint: Some(fingerprint),
600            };
601        }
602
603        Self {
604            payload: bytes,
605            accepted_schema_fingerprint: None,
606        }
607    }
608
609    fn into_bytes(self) -> Vec<u8> {
610        let Some(fingerprint) = self.accepted_schema_fingerprint else {
611            return self.payload;
612        };
613
614        let mut bytes = Vec::with_capacity(RAW_SCHEMA_SNAPSHOT_HEADER_BYTES + self.payload.len());
615        bytes.extend_from_slice(RAW_SCHEMA_SNAPSHOT_MAGIC);
616        bytes.push(RAW_SCHEMA_SNAPSHOT_VALUE_VERSION);
617        bytes.extend_from_slice(&fingerprint);
618        bytes.extend_from_slice(&self.payload);
619
620        bytes
621    }
622
623    const BOUND: StorableBound = StorableBound::Unbounded;
624}
625
626// Validate typed schema snapshots before they are encoded into the raw schema
627// metadata store. This catches caller-side invariant violations separately from
628// raw persisted-byte corruption handled by the codec decode boundary.
629fn validate_typed_schema_snapshot_for_store(
630    snapshot: &PersistedSchemaSnapshot,
631) -> Result<(), InternalError> {
632    if schema_snapshot_integrity_detail(
633        "schema snapshot",
634        snapshot.version(),
635        snapshot.primary_key_field_ids(),
636        snapshot.row_layout(),
637        snapshot.fields(),
638    )
639    .is_some()
640    {
641        return Err(InternalError::store_invariant());
642    }
643
644    Ok(())
645}
646
647///
648/// SchemaStoreCatalogMetadata
649///
650/// Accepted schema-store catalog metadata derived from latest persisted
651/// snapshots. This is diagnostic allocation metadata, not allocation identity.
652///
653
654#[derive(Clone, Copy, Debug, Eq, PartialEq)]
655pub(in crate::db) struct SchemaStoreCatalogMetadata {
656    schema_version: SchemaVersion,
657    schema_fingerprint_method_version: u8,
658    schema_fingerprint: CommitSchemaFingerprint,
659    entity_count: u64,
660}
661
662impl SchemaStoreCatalogMetadata {
663    /// Build catalog metadata from already-derived accepted schema facts.
664    #[must_use]
665    const fn new(
666        schema_version: SchemaVersion,
667        schema_fingerprint_method_version: u8,
668        schema_fingerprint: CommitSchemaFingerprint,
669        entity_count: u64,
670    ) -> Self {
671        Self {
672            schema_version,
673            schema_fingerprint_method_version,
674            schema_fingerprint,
675            entity_count,
676        }
677    }
678
679    /// Return the maximum latest schema version represented in the catalog.
680    #[must_use]
681    pub(in crate::db) const fn schema_version(self) -> SchemaVersion {
682        self.schema_version
683    }
684
685    /// Return the fingerprint method version for this diagnostic metadata row.
686    #[must_use]
687    pub(in crate::db) const fn schema_fingerprint_method_version(self) -> u8 {
688        self.schema_fingerprint_method_version
689    }
690
691    /// Return the deterministic catalog fingerprint for latest accepted
692    /// snapshots.
693    #[must_use]
694    pub(in crate::db) const fn schema_fingerprint(self) -> CommitSchemaFingerprint {
695        self.schema_fingerprint
696    }
697
698    /// Return number of entity schemas represented in this catalog metadata.
699    #[must_use]
700    pub(in crate::db) const fn entity_count(self) -> u64 {
701        self.entity_count
702    }
703}
704
705///
706/// SchemaStoreAllocationMetadata
707///
708/// Role-specific allocation metadata derived from latest accepted schema-store
709/// snapshots. These fingerprints describe the accepted contract that owns each
710/// allocation role; they are diagnostics, not allocation identity.
711///
712
713#[derive(Clone, Copy, Debug, Eq, PartialEq)]
714pub(in crate::db) struct SchemaStoreAllocationMetadata {
715    data: SchemaStoreCatalogMetadata,
716    index: SchemaStoreCatalogMetadata,
717    schema: SchemaStoreCatalogMetadata,
718}
719
720impl SchemaStoreAllocationMetadata {
721    /// Build one role-specific metadata set from already-derived accepted
722    /// schema facts.
723    #[must_use]
724    const fn new(
725        data: SchemaStoreCatalogMetadata,
726        index: SchemaStoreCatalogMetadata,
727        schema: SchemaStoreCatalogMetadata,
728    ) -> Self {
729        Self {
730            data,
731            index,
732            schema,
733        }
734    }
735
736    /// Return accepted row-layout allocation metadata for data memory.
737    #[must_use]
738    pub(in crate::db) const fn data(self) -> SchemaStoreCatalogMetadata {
739        self.data
740    }
741
742    /// Return accepted index-catalog allocation metadata for index memory.
743    #[must_use]
744    pub(in crate::db) const fn index(self) -> SchemaStoreCatalogMetadata {
745        self.index
746    }
747
748    /// Return accepted full schema-catalog allocation metadata for schema
749    /// memory.
750    #[must_use]
751    pub(in crate::db) const fn schema(self) -> SchemaStoreCatalogMetadata {
752        self.schema
753    }
754}
755
756///
757/// PendingRelationActivationDeleteBarrier
758///
759/// Accepted activation identity that blocks target deletion until a candidate
760/// reverse-relation generation is proven and promoted.
761///
762
763pub(in crate::db) struct PendingRelationActivationDeleteBarrier {
764    accepted_schema_fingerprint: CommitSchemaFingerprint,
765    source_entity_tag: EntityTag,
766    constraint_id: ConstraintId,
767}
768
769impl PendingRelationActivationDeleteBarrier {
770    #[must_use]
771    pub(in crate::db) const fn accepted_schema_fingerprint(&self) -> CommitSchemaFingerprint {
772        self.accepted_schema_fingerprint
773    }
774
775    #[must_use]
776    pub(in crate::db) const fn source_entity_tag(&self) -> EntityTag {
777        self.source_entity_tag
778    }
779
780    /// Return the stable accepted constraint identity.
781    #[must_use]
782    pub(in crate::db) const fn constraint_id(&self) -> ConstraintId {
783        self.constraint_id
784    }
785}
786
787///
788/// SchemaStore
789///
790/// Thin persistence wrapper over one journaled or heap schema metadata BTreeMap.
791/// Startup reconciliation writes and validates encoded schema snapshots here
792/// before row/index operations proceed.
793///
794
795pub struct SchemaStore {
796    backend: SchemaStoreBackend,
797    accepted_bundle_cache: RefCell<Option<AcceptedSchemaBundleCache>>,
798    cardinality_header_cache: RefCell<Option<(Vec<u8>, CardinalityGenerationHeader)>>,
799    accepted_catalog_scope: OnceCell<AcceptedStoreCatalogScope>,
800}
801
802struct AcceptedSchemaBundleCache {
803    selection: AcceptedSchemaRootSelection,
804    bundle: AcceptedSchemaRevisionBundle,
805    cardinality_domain: Rc<CardinalityAcceptedDomain>,
806    value_catalog: AcceptedValueCatalogHandle,
807    entity_selections: RefCell<StdBTreeMap<EntityTag, AcceptedCatalogSnapshotSelection>>,
808}
809
810enum SchemaStoreBackend {
811    Heap(StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>),
812    Journaled {
813        canonical:
814            StableBTreeMap<RawSchemaKey, RawSchemaSnapshot, RuntimeMemory<DefaultMemoryImpl>>,
815        live: StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
816        tombstones: BTreeSet<RawSchemaKey>,
817        positions: PositionedOverlayMetadata<RawSchemaKey>,
818    },
819}
820
821/// Control-flow result for schema-store traversal visitors.
822#[derive(Clone, Copy, Debug, Eq, PartialEq)]
823enum SchemaStoreVisit {
824    Continue,
825    #[cfg(test)]
826    Stop,
827}
828
829impl SchemaStoreVisit {
830    const fn should_stop(self) -> bool {
831        match self {
832            Self::Continue => false,
833            #[cfg(test)]
834            Self::Stop => true,
835        }
836    }
837}
838
839#[derive(Clone, Copy)]
840enum IdentityStateStorageView {
841    Effective,
842    Canonical,
843}
844
845/// Fully validated entity-snapshot bytes and identity, prepared before apply.
846/// Fields stay private so consumers cannot bypass the store's encoder.
847pub(in crate::db) struct PreparedSchemaSnapshot {
848    key: RawSchemaKey,
849    snapshot: RawSchemaSnapshot,
850}
851
852/// Candidate and canonical payloads retained from preflight until atomic fold.
853/// Not a reusable plan: preparation and application belong to one batch callback.
854pub(in crate::db) struct PreparedAcceptedSchemaFold {
855    candidate: CandidateSchemaRevision,
856    expected_revision: AcceptedSchemaRevision,
857    snapshots: Vec<PreparedSchemaSnapshot>,
858    identity_updates: Vec<(RawSchemaKey, Vec<u8>)>,
859    retained: BTreeSet<RawSchemaKey>,
860    root_slot: usize,
861}
862
863/// Exact schema/control keys whose live values belong to one journal batch.
864pub(in crate::db) struct PreparedSchemaPositionPublication {
865    keys: Vec<RawSchemaKey>,
866    position: JournalOverlayPosition,
867}
868
869/// Preflighted exact schema/control retirement for one complete journal batch.
870pub(in crate::db) struct PreparedSchemaPositionRetirement {
871    entries: Vec<(RawSchemaKey, PositionedOverlayRetirement)>,
872}
873
874/// Preflighted count-record changes for one isolated cardinality build page.
875pub(in crate::db) struct PreparedCardinalityCountWrites {
876    slot: CardinalityCountSlot,
877    generation: CardinalityGenerationId,
878    entries: Vec<(RawSchemaKey, RawSchemaSnapshot)>,
879    new_count_keys: u64,
880}
881
882impl PreparedCardinalityCountWrites {
883    #[must_use]
884    pub(in crate::db) const fn new_count_keys(&self) -> u64 {
885        self.new_count_keys
886    }
887}
888
889/// Fully preflighted atomic count-and-cursor publication for one build page.
890pub(in crate::db) struct PreparedCardinalityBuildPage {
891    count_entries: Vec<(RawSchemaKey, RawSchemaSnapshot)>,
892    cursor: (RawSchemaKey, RawSchemaSnapshot),
893}
894
895/// Fully preflighted exact count and Ready-watermark transition for one fold.
896pub(in crate::db) struct PreparedCardinalityMaintenance {
897    count_entries: Vec<(RawSchemaKey, Option<RawSchemaSnapshot>)>,
898    header: (RawSchemaKey, RawSchemaSnapshot),
899}
900
901#[derive(Clone, Copy)]
902enum IdentityStateWriteTarget {
903    Durable,
904    Materialized,
905    Canonical,
906}
907
908impl SchemaStore {
909    /// Initialize a volatile heap-backed schema store.
910    #[must_use]
911    pub const fn init_heap() -> Self {
912        Self {
913            backend: SchemaStoreBackend::Heap(StdBTreeMap::new()),
914            accepted_bundle_cache: RefCell::new(None),
915            cardinality_header_cache: RefCell::new(None),
916            accepted_catalog_scope: OnceCell::new(),
917        }
918    }
919
920    /// Initialize a journaled cached-stable schema store.
921    ///
922    /// Normal schema publication writes only the live projection. Canonical
923    /// stable schema history is updated by future journal fold/recovery paths.
924    #[must_use]
925    pub fn init_journaled(memory: RuntimeMemory<DefaultMemoryImpl>) -> Self {
926        Self {
927            backend: SchemaStoreBackend::Journaled {
928                canonical: StableBTreeMap::init(memory),
929                live: StdBTreeMap::new(),
930                tombstones: BTreeSet::new(),
931                positions: PositionedOverlayMetadata::new(),
932            },
933            accepted_bundle_cache: RefCell::new(None),
934            cardinality_header_cache: RefCell::new(None),
935            accepted_catalog_scope: OnceCell::new(),
936        }
937    }
938
939    /// Load the sole current durable cardinality-generation header.
940    pub(in crate::db) fn cardinality_generation_header(
941        &self,
942    ) -> Result<Option<CardinalityGenerationHeader>, InternalError> {
943        let key = RawSchemaKey::from_cardinality_generation_header();
944        let raw = self.get_canonical_raw_value(&key)?;
945        if let Some(raw) = raw.as_ref() {
946            self.decode_cardinality_header_cached(raw).map(Some)
947        } else {
948            self.cardinality_header_cache
949                .try_borrow_mut()
950                .map_err(|_| InternalError::store_invariant())?
951                .take();
952            Ok(None)
953        }
954    }
955
956    /// Load the sole bounded cardinality build cursor.
957    pub(in crate::db) fn cardinality_build_cursor(
958        &self,
959    ) -> Result<Option<CardinalityBuildCursor>, InternalError> {
960        let key = RawSchemaKey::from_cardinality_build_cursor();
961        self.get_canonical_raw_value(&key)?
962            .map(|raw| CardinalityBuildCursor::decode(raw.as_bytes()))
963            .transpose()
964    }
965
966    /// Load the generation header and build cursor through one bounded control range.
967    pub(in crate::db) fn cardinality_generation_control(
968        &self,
969    ) -> Result<
970        (
971            Option<CardinalityGenerationHeader>,
972            Option<CardinalityBuildCursor>,
973        ),
974        InternalError,
975    > {
976        let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
977            return Err(InternalError::store_invariant());
978        };
979        let header_key = RawSchemaKey::from_cardinality_generation_header();
980        let cursor_key = RawSchemaKey::from_cardinality_build_cursor();
981        let mut header = None;
982        let mut cursor = None;
983        for entry in canonical.range(header_key..=cursor_key) {
984            if *entry.key() == header_key {
985                header = Some(self.decode_cardinality_header_cached(&entry.value())?);
986            } else if *entry.key() == cursor_key {
987                cursor = Some(CardinalityBuildCursor::decode(entry.value().as_bytes())?);
988            } else {
989                return Err(InternalError::store_corruption());
990            }
991        }
992        Ok((header, cursor))
993    }
994
995    /// Load only lifecycle identity needed by advisory cardinality consumers.
996    ///
997    /// Cursor payload progress is deliberately not decoded or retained here:
998    /// Building progress does not change whether exact evidence is available.
999    pub(in crate::db) fn cardinality_generation_lifecycle_control(
1000        &self,
1001    ) -> Result<(Option<CardinalityGenerationHeader>, bool), InternalError> {
1002        let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
1003            return Err(InternalError::store_invariant());
1004        };
1005        let header_key = RawSchemaKey::from_cardinality_generation_header();
1006        let cursor_key = RawSchemaKey::from_cardinality_build_cursor();
1007        let mut header = None;
1008        let mut cursor_present = false;
1009        for entry in canonical.range(header_key..=cursor_key) {
1010            if *entry.key() == header_key {
1011                header = Some(self.decode_cardinality_header_cached(&entry.value())?);
1012            } else if *entry.key() == cursor_key {
1013                cursor_present = true;
1014            } else {
1015                return Err(InternalError::store_corruption());
1016            }
1017        }
1018
1019        Ok((header, cursor_present))
1020    }
1021
1022    fn decode_cardinality_header_cached(
1023        &self,
1024        raw: &RawSchemaSnapshot,
1025    ) -> Result<CardinalityGenerationHeader, InternalError> {
1026        if let Some((_, header)) = self
1027            .cardinality_header_cache
1028            .try_borrow()
1029            .map_err(|_| InternalError::store_invariant())?
1030            .as_ref()
1031            .filter(|(bytes, _)| bytes.as_slice() == raw.as_bytes())
1032        {
1033            return Ok(*header);
1034        }
1035        let header = CardinalityGenerationHeader::decode(raw.as_bytes())?;
1036        *self
1037            .cardinality_header_cache
1038            .try_borrow_mut()
1039            .map_err(|_| InternalError::store_invariant())? =
1040            Some((raw.as_bytes().to_vec(), header));
1041        Ok(header)
1042    }
1043
1044    /// Prove that no current cardinality authority or orphaned slot data exists.
1045    pub(in crate::db) fn cardinality_storage_is_pristine(&self) -> Result<bool, InternalError> {
1046        if self.cardinality_generation_header()?.is_some()
1047            || self.cardinality_build_cursor()?.is_some()
1048        {
1049            return Ok(false);
1050        }
1051        Ok(
1052            self.cardinality_count_slot_is_empty(CardinalityCountSlot::A)?
1053                && self.cardinality_count_slot_is_empty(CardinalityCountSlot::B)?,
1054        )
1055    }
1056
1057    /// Persist one already-validated cardinality header into canonical control storage.
1058    pub(in crate::db) fn write_cardinality_generation_header(
1059        &mut self,
1060        header: CardinalityGenerationHeader,
1061    ) -> Result<(), InternalError> {
1062        self.insert_canonical_raw_value(
1063            RawSchemaKey::from_cardinality_generation_header(),
1064            header.encode(),
1065        )
1066    }
1067
1068    /// Replace stale cardinality evidence with a fresh isolated Building generation.
1069    ///
1070    /// Every fallible read, generation increment, and backend check happens
1071    /// before the header switch. Removing the obsolete cursor afterward is a
1072    /// mechanical stable-map operation in the same replicated message.
1073    pub(in crate::db) fn restart_cardinality_generation(
1074        &mut self,
1075        current: CardinalityGenerationHeader,
1076        source: CardinalitySourceIdentity,
1077    ) -> Result<CardinalityGenerationHeader, InternalError> {
1078        if self.cardinality_generation_header()? != Some(current) {
1079            return Err(InternalError::store_corruption());
1080        }
1081        if current.validate_source(source).is_ok() {
1082            return Err(InternalError::store_invariant());
1083        }
1084        if let Some(cursor) = self.cardinality_build_cursor()? {
1085            cursor.validate_header(current)?;
1086        }
1087        let next = CardinalityGenerationHeader::new(
1088            current.generation().checked_next()?,
1089            CardinalityGenerationState::Building,
1090            current.slot().alternate(),
1091            source,
1092        );
1093        let encoded = RawSchemaSnapshot::from_encoded_control_record(next.encode());
1094        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1095            return Err(InternalError::store_invariant());
1096        };
1097        canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1098        canonical.remove(&RawSchemaKey::from_cardinality_build_cursor());
1099        Ok(next)
1100    }
1101
1102    /// Publish exact zero for a physically empty canonical row/index domain.
1103    pub(in crate::db) fn publish_empty_cardinality_generation(
1104        &mut self,
1105        candidate: &EmptyCardinalityReadyCandidate,
1106    ) -> Result<CardinalityGenerationHeader, InternalError> {
1107        if !self.cardinality_storage_is_pristine()? {
1108            return Err(InternalError::store_corruption());
1109        }
1110        let ready = CardinalityGenerationHeader::new(
1111            CardinalityGenerationId::INITIAL,
1112            CardinalityGenerationState::Ready,
1113            CardinalityCountSlot::A,
1114            candidate.source(),
1115        );
1116        let encoded = RawSchemaSnapshot::from_encoded_control_record(ready.encode());
1117        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1118            return Err(InternalError::store_invariant());
1119        };
1120        canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1121        Ok(ready)
1122    }
1123
1124    /// Atomically make one completely exhausted candidate planner-visible.
1125    pub(in crate::db) fn publish_ready_cardinality_generation(
1126        &mut self,
1127        candidate: &CardinalityReadyCandidate,
1128        source: CardinalitySourceIdentity,
1129    ) -> Result<CardinalityGenerationHeader, InternalError> {
1130        let building = candidate.header();
1131        candidate.cursor().validate_header(building)?;
1132        if building.state() != CardinalityGenerationState::Building
1133            || building.validate_source(source).is_err()
1134            || self.cardinality_generation_header()? != Some(building)
1135            || self.cardinality_build_cursor()?.as_ref() != Some(candidate.cursor())
1136        {
1137            return Err(InternalError::store_corruption());
1138        }
1139        let ready = CardinalityGenerationHeader::new(
1140            building.generation(),
1141            CardinalityGenerationState::Ready,
1142            building.slot(),
1143            building.source(),
1144        );
1145        let encoded = RawSchemaSnapshot::from_encoded_control_record(ready.encode());
1146        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1147            return Err(InternalError::store_invariant());
1148        };
1149        canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1150        canonical.remove(&RawSchemaKey::from_cardinality_build_cursor());
1151        Ok(ready)
1152    }
1153
1154    /// Clear at most `limit` records from one inactive count slot.
1155    ///
1156    /// When the slot becomes empty, the initial Rows cursor is installed in
1157    /// the same mutation boundary so a resumed builder never confuses clearing
1158    /// with scanning.
1159    pub(in crate::db) fn clear_cardinality_count_slot_page(
1160        &mut self,
1161        header: CardinalityGenerationHeader,
1162        initial_cursor: &CardinalityBuildCursor,
1163        limit: usize,
1164    ) -> Result<bool, InternalError> {
1165        if limit == 0 {
1166            return Err(InternalError::store_invariant());
1167        }
1168        initial_cursor.validate_header(header)?;
1169        let encoded_cursor = initial_cursor.encode()?;
1170        if self.cardinality_generation_header()? != Some(header)
1171            || self.cardinality_build_cursor()?.is_some()
1172        {
1173            return Err(InternalError::store_corruption());
1174        }
1175        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1176            return Err(InternalError::store_invariant());
1177        };
1178        let bounds = RawSchemaKey::cardinality_count_range_bounds(header.slot());
1179        let collect_limit = limit
1180            .checked_add(1)
1181            .ok_or_else(InternalError::store_unsupported)?;
1182        let mut keys = Vec::new();
1183        keys.try_reserve_exact(collect_limit)
1184            .map_err(|_| InternalError::store_unsupported())?;
1185        for entry in canonical.range(bounds).take(collect_limit) {
1186            keys.push(*entry.key());
1187        }
1188        let has_more = keys.len() > limit;
1189        for key in keys.into_iter().take(limit) {
1190            canonical.remove(&key);
1191        }
1192        if !has_more {
1193            canonical.insert(
1194                RawSchemaKey::from_cardinality_build_cursor(),
1195                RawSchemaSnapshot::from_encoded_control_record(encoded_cursor),
1196            );
1197        }
1198        Ok(has_more)
1199    }
1200
1201    /// Preflight coalesced positive count increments for one isolated page.
1202    pub(in crate::db) fn prepare_cardinality_count_increments(
1203        &self,
1204        slot: CardinalityCountSlot,
1205        generation: CardinalityGenerationId,
1206        increments: &[(CardinalityCountDigest, u64)],
1207    ) -> Result<PreparedCardinalityCountWrites, InternalError> {
1208        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1209            return Err(InternalError::store_invariant());
1210        }
1211        let mut entries = Vec::new();
1212        entries
1213            .try_reserve_exact(increments.len())
1214            .map_err(|_| InternalError::store_unsupported())?;
1215        let mut physical_keys = StdBTreeMap::new();
1216        let mut new_count_keys = 0_u64;
1217        for (digest, increment) in increments {
1218            if *increment == 0 {
1219                return Err(InternalError::store_invariant());
1220            }
1221            let key = RawSchemaKey::from_cardinality_count(slot, *digest);
1222            if let Some(previous_digest) = physical_keys.insert(key, *digest) {
1223                return Err(if previous_digest == *digest {
1224                    InternalError::store_invariant()
1225                } else {
1226                    InternalError::store_corruption()
1227                });
1228            }
1229            let current = self.get_canonical_raw_value(&key)?;
1230            let count = if let Some(raw) = current {
1231                let record = CardinalityCountRecord::decode(raw.as_bytes())?;
1232                record
1233                    .validate_identity(generation, *digest)
1234                    .map_err(|_| InternalError::store_corruption())?
1235                    .checked_add(*increment)
1236                    .ok_or_else(InternalError::store_unsupported)?
1237            } else {
1238                new_count_keys = new_count_keys
1239                    .checked_add(1)
1240                    .ok_or_else(InternalError::store_unsupported)?;
1241                *increment
1242            };
1243            let record = CardinalityCountRecord::new(generation, *digest, count)?;
1244            entries.push((
1245                key,
1246                RawSchemaSnapshot::from_encoded_control_record(record.encode().to_vec()),
1247            ));
1248        }
1249        Ok(PreparedCardinalityCountWrites {
1250            slot,
1251            generation,
1252            entries,
1253            new_count_keys,
1254        })
1255    }
1256
1257    /// Bind preflighted count writes to the exact current cursor transition.
1258    pub(in crate::db) fn prepare_cardinality_build_page(
1259        &self,
1260        header: CardinalityGenerationHeader,
1261        current_cursor: &CardinalityBuildCursor,
1262        counts: PreparedCardinalityCountWrites,
1263        next_cursor: &CardinalityBuildCursor,
1264    ) -> Result<PreparedCardinalityBuildPage, InternalError> {
1265        current_cursor.validate_header(header)?;
1266        next_cursor.validate_header(header)?;
1267        if counts.slot != header.slot() || counts.generation != header.generation() {
1268            return Err(InternalError::store_invariant());
1269        }
1270        if self.cardinality_generation_header()? != Some(header)
1271            || self.cardinality_build_cursor()?.as_ref() != Some(current_cursor)
1272        {
1273            return Err(InternalError::store_corruption());
1274        }
1275        let encoded_cursor = next_cursor.encode()?;
1276        Ok(PreparedCardinalityBuildPage {
1277            count_entries: counts.entries,
1278            cursor: (
1279                RawSchemaKey::from_cardinality_build_cursor(),
1280                RawSchemaSnapshot::from_encoded_control_record(encoded_cursor),
1281            ),
1282        })
1283    }
1284
1285    /// Mechanically publish one fully preflighted count-and-cursor page.
1286    pub(in crate::db) fn apply_prepared_cardinality_build_page(
1287        &mut self,
1288        prepared: PreparedCardinalityBuildPage,
1289    ) -> Result<(), InternalError> {
1290        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1291            return Err(InternalError::store_invariant());
1292        };
1293        for (key, value) in prepared.count_entries {
1294            canonical.insert(key, value);
1295        }
1296        canonical.insert(prepared.cursor.0, prepared.cursor.1);
1297        Ok(())
1298    }
1299
1300    /// Load one exact nonzero count from an isolated generation.
1301    pub(in crate::db) fn cardinality_count(
1302        &self,
1303        slot: CardinalityCountSlot,
1304        generation: CardinalityGenerationId,
1305        digest: CardinalityCountDigest,
1306    ) -> Result<Option<u64>, InternalError> {
1307        let key = RawSchemaKey::from_cardinality_count(slot, digest);
1308        self.get_canonical_raw_value(&key)?
1309            .map(|raw| {
1310                CardinalityCountRecord::decode(raw.as_bytes())?
1311                    .validate_identity(generation, digest)
1312                    .map_err(|_| InternalError::store_corruption())
1313            })
1314            .transpose()
1315    }
1316
1317    /// Preflight exact coalesced count changes and the next complete fold watermark.
1318    pub(in crate::db) fn prepare_cardinality_maintenance(
1319        &self,
1320        current: CardinalityGenerationHeader,
1321        current_source: CardinalitySourceIdentity,
1322        next_source: CardinalitySourceIdentity,
1323        changes: &[(CardinalityCountDigest, i64)],
1324    ) -> Result<PreparedCardinalityMaintenance, InternalError> {
1325        if current.state() != CardinalityGenerationState::Ready
1326            || current.validate_source(current_source).is_err()
1327            || self.cardinality_generation_header()? != Some(current)
1328            || self.cardinality_build_cursor()?.is_some()
1329        {
1330            return Err(InternalError::store_corruption());
1331        }
1332        let next = CardinalityGenerationHeader::new(
1333            current.generation(),
1334            CardinalityGenerationState::Ready,
1335            current.slot(),
1336            next_source,
1337        );
1338        let mut count_entries = Vec::new();
1339        count_entries
1340            .try_reserve_exact(changes.len())
1341            .map_err(|_| InternalError::store_unsupported())?;
1342        let mut physical_keys = StdBTreeMap::new();
1343        for (digest, delta) in changes {
1344            if *delta == 0 {
1345                return Err(InternalError::store_invariant());
1346            }
1347            let key = RawSchemaKey::from_cardinality_count(current.slot(), *digest);
1348            if let Some(previous_digest) = physical_keys.insert(key, *digest) {
1349                return Err(if previous_digest == *digest {
1350                    InternalError::store_invariant()
1351                } else {
1352                    InternalError::store_corruption()
1353                });
1354            }
1355            let base = self
1356                .cardinality_count(current.slot(), current.generation(), *digest)?
1357                .unwrap_or(0);
1358            let count = if *delta > 0 {
1359                base.checked_add(
1360                    u64::try_from(*delta).map_err(|_| InternalError::store_invariant())?,
1361                )
1362            } else {
1363                base.checked_sub(delta.unsigned_abs())
1364            }
1365            .ok_or_else(InternalError::store_corruption)?;
1366            let value = if count == 0 {
1367                None
1368            } else {
1369                Some(RawSchemaSnapshot::from_encoded_control_record(
1370                    CardinalityCountRecord::new(current.generation(), *digest, count)?
1371                        .encode()
1372                        .to_vec(),
1373                ))
1374            };
1375            count_entries.push((key, value));
1376        }
1377        Ok(PreparedCardinalityMaintenance {
1378            count_entries,
1379            header: (
1380                RawSchemaKey::from_cardinality_generation_header(),
1381                RawSchemaSnapshot::from_encoded_control_record(next.encode()),
1382            ),
1383        })
1384    }
1385
1386    /// Mechanically publish one completely preflighted count/watermark transition.
1387    pub(in crate::db) fn apply_prepared_cardinality_maintenance(
1388        &mut self,
1389        prepared: PreparedCardinalityMaintenance,
1390    ) -> Result<(), InternalError> {
1391        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1392            return Err(InternalError::store_invariant());
1393        };
1394        for (key, value) in prepared.count_entries {
1395            if let Some(value) = value {
1396                canonical.insert(key, value);
1397            } else {
1398                canonical.remove(&key);
1399            }
1400        }
1401        canonical.insert(prepared.header.0, prepared.header.1);
1402        Ok(())
1403    }
1404
1405    pub(in crate::db) fn cardinality_count_slot_is_empty(
1406        &self,
1407        slot: CardinalityCountSlot,
1408    ) -> Result<bool, InternalError> {
1409        let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
1410            return Err(InternalError::store_invariant());
1411        };
1412        Ok(canonical
1413            .range(RawSchemaKey::cardinality_count_range_bounds(slot))
1414            .next()
1415            .is_none())
1416    }
1417
1418    /// Prove that recovered journal metadata can fold into canonical storage.
1419    pub(in crate::db) fn preflight_fold_recovered_journal(&self) -> Result<(), InternalError> {
1420        match self.backend {
1421            SchemaStoreBackend::Journaled { .. } => Ok(()),
1422            SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
1423        }
1424    }
1425
1426    fn prepare_identity_state_transition(
1427        &self,
1428        incarnation: DatabaseIncarnationId,
1429        candidate: &CandidateSchemaRevision,
1430        view: IdentityStateStorageView,
1431    ) -> Result<IdentityStateTransition, InternalError> {
1432        let current = match view {
1433            IdentityStateStorageView::Effective => self
1434                .current_accepted_schema_bundle_ref()?
1435                .as_ref()
1436                .map(|bundle| (*bundle).clone()),
1437            IdentityStateStorageView::Canonical => {
1438                self.current_canonical_accepted_schema_bundle()?
1439            }
1440        };
1441        let inventory = self.identity_state_inventory(view)?;
1442        prepare_identity_state_transition(
1443            incarnation,
1444            current.as_ref(),
1445            candidate.bundle(),
1446            inventory,
1447        )
1448    }
1449
1450    fn validate_identity_state_closure(
1451        &self,
1452        bundle: &AcceptedSchemaRevisionBundle,
1453    ) -> Result<(), InternalError> {
1454        let inventory = self.identity_state_inventory(IdentityStateStorageView::Effective)?;
1455        validate_identity_state_closure(bundle, &inventory)
1456    }
1457
1458    /// Read one accepted active Identity owner into statement-local allocation state.
1459    pub(in crate::db) fn identity_statement_cursor(
1460        &self,
1461        database_incarnation_id: DatabaseIncarnationId,
1462        entity_tag: EntityTag,
1463        field_id: FieldId,
1464        accepted_kind: &AcceptedFieldKind,
1465    ) -> Result<IdentityStatementCursor, InternalError> {
1466        let key = RawSchemaKey::from_identity_state(entity_tag, field_id);
1467        let raw = self
1468            .get_raw_snapshot(&key)
1469            .ok_or_else(InternalError::identity_state_corruption)?;
1470        let state = decode_identity_state(raw.as_bytes())?;
1471        let owner = state.owner();
1472        if owner.database_incarnation_id() != database_incarnation_id
1473            || owner.entity_tag() != entity_tag
1474            || owner.field_id() != field_id
1475            || state.accepted_kind() != accepted_kind
1476            || state.lifecycle() != IdentityStateLifecycle::Active
1477        {
1478            return Err(InternalError::identity_state_corruption());
1479        }
1480        IdentityStatementCursor::from_active_state(&state)
1481    }
1482
1483    /// Read one quiescent materialized high-water for bounded row integrity.
1484    pub(in crate::db) fn identity_high_water_for_integrity(
1485        &self,
1486        database_incarnation_id: DatabaseIncarnationId,
1487        entity_tag: EntityTag,
1488        field_id: FieldId,
1489        accepted_kind: &AcceptedFieldKind,
1490    ) -> Result<u128, InternalError> {
1491        let key = RawSchemaKey::from_identity_state(entity_tag, field_id);
1492        let raw = self
1493            .get_raw_snapshot(&key)
1494            .ok_or_else(InternalError::identity_state_corruption)?;
1495        let state = decode_identity_state(raw.as_bytes())?;
1496        let owner = state.owner();
1497        if owner.database_incarnation_id() != database_incarnation_id
1498            || owner.entity_tag() != entity_tag
1499            || owner.field_id() != field_id
1500            || state.accepted_kind() != accepted_kind
1501            || state.lifecycle() != IdentityStateLifecycle::Active
1502        {
1503            return Err(InternalError::identity_state_corruption());
1504        }
1505        Ok(state.materialized_high_water())
1506    }
1507
1508    /// Revalidate one tentative range against the quiescent effective state.
1509    pub(in crate::db) fn preflight_identity_range_advance(
1510        &self,
1511        range: IdentityRangeAdvance,
1512    ) -> Result<(), InternalError> {
1513        let state =
1514            self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Effective)?;
1515        state.preflight_range_advance(range)
1516    }
1517
1518    /// Materialize one marker-owned range in the effective live projection.
1519    pub(in crate::db) fn apply_identity_range_advance(
1520        &mut self,
1521        range: IdentityRangeAdvance,
1522        advance_id: IdentityAdvanceId,
1523    ) -> Result<(), InternalError> {
1524        self.apply_identity_range_advance_to(
1525            range,
1526            advance_id,
1527            IdentityStateStorageView::Effective,
1528            IdentityStateWriteTarget::Materialized,
1529        )
1530    }
1531
1532    /// Fold one marker-owned range into canonical journaled state.
1533    pub(in crate::db) fn fold_identity_range_advance(
1534        &mut self,
1535        range: IdentityRangeAdvance,
1536        advance_id: IdentityAdvanceId,
1537    ) -> Result<(), InternalError> {
1538        self.apply_identity_range_advance_to(
1539            range,
1540            advance_id,
1541            IdentityStateStorageView::Canonical,
1542            IdentityStateWriteTarget::Canonical,
1543        )
1544    }
1545
1546    /// Preflight one canonical Identity range fold without changing storage.
1547    pub(in crate::db) fn preflight_fold_identity_range_advance(
1548        &self,
1549        range: IdentityRangeAdvance,
1550        advance_id: IdentityAdvanceId,
1551    ) -> Result<(), InternalError> {
1552        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1553            return Err(InternalError::store_invariant());
1554        }
1555        let state =
1556            self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Canonical)?;
1557        let advanced = state.apply_range_advance(range, advance_id)?;
1558        let _encoded = encode_identity_state(&advanced)?;
1559        Ok(())
1560    }
1561
1562    /// Verify one exact range identity against effective state.
1563    pub(in crate::db) fn verify_identity_range_advance(
1564        &self,
1565        range: IdentityRangeAdvance,
1566        advance_id: IdentityAdvanceId,
1567    ) -> Result<(), InternalError> {
1568        let state =
1569            self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Effective)?;
1570        if state.materialized_high_water() != range.new_high_water()
1571            || state.last_applied_advance() != Some(advance_id)
1572        {
1573            return Err(InternalError::recovery_effect_verification_failed());
1574        }
1575        Ok(())
1576    }
1577
1578    /// Resolve committed versus materialized range state without changing it.
1579    pub(in crate::db) fn identity_range_commit_state(
1580        &self,
1581        range: IdentityRangeAdvance,
1582        advance_id: IdentityAdvanceId,
1583        canonical: bool,
1584    ) -> Result<IdentityRangeCommitState, InternalError> {
1585        let view = if canonical {
1586            IdentityStateStorageView::Canonical
1587        } else {
1588            IdentityStateStorageView::Effective
1589        };
1590        self.identity_state_for_owner(range.owner(), view)?
1591            .range_commit_state(range, advance_id)
1592    }
1593
1594    /// Enumerate and validate the complete current-form active/retired state
1595    /// inventory for bounded database-wide integrity inspection.
1596    pub(in crate::db) fn identity_state_inventory_for_integrity(
1597        &self,
1598        incarnation: DatabaseIncarnationId,
1599    ) -> Result<Vec<IdentityState>, InternalError> {
1600        let has_accepted_bundle = self.current_accepted_schema_bundle_ref()?.is_some();
1601        let inventory = self.identity_state_inventory(IdentityStateStorageView::Effective)?;
1602        if !has_accepted_bundle && !inventory.is_empty() {
1603            return Err(InternalError::identity_state_corruption());
1604        }
1605        if inventory
1606            .values()
1607            .any(|state| state.owner().database_incarnation_id() != incarnation)
1608        {
1609            return Err(InternalError::identity_state_corruption());
1610        }
1611        Ok(inventory.into_values().collect())
1612    }
1613
1614    fn identity_state_for_owner(
1615        &self,
1616        owner: crate::db::schema::identity_state::IdentityStateOwner,
1617        view: IdentityStateStorageView,
1618    ) -> Result<IdentityState, InternalError> {
1619        let key = RawSchemaKey::from_identity_state(owner.entity_tag(), owner.field_id());
1620        let raw = match view {
1621            IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
1622            IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
1623        }
1624        .ok_or_else(InternalError::identity_state_corruption)?;
1625        let state = decode_identity_state(raw.as_bytes())?;
1626        if state.owner() != owner {
1627            return Err(InternalError::identity_state_corruption());
1628        }
1629        Ok(state)
1630    }
1631
1632    fn apply_identity_range_advance_to(
1633        &mut self,
1634        range: IdentityRangeAdvance,
1635        advance_id: IdentityAdvanceId,
1636        view: IdentityStateStorageView,
1637        target: IdentityStateWriteTarget,
1638    ) -> Result<(), InternalError> {
1639        let state = self.identity_state_for_owner(range.owner(), view)?;
1640        let advanced = state.apply_range_advance(range, advance_id)?;
1641        let key = RawSchemaKey::from_identity_state(
1642            advanced.owner().entity_tag(),
1643            advanced.owner().field_id(),
1644        );
1645        let bytes = encode_identity_state(&advanced)?;
1646        match target {
1647            IdentityStateWriteTarget::Materialized => {
1648                self.insert_raw_snapshot(
1649                    key,
1650                    RawSchemaSnapshot::from_encoded_control_record(bytes),
1651                );
1652            }
1653            IdentityStateWriteTarget::Canonical => {
1654                self.insert_canonical_raw_value(key, bytes)?;
1655            }
1656            IdentityStateWriteTarget::Durable => {
1657                return Err(InternalError::store_invariant());
1658            }
1659        }
1660        Ok(())
1661    }
1662
1663    fn identity_state_inventory(
1664        &self,
1665        view: IdentityStateStorageView,
1666    ) -> Result<IdentityStateInventory, InternalError> {
1667        let bounds = RawSchemaKey::all_identity_state_range_bounds();
1668        let mut inventory = StdBTreeMap::new();
1669        let mut collect = |key: &RawSchemaKey,
1670                           raw: &RawSchemaSnapshot|
1671         -> Result<SchemaStoreVisit, InternalError> {
1672            if inventory.len() >= MAX_IDENTITY_STATE_RECORDS_PER_DATABASE {
1673                return Err(InternalError::identity_state_corruption());
1674            }
1675            let state = decode_identity_state(raw.as_bytes())?;
1676            let state_key = (key.entity_tag(), FieldId::new(key.version()));
1677            if !key.is_identity_state()
1678                || state.owner().entity_tag() != state_key.0
1679                || state.owner().field_id() != state_key.1
1680                || inventory.insert(state_key, state).is_some()
1681            {
1682                return Err(InternalError::identity_state_corruption());
1683            }
1684            Ok(SchemaStoreVisit::Continue)
1685        };
1686
1687        match (&self.backend, view) {
1688            (SchemaStoreBackend::Heap(map), IdentityStateStorageView::Effective) => {
1689                for (key, raw) in map.range((bounds.0, bounds.1)) {
1690                    collect(key, raw)?;
1691                }
1692            }
1693            (
1694                SchemaStoreBackend::Journaled {
1695                    canonical,
1696                    live,
1697                    tombstones,
1698                    ..
1699                },
1700                IdentityStateStorageView::Effective,
1701            ) => Self::visit_journaled_raw_snapshot_range(
1702                canonical,
1703                live,
1704                tombstones,
1705                bounds,
1706                Direction::Asc,
1707                &mut collect,
1708            )?,
1709            (
1710                SchemaStoreBackend::Journaled { canonical, .. },
1711                IdentityStateStorageView::Canonical,
1712            ) => {
1713                for entry in canonical.range((bounds.0, bounds.1)) {
1714                    collect(entry.key(), &entry.value())?;
1715                }
1716            }
1717            (SchemaStoreBackend::Heap(_), IdentityStateStorageView::Canonical) => {
1718                return Err(InternalError::store_invariant());
1719            }
1720        }
1721
1722        Ok(inventory)
1723    }
1724
1725    fn apply_identity_state_transition(
1726        &mut self,
1727        transition: IdentityStateTransition,
1728        target: IdentityStateWriteTarget,
1729    ) -> Result<(), InternalError> {
1730        for state in transition.into_updates() {
1731            let key = RawSchemaKey::from_identity_state(
1732                state.owner().entity_tag(),
1733                state.owner().field_id(),
1734            );
1735            let bytes = encode_identity_state(&state)?;
1736            match target {
1737                IdentityStateWriteTarget::Durable => {
1738                    self.insert_durable_raw_value(key, bytes);
1739                }
1740                IdentityStateWriteTarget::Materialized => {
1741                    self.insert_raw_snapshot(
1742                        key,
1743                        RawSchemaSnapshot::from_encoded_control_record(bytes),
1744                    );
1745                }
1746                IdentityStateWriteTarget::Canonical => {
1747                    self.insert_canonical_raw_value(key, bytes)?;
1748                }
1749            }
1750        }
1751        Ok(())
1752    }
1753
1754    pub(in crate::db) fn current_canonical_accepted_schema_bundle(
1755        &self,
1756    ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
1757        self.current_canonical_accepted_schema_authority()
1758            .map(|authority| authority.map(|(_, bundle)| bundle))
1759    }
1760
1761    /// Load one canonical accepted root and its verified immutable bundle.
1762    pub(in crate::db) fn current_canonical_accepted_schema_authority(
1763        &self,
1764    ) -> Result<Option<(AcceptedSchemaRootSelection, AcceptedSchemaRevisionBundle)>, InternalError>
1765    {
1766        let Some(selection) = self.current_canonical_accepted_schema_root()? else {
1767            return Ok(None);
1768        };
1769        let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
1770        let raw = self
1771            .get_canonical_raw_value(&bundle_key)?
1772            .ok_or_else(InternalError::store_corruption)?;
1773        let bundle =
1774            decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
1775        Ok(Some((selection, bundle)))
1776    }
1777
1778    /// Return the accepted root selected only from canonical predecessor slots.
1779    pub(in crate::db) fn current_canonical_accepted_schema_root(
1780        &self,
1781    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
1782        let first = self.canonical_root_slot_bytes(0)?;
1783        let second = self.canonical_root_slot_bytes(1)?;
1784        select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
1785    }
1786
1787    /// Select effective and canonical roots from one canonical slot read.
1788    pub(in crate::db) fn current_effective_and_canonical_accepted_schema_roots(
1789        &self,
1790    ) -> Result<
1791        (
1792            Option<AcceptedSchemaRootSelection>,
1793            Option<AcceptedSchemaRootSelection>,
1794        ),
1795        InternalError,
1796    > {
1797        let SchemaStoreBackend::Journaled {
1798            canonical,
1799            live,
1800            tombstones,
1801            ..
1802        } = &self.backend
1803        else {
1804            return Err(InternalError::store_invariant());
1805        };
1806        let first_key = RawSchemaKey::from_accepted_root_slot(0)?;
1807        let second_key = RawSchemaKey::from_accepted_root_slot(1)?;
1808        let mut canonical_first = None;
1809        let mut canonical_second = None;
1810        for entry in canonical.range(first_key..=second_key) {
1811            if *entry.key() == first_key {
1812                canonical_first = Some(entry.value().clone());
1813            } else if *entry.key() == second_key {
1814                canonical_second = Some(entry.value().clone());
1815            } else {
1816                return Err(InternalError::store_corruption());
1817            }
1818        }
1819        let effective_first = if tombstones.contains(&first_key) {
1820            None
1821        } else {
1822            live.get(&first_key)
1823                .cloned()
1824                .or_else(|| canonical_first.clone())
1825        };
1826        let effective_second = if tombstones.contains(&second_key) {
1827            None
1828        } else {
1829            live.get(&second_key)
1830                .cloned()
1831                .or_else(|| canonical_second.clone())
1832        };
1833        let effective_first = effective_first.map(RawSchemaSnapshot::into_bytes);
1834        let effective_second = effective_second.map(RawSchemaSnapshot::into_bytes);
1835        let canonical_first = canonical_first.map(RawSchemaSnapshot::into_bytes);
1836        let canonical_second = canonical_second.map(RawSchemaSnapshot::into_bytes);
1837        Ok((
1838            select_current_accepted_schema_root([
1839                effective_first.as_deref(),
1840                effective_second.as_deref(),
1841            ])?,
1842            select_current_accepted_schema_root([
1843                canonical_first.as_deref(),
1844                canonical_second.as_deref(),
1845            ])?,
1846        ))
1847    }
1848
1849    /// Insert or replace one typed persisted schema snapshot.
1850    pub(in crate::db) fn insert_persisted_snapshot(
1851        &mut self,
1852        entity: EntityTag,
1853        snapshot: &PersistedSchemaSnapshot,
1854    ) -> Result<(), InternalError> {
1855        let prepared = Self::prepare_persisted_snapshot(entity, snapshot)?;
1856        self.apply_prepared_persisted_snapshot(prepared);
1857
1858        Ok(())
1859    }
1860
1861    /// Finish snapshot validation, fingerprinting and encoding before mutation.
1862    pub(in crate::db) fn prepare_persisted_snapshot(
1863        entity: EntityTag,
1864        snapshot: &PersistedSchemaSnapshot,
1865    ) -> Result<PreparedSchemaSnapshot, InternalError> {
1866        Ok(PreparedSchemaSnapshot {
1867            key: RawSchemaKey::from_entity_version(entity, snapshot.version()),
1868            snapshot: RawSchemaSnapshot::from_persisted_snapshot(snapshot)?,
1869        })
1870    }
1871
1872    /// Publish prepared bytes without repeating semantic construction.
1873    pub(in crate::db) fn apply_prepared_persisted_snapshot(
1874        &mut self,
1875        prepared: PreparedSchemaSnapshot,
1876    ) {
1877        let _ = self.insert_raw_snapshot(prepared.key, prepared.snapshot);
1878    }
1879
1880    /// Load one schema-owned constraint validation job.
1881    pub(in crate::db) fn constraint_validation_job(
1882        &self,
1883        entity: EntityTag,
1884        constraint_id: ConstraintId,
1885    ) -> Result<Option<ConstraintValidationJob>, InternalError> {
1886        let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1887        self.get_raw_snapshot(&key)
1888            .map(|raw| decode_constraint_validation_job(raw.as_bytes()))
1889            .transpose()
1890    }
1891
1892    /// Apply one marker-authorized validation job to the live schema projection.
1893    pub(in crate::db) fn apply_constraint_validation_job(
1894        &mut self,
1895        job: &ConstraintValidationJob,
1896    ) -> Result<(), InternalError> {
1897        let key =
1898            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1899        let bytes = encode_constraint_validation_job(job)?;
1900        let _ =
1901            self.insert_raw_snapshot(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1902        Ok(())
1903    }
1904
1905    /// Remove one marker-authorized validation job from the live projection.
1906    #[expect(
1907        clippy::unnecessary_wraps,
1908        reason = "marker apply operations share one fallible callback contract"
1909    )]
1910    pub(in crate::db) fn apply_constraint_validation_job_removal(
1911        &mut self,
1912        entity: EntityTag,
1913        constraint_id: ConstraintId,
1914    ) -> Result<(), InternalError> {
1915        let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1916        match &mut self.backend {
1917            SchemaStoreBackend::Heap(map) => {
1918                map.remove(&key);
1919            }
1920            SchemaStoreBackend::Journaled {
1921                live, tombstones, ..
1922            } => {
1923                live.remove(&key);
1924                tombstones.insert(key);
1925            }
1926        }
1927        Ok(())
1928    }
1929
1930    /// Fold one committed validation job into the canonical stable base.
1931    pub(in crate::db) fn fold_constraint_validation_job(
1932        &mut self,
1933        job: &ConstraintValidationJob,
1934    ) -> Result<(), InternalError> {
1935        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1936            return Err(InternalError::store_invariant());
1937        };
1938        let key =
1939            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1940        let bytes = encode_constraint_validation_job(job)?;
1941        canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1942        Ok(())
1943    }
1944
1945    /// Preflight one canonical validation-job fold without changing storage.
1946    pub(in crate::db) fn preflight_fold_constraint_validation_job(
1947        &self,
1948        job: &ConstraintValidationJob,
1949    ) -> Result<(), InternalError> {
1950        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1951            return Err(InternalError::store_invariant());
1952        }
1953        let _encoded = encode_constraint_validation_job(job)?;
1954        Ok(())
1955    }
1956
1957    /// Fold one committed validation-job removal into the canonical stable base.
1958    pub(in crate::db) fn fold_constraint_validation_job_removal(
1959        &mut self,
1960        entity: EntityTag,
1961        constraint_id: ConstraintId,
1962    ) -> Result<(), InternalError> {
1963        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1964            return Err(InternalError::store_invariant());
1965        };
1966        canonical.remove(&RawSchemaKey::from_constraint_validation_job(
1967            entity,
1968            constraint_id,
1969        ));
1970        Ok(())
1971    }
1972
1973    /// Preflight one canonical validation-job removal without changing storage.
1974    pub(in crate::db) fn preflight_fold_constraint_validation_job_removal(
1975        &self,
1976    ) -> Result<(), InternalError> {
1977        match self.backend {
1978            SchemaStoreBackend::Journaled { .. } => Ok(()),
1979            SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
1980        }
1981    }
1982
1983    /// Reset the volatile projection for journaled recovery without mutating
1984    /// the canonical stable schema base.
1985    pub(in crate::db) fn reset_journaled_live_projection(&mut self) -> Result<(), InternalError> {
1986        let SchemaStoreBackend::Journaled {
1987            live,
1988            tombstones,
1989            positions,
1990            ..
1991        } = &mut self.backend
1992        else {
1993            return Err(InternalError::store_invariant());
1994        };
1995
1996        live.clear();
1997        tombstones.clear();
1998        positions.clear();
1999        self.accepted_bundle_cache.get_mut().take();
2000
2001        Ok(())
2002    }
2003
2004    /// Preflight every schema/control position represented by one online batch.
2005    pub(in crate::db) fn prepare_positioned_journal_batch_publication(
2006        &self,
2007        incarnation: DatabaseIncarnationId,
2008        batch: &JournalBatch,
2009        position: JournalOverlayPosition,
2010    ) -> Result<PreparedSchemaPositionPublication, InternalError> {
2011        let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2012            return Err(InternalError::store_invariant());
2013        };
2014        let keys = self.positioned_journal_batch_keys(
2015            incarnation,
2016            batch,
2017            IdentityStateStorageView::Effective,
2018        )?;
2019        for key in &keys {
2020            positions.preflight_publish(key, position)?;
2021        }
2022        Ok(PreparedSchemaPositionPublication {
2023            keys: keys.into_iter().collect(),
2024            position,
2025        })
2026    }
2027
2028    /// Preflight exact schema/control retirement before canonical mutation.
2029    pub(in crate::db) fn prepare_positioned_journal_batch_retirement(
2030        &self,
2031        incarnation: DatabaseIncarnationId,
2032        batch: &JournalBatch,
2033        position: JournalOverlayPosition,
2034    ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2035        let keys = self.positioned_journal_batch_keys(
2036            incarnation,
2037            batch,
2038            IdentityStateStorageView::Canonical,
2039        )?;
2040        self.prepare_positioned_key_retirements(keys, position)
2041    }
2042
2043    fn prepare_positioned_key_retirements(
2044        &self,
2045        keys: impl IntoIterator<Item = RawSchemaKey>,
2046        position: JournalOverlayPosition,
2047    ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2048        let SchemaStoreBackend::Journaled {
2049            live,
2050            tombstones,
2051            positions,
2052            ..
2053        } = &self.backend
2054        else {
2055            return Err(InternalError::store_invariant());
2056        };
2057        let mut entries = Vec::new();
2058        for key in keys {
2059            if !positions.is_positioned(&key) {
2060                // A prior row fold can create canonical-only derived metadata
2061                // after this older schema batch publishes. Its canonical fold
2062                // owns that key; there is no overlay for this batch to retire.
2063                if live.contains_key(&key) || tombstones.contains(&key) {
2064                    return Err(InternalError::store_invariant());
2065                }
2066                continue;
2067            }
2068            let retirement = positions.preflight_retirement(&key, position)?;
2069            entries.push((key, retirement));
2070        }
2071        Ok(PreparedSchemaPositionRetirement { entries })
2072    }
2073
2074    /// Publish schema positions after their values have been mechanically applied.
2075    pub(in crate::db) fn publish_prepared_journal_batch_positions(
2076        &mut self,
2077        prepared: PreparedSchemaPositionPublication,
2078    ) {
2079        let SchemaStoreBackend::Journaled { positions, .. } = &mut self.backend else {
2080            debug_assert!(
2081                false,
2082                "preflighted schema positions require a journaled store"
2083            );
2084            return;
2085        };
2086        for key in prepared.keys {
2087            positions.publish_preflighted(key, prepared.position);
2088        }
2089    }
2090
2091    /// Retire only exact schema/control overlays after canonical mutation.
2092    pub(in crate::db) fn apply_prepared_journal_batch_retirement(
2093        &mut self,
2094        prepared: PreparedSchemaPositionRetirement,
2095    ) {
2096        for (key, retirement) in prepared.entries {
2097            if retirement != PositionedOverlayRetirement::Exact {
2098                continue;
2099            }
2100            self.invalidate_accepted_bundle_cache_for_key(key);
2101            let SchemaStoreBackend::Journaled {
2102                live,
2103                tombstones,
2104                positions,
2105                ..
2106            } = &mut self.backend
2107            else {
2108                debug_assert!(
2109                    false,
2110                    "preflighted schema retirement requires a journaled store"
2111                );
2112                return;
2113            };
2114            live.remove(&key);
2115            tombstones.remove(&key);
2116            positions.retire_preflighted(&key, retirement);
2117        }
2118    }
2119
2120    #[cfg(test)]
2121    fn publish_positioned_journal_entry(
2122        &mut self,
2123        key: RawSchemaKey,
2124        snapshot: Option<RawSchemaSnapshot>,
2125        position: JournalOverlayPosition,
2126    ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
2127        let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2128            return Err(InternalError::store_invariant());
2129        };
2130        positions.preflight_publish(&key, position)?;
2131        self.invalidate_accepted_bundle_cache_for_key(key);
2132        let SchemaStoreBackend::Journaled {
2133            canonical,
2134            live,
2135            tombstones,
2136            positions,
2137        } = &mut self.backend
2138        else {
2139            return Err(InternalError::store_invariant());
2140        };
2141        let previous = if tombstones.contains(&key) {
2142            None
2143        } else {
2144            live.get(&key).cloned().or_else(|| canonical.get(&key))
2145        };
2146        if let Some(snapshot) = snapshot {
2147            tombstones.remove(&key);
2148            live.insert(key, snapshot);
2149        } else {
2150            live.remove(&key);
2151            tombstones.insert(key);
2152        }
2153        positions.publish_preflighted(key, position);
2154        Ok(previous)
2155    }
2156
2157    #[cfg(test)]
2158    fn retire_positioned_journal_effect(
2159        &mut self,
2160        key: RawSchemaKey,
2161        position: JournalOverlayPosition,
2162    ) -> Result<PositionedOverlayRetirement, InternalError> {
2163        let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2164            return Err(InternalError::store_invariant());
2165        };
2166        let retirement = positions.preflight_retirement(&key, position)?;
2167        let prepared = PreparedSchemaPositionRetirement {
2168            entries: vec![(key, retirement)],
2169        };
2170        self.apply_prepared_journal_batch_retirement(prepared);
2171        Ok(retirement)
2172    }
2173
2174    /// Seed a test's canonical snapshot through the maintained prepared handoff.
2175    #[cfg(test)]
2176    pub(in crate::db) fn fold_persisted_snapshot(
2177        &mut self,
2178        entity: EntityTag,
2179        snapshot: &PersistedSchemaSnapshot,
2180    ) -> Result<(), InternalError> {
2181        let prepared = self.prepare_fold_persisted_snapshot(entity, snapshot)?;
2182        self.apply_prepared_fold_persisted_snapshot(prepared)
2183    }
2184
2185    /// Apply only the encoded snapshot prepared for this journaled store.
2186    pub(in crate::db) fn apply_prepared_fold_persisted_snapshot(
2187        &mut self,
2188        prepared: PreparedSchemaSnapshot,
2189    ) -> Result<(), InternalError> {
2190        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
2191            return Err(InternalError::store_invariant());
2192        };
2193        canonical.insert(prepared.key, prepared.snapshot);
2194
2195        Ok(())
2196    }
2197
2198    /// Prepare one canonical fold, retaining encoded bytes without changing storage.
2199    pub(in crate::db) fn prepare_fold_persisted_snapshot(
2200        &self,
2201        entity: EntityTag,
2202        snapshot: &PersistedSchemaSnapshot,
2203    ) -> Result<PreparedSchemaSnapshot, InternalError> {
2204        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2205            return Err(InternalError::store_invariant());
2206        }
2207        Self::prepare_persisted_snapshot(entity, snapshot)
2208    }
2209
2210    /// Return the current accepted store root selected from its two checksummed slots.
2211    pub(in crate::db) fn current_accepted_schema_root(
2212        &self,
2213    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
2214        let first = self.accepted_root_slot_bytes(0)?;
2215        let second = self.accepted_root_slot_bytes(1)?;
2216        select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
2217    }
2218
2219    /// Load and verify the immutable bundle referenced by the current root.
2220    pub(in crate::db) fn current_accepted_schema_bundle(
2221        &self,
2222    ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
2223        self.borrow_current_accepted_schema_bundle()
2224            .map(|bundle| bundle.map(|bundle| bundle.clone()))
2225    }
2226
2227    /// Borrow the current verified bundle, rechecking durable job closure even
2228    /// on a cache hit. The borrow must end before mutating schema authority.
2229    pub(in crate::db) fn borrow_current_accepted_schema_bundle(
2230        &self,
2231    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
2232        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2233            return Ok(None);
2234        };
2235        self.validate_constraint_validation_job_closure(&bundle)?;
2236
2237        Ok(Some(bundle))
2238    }
2239
2240    /// Borrow the verified bundle selected by a current catalog operation.
2241    /// Root publication invalidates this store-owned cache before replacement;
2242    /// exact authority equality also binds store scope, revision and fingerprint.
2243    /// Mutable identity and validation-job records still require live checks.
2244    pub(in crate::db) fn borrow_accepted_schema_bundle_for_authority(
2245        &self,
2246        expected: &AcceptedSchemaAuthority,
2247    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
2248        let Some(selection) = self.current_accepted_selection_for_authority(expected)? else {
2249            return Ok(None);
2250        };
2251        let Some((_, bundle)) =
2252            self.accepted_schema_authority_ref_for_selection(Some(selection))?
2253        else {
2254            return Ok(None);
2255        };
2256        self.validate_constraint_validation_job_closure(&bundle)?;
2257
2258        Ok(Some(bundle))
2259    }
2260
2261    /// Project current accepted entity identity onto one registry-owned store path.
2262    pub(in crate::db) fn current_accepted_runtime_entities(
2263        &self,
2264        registered_store_path: &'static str,
2265    ) -> Result<Vec<AcceptedRuntimeEntity>, InternalError> {
2266        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2267            return Ok(Vec::new());
2268        };
2269        if bundle.store_path() != registered_store_path {
2270            return Err(InternalError::store_corruption());
2271        }
2272
2273        bundle
2274            .entity_snapshots()
2275            .iter()
2276            .map(|(entity_tag, snapshot)| {
2277                AcceptedRuntimeEntity::from_accepted_snapshot(
2278                    &bundle,
2279                    *entity_tag,
2280                    snapshot,
2281                    registered_store_path,
2282                )
2283            })
2284            .collect()
2285    }
2286
2287    /// Resolve one accepted entity tag without materializing the full store catalog.
2288    pub(in crate::db) fn current_accepted_runtime_entity_for_tag(
2289        &self,
2290        registered_store_path: &'static str,
2291        entity_tag: EntityTag,
2292    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2293        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2294            return Ok(None);
2295        };
2296        if bundle.store_path() != registered_store_path {
2297            return Err(InternalError::store_corruption());
2298        }
2299        let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2300            return Ok(None);
2301        };
2302
2303        AcceptedRuntimeEntity::from_accepted_snapshot(
2304            &bundle,
2305            entity_tag,
2306            snapshot,
2307            registered_store_path,
2308        )
2309        .map(Some)
2310    }
2311
2312    /// Resolve one entity tag from the canonical accepted predecessor.
2313    pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_tag(
2314        &self,
2315        registered_store_path: &'static str,
2316        entity_tag: EntityTag,
2317    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2318        let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2319            return Ok(None);
2320        };
2321        if bundle.store_path() != registered_store_path {
2322            return Err(InternalError::store_corruption());
2323        }
2324        let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2325            return Ok(None);
2326        };
2327
2328        AcceptedRuntimeEntity::from_accepted_snapshot(
2329            &bundle,
2330            entity_tag,
2331            snapshot,
2332            registered_store_path,
2333        )
2334        .map(Some)
2335    }
2336
2337    /// Resolve one accepted entity source path without materializing the full store catalog.
2338    pub(in crate::db) fn current_accepted_runtime_entity_for_path(
2339        &self,
2340        registered_store_path: &'static str,
2341        entity_path: &str,
2342    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2343        self.current_accepted_runtime_entity_matching(registered_store_path, |snapshot_path, _| {
2344            snapshot_path == entity_path
2345        })
2346    }
2347
2348    /// Resolve one entity path from the canonical accepted predecessor.
2349    pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_path(
2350        &self,
2351        registered_store_path: &'static str,
2352        entity_path: &str,
2353    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2354        let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2355            return Ok(None);
2356        };
2357        if bundle.store_path() != registered_store_path {
2358            return Err(InternalError::store_corruption());
2359        }
2360
2361        let mut matched = None;
2362        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2363            if snapshot.entity_path() != entity_path {
2364                continue;
2365            }
2366            let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2367                &bundle,
2368                *entity_tag,
2369                snapshot,
2370                registered_store_path,
2371            )?;
2372            if matched.replace(entity).is_some() {
2373                return Err(InternalError::store_corruption());
2374            }
2375        }
2376
2377        Ok(matched)
2378    }
2379
2380    /// Resolve one accepted entity display name without materializing the full store catalog.
2381    #[cfg(test)]
2382    pub(in crate::db) fn current_accepted_runtime_entity_for_name(
2383        &self,
2384        registered_store_path: &'static str,
2385        entity_name: &str,
2386    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2387        self.current_accepted_runtime_entity_matching(registered_store_path, |_, snapshot_name| {
2388            snapshot_name == entity_name
2389        })
2390    }
2391
2392    fn current_accepted_runtime_entity_matching(
2393        &self,
2394        registered_store_path: &'static str,
2395        mut predicate: impl FnMut(&str, &str) -> bool,
2396    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2397        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2398            return Ok(None);
2399        };
2400        if bundle.store_path() != registered_store_path {
2401            return Err(InternalError::store_corruption());
2402        }
2403
2404        let mut matched = None;
2405        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2406            if !predicate(snapshot.entity_path(), snapshot.entity_name()) {
2407                continue;
2408            }
2409            let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2410                &bundle,
2411                *entity_tag,
2412                snapshot,
2413                registered_store_path,
2414            )?;
2415            if matched.replace(entity).is_some() {
2416                return Err(InternalError::store_corruption());
2417            }
2418        }
2419
2420        Ok(matched)
2421    }
2422
2423    /// Return the current accepted revision without decoding its bundle.
2424    pub(in crate::db) fn current_accepted_schema_revision(
2425        &self,
2426    ) -> Result<Option<AcceptedSchemaRevision>, InternalError> {
2427        Ok(self
2428            .current_accepted_schema_root()?
2429            .map(|selection| selection.root().revision()))
2430    }
2431
2432    /// Return the pending relation activation that blocks deletes from one target.
2433    ///
2434    /// This reads the immutable accepted-bundle cache directly so ordinary
2435    /// deletes do not decode and clone every store catalog merely to prove that
2436    /// no candidate reverse generation targets the deleted entity.
2437    pub(in crate::db) fn pending_relation_activation_for_target(
2438        &self,
2439        target_path: &str,
2440    ) -> Result<Option<PendingRelationActivationDeleteBarrier>, InternalError> {
2441        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2442            return Ok(None);
2443        };
2444        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2445            let Some(candidate) = snapshot
2446                .candidate_relations()
2447                .iter()
2448                .find(|candidate| candidate.target_path() == target_path)
2449            else {
2450                continue;
2451            };
2452            let activation = snapshot
2453                .constraint_activations()
2454                .iter()
2455                .find(|activation| {
2456                    matches!(
2457                        activation.kind(),
2458                        ConstraintActivationKind::Relation { relation_id }
2459                            if *relation_id == candidate.id()
2460                    )
2461                })
2462                .ok_or_else(InternalError::store_corruption)?;
2463            return Ok(Some(PendingRelationActivationDeleteBarrier {
2464                accepted_schema_fingerprint:
2465                    accepted_schema_cache_fingerprint_for_persisted_snapshot(snapshot)?,
2466                source_entity_tag: *entity_tag,
2467                constraint_id: activation.id(),
2468            }));
2469        }
2470
2471        Ok(None)
2472    }
2473
2474    /// Return whether one accepted source entity owns a live relation to a target.
2475    pub(in crate::db) fn entity_has_relation_to_target(
2476        &self,
2477        source_entity: EntityTag,
2478        target_path: &str,
2479    ) -> Result<bool, InternalError> {
2480        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2481            return Ok(false);
2482        };
2483        let Some(snapshot) = bundle.entity_snapshots().get(&source_entity) else {
2484            return Ok(false);
2485        };
2486
2487        Ok(snapshot
2488            .relations()
2489            .iter()
2490            .any(|relation| relation.target_path() == target_path))
2491    }
2492
2493    /// Reject any same-entity schema change beside one exact activation lifecycle step.
2494    pub(in crate::db) fn validate_live_activation_transition(
2495        &self,
2496        candidate: &AcceptedSchemaRevisionBundle,
2497    ) -> Result<(), InternalError> {
2498        let Some(current) = self.current_accepted_schema_bundle()? else {
2499            return Ok(());
2500        };
2501        Self::validate_activation_transition_from(&current, candidate)
2502    }
2503
2504    /// Validate one transition against the canonical accepted predecessor.
2505    pub(in crate::db) fn validate_canonical_activation_transition(
2506        &self,
2507        candidate: &AcceptedSchemaRevisionBundle,
2508    ) -> Result<(), InternalError> {
2509        let Some(current) = self.current_canonical_accepted_schema_bundle()? else {
2510            return Ok(());
2511        };
2512        Self::validate_activation_transition_from(&current, candidate)
2513    }
2514
2515    fn validate_activation_transition_from(
2516        current: &AcceptedSchemaRevisionBundle,
2517        candidate: &AcceptedSchemaRevisionBundle,
2518    ) -> Result<(), InternalError> {
2519        for (entity_tag, before) in current.entity_snapshots() {
2520            if before.constraint_activations().is_empty() {
2521                continue;
2522            }
2523            let after = candidate
2524                .entity_snapshots()
2525                .get(entity_tag)
2526                .ok_or_else(InternalError::store_invariant)?;
2527            if before == after {
2528                continue;
2529            }
2530            let expected_shape = before
2531                .clone()
2532                .with_constraint_catalog(after.constraint_catalog().clone());
2533            let catalog_only_transition = expected_shape == *after
2534                && before
2535                    .constraint_catalog()
2536                    .permits_live_activation_transition_to(after.constraint_catalog());
2537            let sql_row_local_abort_with_version =
2538                before.constraint_activations().iter().any(|activation| {
2539                    activation.origin() == ConstraintOrigin::SqlDdl
2540                        && matches!(
2541                            activation.kind(),
2542                            ConstraintActivationKind::Check { .. }
2543                                | ConstraintActivationKind::NotNull { .. }
2544                        )
2545                        && before.version().get().checked_add(1) == Some(after.version().get())
2546                        && before
2547                            .constraint_catalog()
2548                            .clone()
2549                            .with_aborted_activation(activation.id())
2550                            .is_ok_and(|catalog| catalog == *after.constraint_catalog())
2551                        && before
2552                            .clone()
2553                            .with_constraint_catalog(after.constraint_catalog().clone())
2554                            .with_schema_version(after.version())
2555                            == *after
2556                });
2557            let sql_unique_abort_with_version =
2558                before.constraint_activations().iter().any(|activation| {
2559                    activation.origin() == ConstraintOrigin::SqlDdl
2560                        && matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2561                        && before.version().get().checked_add(1) == Some(after.version().get())
2562                        && before
2563                            .with_aborted_unique_activation(activation.id(), after.version())
2564                            .is_ok_and(|expected| expected == *after)
2565                });
2566            let not_null_promotion = before.constraint_activations().iter().any(|activation| {
2567                matches!(activation.kind(), ConstraintActivationKind::NotNull { .. })
2568                    && before
2569                        .with_promoted_not_null_activation(activation.id(), after.version())
2570                        .is_ok_and(|expected| expected == *after)
2571            });
2572            let unique_promotion = before.constraint_activations().iter().any(|activation| {
2573                matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2574                    && before
2575                        .with_promoted_unique_activation(activation.id(), after.version())
2576                        .is_ok_and(|expected| expected == *after)
2577            });
2578            let relation_promotion = before.constraint_activations().iter().any(|activation| {
2579                matches!(activation.kind(), ConstraintActivationKind::Relation { .. })
2580                    && before
2581                        .with_promoted_relation_activation(activation.id(), after.version())
2582                        .is_ok_and(|expected| expected == *after)
2583            });
2584            if !catalog_only_transition
2585                && !sql_row_local_abort_with_version
2586                && !sql_unique_abort_with_version
2587                && !not_null_promotion
2588                && !unique_promotion
2589                && !relation_promotion
2590            {
2591                return Err(InternalError::store_invariant());
2592            }
2593        }
2594        Ok(())
2595    }
2596
2597    /// Prove exact pairing between live activations and durable validation jobs.
2598    pub(in crate::db) fn validate_constraint_validation_job_closure(
2599        &self,
2600        bundle: &AcceptedSchemaRevisionBundle,
2601    ) -> Result<(), InternalError> {
2602        self.validate_constraint_validation_job_closure_with_change(bundle, None, None)
2603    }
2604
2605    /// Prove the activation/job closure that would exist after one bounded
2606    /// marker-owned job replacement or removal.
2607    pub(in crate::db) fn validate_constraint_validation_job_closure_with_change(
2608        &self,
2609        bundle: &AcceptedSchemaRevisionBundle,
2610        replacement: Option<&ConstraintValidationJob>,
2611        removal: Option<(EntityTag, ConstraintId)>,
2612    ) -> Result<(), InternalError> {
2613        self.validate_constraint_validation_job_closure_with_change_in_view(
2614            bundle,
2615            replacement,
2616            removal,
2617            IdentityStateStorageView::Effective,
2618        )
2619    }
2620
2621    /// Prove activation/job closure against the canonical predecessor view.
2622    pub(in crate::db) fn validate_canonical_constraint_validation_job_closure_with_change(
2623        &self,
2624        bundle: &AcceptedSchemaRevisionBundle,
2625        replacement: Option<&ConstraintValidationJob>,
2626        removal: Option<(EntityTag, ConstraintId)>,
2627    ) -> Result<(), InternalError> {
2628        self.validate_constraint_validation_job_closure_with_change_in_view(
2629            bundle,
2630            replacement,
2631            removal,
2632            IdentityStateStorageView::Canonical,
2633        )
2634    }
2635
2636    fn validate_constraint_validation_job_closure_with_change_in_view(
2637        &self,
2638        bundle: &AcceptedSchemaRevisionBundle,
2639        replacement: Option<&ConstraintValidationJob>,
2640        removal: Option<(EntityTag, ConstraintId)>,
2641        view: IdentityStateStorageView,
2642    ) -> Result<(), InternalError> {
2643        if replacement.is_some() && removal.is_some() {
2644            return Err(InternalError::store_invariant());
2645        }
2646        let replacement_key = replacement.map(|job| {
2647            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id())
2648        });
2649        let removal_key = removal.map(|(entity_tag, constraint_id)| {
2650            RawSchemaKey::from_constraint_validation_job(entity_tag, constraint_id)
2651        });
2652        let mut expected = BTreeSet::new();
2653        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2654            for activation in snapshot.constraint_activations() {
2655                let key =
2656                    RawSchemaKey::from_constraint_validation_job(*entity_tag, activation.id());
2657                match activation.state() {
2658                    ConstraintActivationState::EnforcingNewWrites => {
2659                        if self
2660                            .constraint_validation_job_after_change(
2661                                key,
2662                                replacement,
2663                                replacement_key,
2664                                removal_key,
2665                                view,
2666                            )?
2667                            .is_some()
2668                        {
2669                            return Err(InternalError::store_corruption());
2670                        }
2671                    }
2672                    ConstraintActivationState::Validating => {
2673                        let job = self
2674                            .constraint_validation_job_after_change(
2675                                key,
2676                                replacement,
2677                                replacement_key,
2678                                removal_key,
2679                                view,
2680                            )?
2681                            .ok_or_else(InternalError::store_corruption)?;
2682                        if job.entity_tag() != *entity_tag
2683                            || job.entity_path() != snapshot.entity_path()
2684                        {
2685                            return Err(InternalError::store_corruption());
2686                        }
2687                        job.validate(Some(activation))?;
2688                        expected.insert(key);
2689                    }
2690                }
2691            }
2692        }
2693
2694        self.visit_constraint_validation_jobs_in_view(view, |key, raw| {
2695            if removal_key == Some(*key) || replacement_key == Some(*key) {
2696                return Ok(SchemaStoreVisit::Continue);
2697            }
2698            if !expected.contains(key) {
2699                return Err(InternalError::store_corruption());
2700            }
2701            let job = decode_constraint_validation_job(raw.as_bytes())?;
2702            if job.entity_tag() != key.entity_tag()
2703                || key.constraint_id() != Some(job.constraint_id())
2704            {
2705                return Err(InternalError::store_corruption());
2706            }
2707            Ok(SchemaStoreVisit::Continue)
2708        })?;
2709
2710        if let Some(key) = replacement_key
2711            && !expected.contains(&key)
2712        {
2713            return Err(InternalError::store_corruption());
2714        }
2715        if let Some(key) = removal_key
2716            && expected.contains(&key)
2717        {
2718            return Err(InternalError::store_corruption());
2719        }
2720
2721        Ok(())
2722    }
2723
2724    fn constraint_validation_job_after_change(
2725        &self,
2726        key: RawSchemaKey,
2727        replacement: Option<&ConstraintValidationJob>,
2728        replacement_key: Option<RawSchemaKey>,
2729        removal_key: Option<RawSchemaKey>,
2730        view: IdentityStateStorageView,
2731    ) -> Result<Option<ConstraintValidationJob>, InternalError> {
2732        if removal_key == Some(key) {
2733            return Ok(None);
2734        }
2735        if replacement_key == Some(key) {
2736            return Ok(replacement.cloned());
2737        }
2738        let raw = match view {
2739            IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
2740            IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
2741        };
2742        raw.map(|raw| decode_constraint_validation_job(raw.as_bytes()))
2743            .transpose()
2744    }
2745
2746    /// Return whether one retained schema authority still names this store's
2747    /// current immutable accepted root.
2748    pub(in crate::db) fn current_accepted_schema_authority_matches(
2749        &self,
2750        expected: &AcceptedSchemaAuthority,
2751    ) -> Result<bool, InternalError> {
2752        self.current_accepted_selection_for_authority(expected)
2753            .map(|selection| selection.is_some())
2754    }
2755
2756    // One store-scoped authority check for execution and borrowed metadata.
2757    // A dropped cache is not lost authority: reload its durable selection.
2758    fn current_accepted_selection_for_authority(
2759        &self,
2760        expected: &AcceptedSchemaAuthority,
2761    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
2762        let Some(store_scope) = self.accepted_catalog_scope.get() else {
2763            return Ok(None);
2764        };
2765
2766        // Root-writing primitives invalidate this cache before publication,
2767        // so a retained selection is the current in-memory authority.
2768        let cached = self
2769            .accepted_bundle_cache
2770            .try_borrow()
2771            .map_err(|_| InternalError::store_invariant())?
2772            .as_ref()
2773            .map(|cached| cached.selection);
2774        let selection = match cached {
2775            Some(selection) => Some(selection),
2776            None => self.current_accepted_schema_root()?,
2777        };
2778        Ok(selection.filter(|selection| {
2779            let root = selection.root();
2780            expected.matches_store_root(store_scope, root.revision(), root.fingerprint())
2781        }))
2782    }
2783
2784    /// Publish a candidate directly into its canonical schema allocation.
2785    ///
2786    /// Journaled online revisions must use
2787    /// `apply_journaled_accepted_schema_candidate`; this path owns initial
2788    /// bootstrap and marker-owned live-projection updates.
2789    pub(in crate::db) fn publish_accepted_schema_candidate(
2790        &mut self,
2791        incarnation: DatabaseIncarnationId,
2792        expected_revision: AcceptedSchemaRevision,
2793        candidate: &CandidateSchemaRevision,
2794    ) -> Result<(), InternalError> {
2795        let identity_transition = self.prepare_identity_state_transition(
2796            incarnation,
2797            candidate,
2798            IdentityStateStorageView::Effective,
2799        )?;
2800        if self.current_root_matches_candidate(candidate)? {
2801            if !identity_transition.is_empty() {
2802                return Err(InternalError::identity_state_corruption());
2803            }
2804            let selection = self
2805                .current_accepted_schema_root()?
2806                .ok_or_else(InternalError::store_corruption)?;
2807            self.retain_durable_candidate_entries(candidate, selection.slot())?;
2808            return Ok(());
2809        }
2810        let first = self.accepted_root_slot_bytes(0)?;
2811        let second = self.accepted_root_slot_bytes(1)?;
2812        prepare_accepted_schema_root_publication(
2813            [first.as_deref(), second.as_deref()],
2814            expected_revision,
2815            candidate,
2816        )
2817        .map_err(map_schema_publication_error)?;
2818
2819        self.insert_durable_candidate_snapshots(candidate)?;
2820        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2821        self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2822        let persisted_bundle = self
2823            .get_raw_snapshot(&bundle_key)
2824            .ok_or_else(InternalError::store_corruption)?;
2825        let _verified = decode_verified_accepted_schema_revision_bundle(
2826            candidate.root(),
2827            persisted_bundle.as_bytes(),
2828        )?;
2829        self.apply_identity_state_transition(
2830            identity_transition,
2831            IdentityStateWriteTarget::Durable,
2832        )?;
2833
2834        // Re-read the root immediately before the inactive-slot write. This is
2835        // the compare-and-swap check after candidate persistence.
2836        let first = self.accepted_root_slot_bytes(0)?;
2837        let second = self.accepted_root_slot_bytes(1)?;
2838        let publication = prepare_accepted_schema_root_publication(
2839            [first.as_deref(), second.as_deref()],
2840            expected_revision,
2841            candidate,
2842        )
2843        .map_err(map_schema_publication_error)?;
2844        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
2845        self.insert_durable_raw_value(root_key, publication.encoded_root().to_vec());
2846
2847        let selected = self
2848            .current_accepted_schema_root()?
2849            .ok_or_else(InternalError::store_corruption)?;
2850        if selected.root() != candidate.root() {
2851            return Err(InternalError::store_corruption());
2852        }
2853        self.retain_durable_candidate_entries(candidate, selected.slot())?;
2854        Ok(())
2855    }
2856
2857    /// Restore one current accepted candidate into an empty live-only schema
2858    /// store from its durable database-control checkpoint.
2859    pub(in crate::db) fn restore_live_accepted_schema_checkpoint(
2860        &mut self,
2861        incarnation: DatabaseIncarnationId,
2862        candidate: &CandidateSchemaRevision,
2863        checkpoint_identity_states: &IdentityStateInventory,
2864    ) -> Result<(), InternalError> {
2865        if !matches!(self.backend, SchemaStoreBackend::Heap(_)) {
2866            return Err(InternalError::store_invariant());
2867        }
2868        let checkpoint_validation = prepare_identity_state_transition(
2869            incarnation,
2870            Some(candidate.bundle()),
2871            candidate.bundle(),
2872            checkpoint_identity_states.clone(),
2873        )?;
2874        if !checkpoint_validation.is_empty() {
2875            return Err(InternalError::identity_state_corruption());
2876        }
2877        if self.current_root_matches_candidate(candidate)? {
2878            for state in checkpoint_identity_states.values() {
2879                let key = RawSchemaKey::from_identity_state(
2880                    state.owner().entity_tag(),
2881                    state.owner().field_id(),
2882                );
2883                self.insert_durable_raw_value(key, encode_identity_state(state)?);
2884            }
2885            if self.identity_state_inventory(IdentityStateStorageView::Effective)?
2886                != *checkpoint_identity_states
2887            {
2888                return Err(InternalError::identity_state_corruption());
2889            }
2890            let selection = self
2891                .current_accepted_schema_root()?
2892                .ok_or_else(InternalError::store_corruption)?;
2893            self.retain_durable_candidate_entries(candidate, selection.slot())?;
2894            return Ok(());
2895        }
2896        if self.current_accepted_schema_root()?.is_some()
2897            || !self
2898                .identity_state_inventory(IdentityStateStorageView::Effective)?
2899                .is_empty()
2900        {
2901            return Err(InternalError::store_corruption());
2902        }
2903
2904        self.insert_durable_candidate_snapshots(candidate)?;
2905        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2906        self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2907        for state in checkpoint_identity_states.values() {
2908            let key = RawSchemaKey::from_identity_state(
2909                state.owner().entity_tag(),
2910                state.owner().field_id(),
2911            );
2912            self.insert_durable_raw_value(key, encode_identity_state(state)?);
2913        }
2914        let root_key = RawSchemaKey::from_accepted_root_slot(0)?;
2915        self.insert_durable_raw_value(root_key, candidate.encoded_root().to_vec());
2916
2917        let selected = self
2918            .current_accepted_schema_root()?
2919            .ok_or_else(InternalError::store_corruption)?;
2920        if selected.root() != candidate.root() {
2921            return Err(InternalError::store_corruption());
2922        }
2923        self.retain_durable_candidate_entries(candidate, selected.slot())?;
2924        Ok(())
2925    }
2926
2927    /// Preflight one accepted candidate without changing durable or live
2928    /// schema state.
2929    ///
2930    /// Returns `true` only when this exact candidate is already authoritative.
2931    /// Multi-store publication uses that distinction to reject partial replay
2932    /// before opening one marker-owned commit window.
2933    pub(in crate::db) fn preflight_accepted_schema_candidate(
2934        &self,
2935        incarnation: DatabaseIncarnationId,
2936        expected_revision: AcceptedSchemaRevision,
2937        candidate: &CandidateSchemaRevision,
2938    ) -> Result<bool, InternalError> {
2939        let identity_transition = self.prepare_identity_state_transition(
2940            incarnation,
2941            candidate,
2942            IdentityStateStorageView::Effective,
2943        )?;
2944        if self.current_root_matches_candidate(candidate)? {
2945            if !identity_transition.is_empty() {
2946                return Err(InternalError::identity_state_corruption());
2947            }
2948            return Ok(true);
2949        }
2950        let first = self.accepted_root_slot_bytes(0)?;
2951        let second = self.accepted_root_slot_bytes(1)?;
2952        prepare_accepted_schema_root_publication(
2953            [first.as_deref(), second.as_deref()],
2954            expected_revision,
2955            candidate,
2956        )
2957        .map_err(map_schema_publication_error)?;
2958
2959        Ok(false)
2960    }
2961
2962    /// Prepare one accepted candidate against canonical journaled authority.
2963    pub(in crate::db) fn prepare_fold_journaled_accepted_schema_candidate(
2964        &self,
2965        incarnation: DatabaseIncarnationId,
2966        expected_revision: AcceptedSchemaRevision,
2967        candidate: CandidateSchemaRevision,
2968    ) -> Result<PreparedAcceptedSchemaFold, InternalError> {
2969        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2970            return Err(InternalError::store_invariant());
2971        }
2972        let identity_transition = self.prepare_identity_state_transition(
2973            incarnation,
2974            &candidate,
2975            IdentityStateStorageView::Canonical,
2976        )?;
2977        let candidate_is_current = self.canonical_root_matches_candidate(&candidate)?;
2978        if candidate_is_current && !identity_transition.is_empty() {
2979            return Err(InternalError::identity_state_corruption());
2980        }
2981        let identity_updates = identity_transition
2982            .into_updates()
2983            .into_iter()
2984            .map(|state| {
2985                Ok((
2986                    RawSchemaKey::from_identity_state(
2987                        state.owner().entity_tag(),
2988                        state.owner().field_id(),
2989                    ),
2990                    encode_identity_state(&state)?,
2991                ))
2992            })
2993            .collect::<Result<Vec<_>, InternalError>>()?;
2994        let snapshots = candidate
2995            .bundle()
2996            .entity_snapshots()
2997            .iter()
2998            .map(|(entity, snapshot)| Self::prepare_persisted_snapshot(*entity, snapshot))
2999            .collect::<Result<Vec<_>, _>>()?;
3000
3001        let first = self.canonical_root_slot_bytes(0)?;
3002        let second = self.canonical_root_slot_bytes(1)?;
3003        let root_slot = if candidate_is_current {
3004            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3005                .ok_or_else(InternalError::store_corruption)?
3006                .slot()
3007        } else {
3008            prepare_accepted_schema_root_publication(
3009                [first.as_deref(), second.as_deref()],
3010                expected_revision,
3011                &candidate,
3012            )
3013            .map_err(map_schema_publication_error)?
3014            .target_slot()
3015        };
3016        let retained = Self::candidate_entry_keys(&candidate, root_slot)?;
3017        Ok(PreparedAcceptedSchemaFold {
3018            candidate,
3019            expected_revision,
3020            snapshots,
3021            identity_updates,
3022            retained,
3023            root_slot,
3024        })
3025    }
3026
3027    /// Return the retained Identity owner count after admitting one candidate.
3028    pub(in crate::db) fn projected_identity_state_count(
3029        &self,
3030        incarnation: DatabaseIncarnationId,
3031        candidate: &CandidateSchemaRevision,
3032    ) -> Result<usize, InternalError> {
3033        Ok(self
3034            .prepare_identity_state_transition(
3035                incarnation,
3036                candidate,
3037                IdentityStateStorageView::Effective,
3038            )?
3039            .projected_inventory_len())
3040    }
3041
3042    /// Apply one marker-bound schema candidate to the journaled live projection.
3043    pub(in crate::db) fn apply_journaled_accepted_schema_candidate(
3044        &mut self,
3045        incarnation: DatabaseIncarnationId,
3046        expected_revision: AcceptedSchemaRevision,
3047        candidate: &CandidateSchemaRevision,
3048    ) -> Result<(), InternalError> {
3049        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3050            return Err(InternalError::store_invariant());
3051        }
3052        let identity_transition = self.prepare_identity_state_transition(
3053            incarnation,
3054            candidate,
3055            IdentityStateStorageView::Effective,
3056        )?;
3057        if self.current_root_matches_candidate(candidate)? {
3058            if !identity_transition.is_empty() {
3059                return Err(InternalError::identity_state_corruption());
3060            }
3061            let selection = self
3062                .current_accepted_schema_root()?
3063                .ok_or_else(InternalError::store_corruption)?;
3064            self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3065            return Ok(());
3066        }
3067
3068        let first = self.accepted_root_slot_bytes(0)?;
3069        let second = self.accepted_root_slot_bytes(1)?;
3070        prepare_accepted_schema_root_publication(
3071            [first.as_deref(), second.as_deref()],
3072            expected_revision,
3073            candidate,
3074        )
3075        .map_err(map_schema_publication_error)?;
3076
3077        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3078            self.insert_persisted_snapshot(*entity_tag, snapshot)?;
3079        }
3080        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3081        self.insert_raw_snapshot(
3082            bundle_key,
3083            RawSchemaSnapshot::from_encoded_control_record(candidate.encoded_bundle().to_vec()),
3084        );
3085        let persisted_bundle = self
3086            .get_raw_snapshot(&bundle_key)
3087            .ok_or_else(InternalError::store_corruption)?;
3088        let _verified = decode_verified_accepted_schema_revision_bundle(
3089            candidate.root(),
3090            persisted_bundle.as_bytes(),
3091        )?;
3092        self.apply_identity_state_transition(
3093            identity_transition,
3094            IdentityStateWriteTarget::Materialized,
3095        )?;
3096
3097        let first = self.accepted_root_slot_bytes(0)?;
3098        let second = self.accepted_root_slot_bytes(1)?;
3099        let publication = prepare_accepted_schema_root_publication(
3100            [first.as_deref(), second.as_deref()],
3101            expected_revision,
3102            candidate,
3103        )
3104        .map_err(map_schema_publication_error)?;
3105        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3106        self.insert_raw_snapshot(
3107            root_key,
3108            RawSchemaSnapshot::from_encoded_control_record(publication.encoded_root().to_vec()),
3109        );
3110
3111        if !self.current_root_matches_candidate(candidate)? {
3112            return Err(InternalError::store_corruption());
3113        }
3114        let selection = self
3115            .current_accepted_schema_root()?
3116            .ok_or_else(InternalError::store_corruption)?;
3117        self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3118        Ok(())
3119    }
3120
3121    /// Consume one candidate's preflight payloads in the same atomic fold callback.
3122    pub(in crate::db) fn apply_prepared_accepted_schema_fold(
3123        &mut self,
3124        prepared: PreparedAcceptedSchemaFold,
3125    ) -> Result<(), InternalError> {
3126        let PreparedAcceptedSchemaFold {
3127            candidate,
3128            expected_revision,
3129            snapshots,
3130            identity_updates,
3131            retained,
3132            root_slot,
3133        } = prepared;
3134        if self.canonical_root_matches_candidate(&candidate)? {
3135            if !identity_updates.is_empty() {
3136                return Err(InternalError::identity_state_corruption());
3137            }
3138            let first = self.canonical_root_slot_bytes(0)?;
3139            let second = self.canonical_root_slot_bytes(1)?;
3140            let selection =
3141                select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3142                    .ok_or_else(InternalError::store_corruption)?;
3143            if selection.slot() != root_slot {
3144                return Err(InternalError::store_invariant());
3145            }
3146            self.retain_canonical_candidate_entries(&retained)?;
3147            return Ok(());
3148        }
3149
3150        let first = self.canonical_root_slot_bytes(0)?;
3151        let second = self.canonical_root_slot_bytes(1)?;
3152        let publication = prepare_accepted_schema_root_publication(
3153            [first.as_deref(), second.as_deref()],
3154            expected_revision,
3155            &candidate,
3156        )
3157        .map_err(map_schema_publication_error)?;
3158        if publication.target_slot() != root_slot {
3159            return Err(InternalError::store_invariant());
3160        }
3161
3162        for snapshot in snapshots {
3163            self.apply_prepared_fold_persisted_snapshot(snapshot)?;
3164        }
3165        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3166        self.insert_canonical_raw_value(bundle_key, candidate.encoded_bundle().to_vec())?;
3167        let persisted_bundle = self
3168            .get_canonical_raw_value(&bundle_key)?
3169            .ok_or_else(InternalError::store_corruption)?;
3170        let _verified = decode_verified_accepted_schema_revision_bundle(
3171            candidate.root(),
3172            persisted_bundle.as_bytes(),
3173        )?;
3174        for (key, bytes) in identity_updates {
3175            self.insert_canonical_raw_value(key, bytes)?;
3176        }
3177
3178        let first = self.canonical_root_slot_bytes(0)?;
3179        let second = self.canonical_root_slot_bytes(1)?;
3180        let publication = prepare_accepted_schema_root_publication(
3181            [first.as_deref(), second.as_deref()],
3182            expected_revision,
3183            &candidate,
3184        )
3185        .map_err(map_schema_publication_error)?;
3186        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3187        self.insert_canonical_raw_value(root_key, publication.encoded_root().to_vec())?;
3188
3189        if !self.canonical_root_matches_candidate(&candidate)? {
3190            return Err(InternalError::store_corruption());
3191        }
3192        let first = self.canonical_root_slot_bytes(0)?;
3193        let second = self.canonical_root_slot_bytes(1)?;
3194        let selection = select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3195            .ok_or_else(InternalError::store_corruption)?;
3196        if selection.slot() != root_slot {
3197            return Err(InternalError::store_invariant());
3198        }
3199        self.retain_canonical_candidate_entries(&retained)?;
3200        Ok(())
3201    }
3202
3203    /// Load and decode one typed persisted schema snapshot.
3204    pub(in crate::db) fn get_persisted_snapshot(
3205        &self,
3206        entity: EntityTag,
3207        version: SchemaVersion,
3208    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3209        let key = RawSchemaKey::from_entity_version(entity, version);
3210        self.get_raw_snapshot(&key)
3211            .map(|snapshot| snapshot.decode_persisted_snapshot())
3212            .transpose()
3213    }
3214
3215    #[cfg(test)]
3216    fn latest_staged_persisted_snapshot(
3217        &self,
3218        entity: EntityTag,
3219    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3220        self.latest_raw_snapshots_by_entity()
3221            .remove(&entity)
3222            .map(|(_, snapshot)| snapshot.decode_persisted_snapshot())
3223            .transpose()
3224    }
3225
3226    /// Load one entity snapshot from the immutable bundle selected by the
3227    /// current accepted root.
3228    pub(in crate::db) fn current_accepted_persisted_snapshot(
3229        &self,
3230        entity: EntityTag,
3231    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3232        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3233            return Ok(None);
3234        };
3235
3236        Ok(bundle.entity_snapshots().get(&entity).cloned())
3237    }
3238
3239    /// Return one accepted catalog selection from the current immutable root.
3240    pub(in crate::db) fn current_accepted_catalog_selection(
3241        &self,
3242        entity: EntityTag,
3243        entity_path: &str,
3244        store_path: &'static str,
3245    ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3246        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3247            return Ok(None);
3248        };
3249        if bundle.store_path() != store_path {
3250            return Err(InternalError::store_corruption());
3251        }
3252        let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3253            return Ok(None);
3254        };
3255        if snapshot.entity_path() != entity_path {
3256            return Err(InternalError::store_corruption());
3257        }
3258
3259        let cache = self
3260            .accepted_bundle_cache
3261            .try_borrow()
3262            .map_err(|_| InternalError::store_invariant())?;
3263        let cached = cache.as_ref().ok_or_else(InternalError::store_invariant)?;
3264        if let Some(selection) = cached
3265            .entity_selections
3266            .try_borrow()
3267            .map_err(|_| InternalError::store_invariant())?
3268            .get(&entity)
3269            .cloned()
3270        {
3271            return Ok(Some(selection));
3272        }
3273
3274        let selected = AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3275            entity,
3276            store_path,
3277            snapshot,
3278            cached.value_catalog.clone(),
3279        )?;
3280        cached
3281            .entity_selections
3282            .try_borrow_mut()
3283            .map_err(|_| InternalError::store_invariant())?
3284            .insert(entity, selected.clone());
3285
3286        Ok(Some(selected))
3287    }
3288
3289    /// Return one accepted catalog selection from the canonical journal base.
3290    /// Recovery uses this while folding historical row batches whose schema
3291    /// revision can precede the current live accepted root.
3292    pub(in crate::db) fn current_canonical_accepted_catalog_selection(
3293        &self,
3294        entity: EntityTag,
3295        entity_path: &str,
3296        store_path: &'static str,
3297    ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3298        let first = self.canonical_root_slot_bytes(0)?;
3299        let second = self.canonical_root_slot_bytes(1)?;
3300        let Some(selection) =
3301            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3302        else {
3303            return Ok(None);
3304        };
3305        let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3306        let raw_bundle = self
3307            .get_canonical_raw_value(&bundle_key)?
3308            .ok_or_else(InternalError::store_corruption)?;
3309        let bundle = decode_verified_accepted_schema_revision_bundle(
3310            selection.root(),
3311            raw_bundle.as_bytes(),
3312        )?;
3313        if bundle.store_path() != store_path {
3314            return Err(InternalError::store_corruption());
3315        }
3316        let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3317            return Ok(None);
3318        };
3319        if snapshot.entity_path() != entity_path {
3320            return Err(InternalError::store_corruption());
3321        }
3322
3323        AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3324            entity,
3325            store_path,
3326            snapshot,
3327            AcceptedValueCatalogHandle::new(
3328                bundle.enum_catalog().clone(),
3329                bundle.composite_catalog().clone(),
3330                self.accepted_catalog_scope
3331                    .get_or_init(AcceptedStoreCatalogScope::new)
3332                    .clone(),
3333                bundle.revision(),
3334                selection.root().fingerprint(),
3335            ),
3336        )
3337        .map(Some)
3338    }
3339
3340    /// Derive accepted catalog metadata from latest persisted schema snapshots.
3341    ///
3342    /// This function intentionally reads only the persisted schema store. It
3343    /// does not reconstruct metadata from generated models when the store has
3344    /// no accepted snapshots.
3345    #[cfg(test)]
3346    pub(in crate::db) fn catalog_metadata(
3347        &self,
3348    ) -> Result<Option<SchemaStoreCatalogMetadata>, InternalError> {
3349        Ok(self
3350            .allocation_metadata()?
3351            .map(SchemaStoreAllocationMetadata::schema))
3352    }
3353
3354    /// Derive role-specific allocation metadata from latest persisted schema
3355    /// snapshots.
3356    ///
3357    /// This function intentionally reads only accepted schema-store payloads.
3358    /// It never reconstructs metadata from generated models when the store has
3359    /// no accepted snapshots.
3360    pub(in crate::db) fn allocation_metadata(
3361        &self,
3362    ) -> Result<Option<SchemaStoreAllocationMetadata>, InternalError> {
3363        let latest_by_entity = self.latest_raw_snapshots_by_entity();
3364        if latest_by_entity.is_empty() {
3365            return Ok(None);
3366        }
3367
3368        Ok(Some(SchemaStoreAllocationMetadata::new(
3369            derive_data_allocation_metadata(&latest_by_entity)?,
3370            derive_index_allocation_metadata(&latest_by_entity)?,
3371            derive_schema_catalog_metadata(&latest_by_entity)?,
3372        )))
3373    }
3374
3375    /// Insert or replace one raw schema snapshot.
3376    fn insert_raw_snapshot(
3377        &mut self,
3378        key: RawSchemaKey,
3379        snapshot: RawSchemaSnapshot,
3380    ) -> Option<RawSchemaSnapshot> {
3381        self.invalidate_accepted_bundle_cache_for_key(key);
3382        let previous_journaled = if matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3383            self.get_raw_snapshot_for_backend(&key)
3384        } else {
3385            None
3386        };
3387        match &mut self.backend {
3388            SchemaStoreBackend::Heap(map) => map.insert(key, snapshot),
3389            SchemaStoreBackend::Journaled {
3390                live, tombstones, ..
3391            } => {
3392                tombstones.remove(&key);
3393                live.insert(key, snapshot);
3394                previous_journaled
3395            }
3396        }
3397    }
3398
3399    /// Load one raw schema snapshot by key.
3400    #[must_use]
3401    fn get_raw_snapshot(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
3402        match &self.backend {
3403            SchemaStoreBackend::Heap(map) => map.get(key).cloned(),
3404            SchemaStoreBackend::Journaled { .. } => self.get_raw_snapshot_for_backend(key),
3405        }
3406    }
3407
3408    fn accepted_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3409        let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3410        Ok(self
3411            .get_raw_snapshot(&key)
3412            .map(RawSchemaSnapshot::into_bytes))
3413    }
3414
3415    fn canonical_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3416        let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3417        Ok(self
3418            .get_canonical_raw_value(&key)?
3419            .map(RawSchemaSnapshot::into_bytes))
3420    }
3421
3422    fn current_root_matches_candidate(
3423        &self,
3424        candidate: &CandidateSchemaRevision,
3425    ) -> Result<bool, InternalError> {
3426        let Some(selection) = self.current_accepted_schema_root()? else {
3427            return Ok(false);
3428        };
3429        if selection.root() != candidate.root() {
3430            return Ok(false);
3431        }
3432        let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3433        let bundle = self
3434            .get_raw_snapshot(&key)
3435            .ok_or_else(InternalError::store_corruption)?;
3436        let _verified =
3437            decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3438        Ok(true)
3439    }
3440
3441    fn canonical_root_matches_candidate(
3442        &self,
3443        candidate: &CandidateSchemaRevision,
3444    ) -> Result<bool, InternalError> {
3445        let first = self.canonical_root_slot_bytes(0)?;
3446        let second = self.canonical_root_slot_bytes(1)?;
3447        let Some(selection) =
3448            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3449        else {
3450            return Ok(false);
3451        };
3452        if selection.root() != candidate.root() {
3453            return Ok(false);
3454        }
3455        let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3456        let bundle = self
3457            .get_canonical_raw_value(&key)?
3458            .ok_or_else(InternalError::store_corruption)?;
3459        let _verified =
3460            decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3461        Ok(true)
3462    }
3463
3464    fn get_canonical_raw_value(
3465        &self,
3466        key: &RawSchemaKey,
3467    ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
3468        match &self.backend {
3469            SchemaStoreBackend::Journaled { canonical, .. } => Ok(canonical.get(key)),
3470            SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
3471        }
3472    }
3473
3474    fn insert_canonical_raw_value(
3475        &mut self,
3476        key: RawSchemaKey,
3477        bytes: Vec<u8>,
3478    ) -> Result<(), InternalError> {
3479        self.invalidate_accepted_bundle_cache_for_key(key);
3480        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3481            return Err(InternalError::store_invariant());
3482        };
3483        canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
3484        Ok(())
3485    }
3486
3487    // Initial accepted-catalog bootstrap persists immutable bundle/root values
3488    // directly in the schema allocation. Later online schema mutation will
3489    // carry the same values through the journal before calling this primitive.
3490    fn insert_durable_raw_value(&mut self, key: RawSchemaKey, bytes: Vec<u8>) {
3491        self.invalidate_accepted_bundle_cache_for_key(key);
3492        let value = RawSchemaSnapshot::from_encoded_control_record(bytes);
3493        match &mut self.backend {
3494            SchemaStoreBackend::Heap(map) => {
3495                map.insert(key, value);
3496            }
3497            SchemaStoreBackend::Journaled {
3498                canonical,
3499                live,
3500                tombstones,
3501                ..
3502            } => {
3503                live.remove(&key);
3504                tombstones.remove(&key);
3505                canonical.insert(key, value);
3506            }
3507        }
3508    }
3509
3510    fn invalidate_accepted_bundle_cache_for_key(&mut self, key: RawSchemaKey) {
3511        if key.is_accepted_root() {
3512            self.accepted_bundle_cache.get_mut().take();
3513        }
3514    }
3515
3516    fn insert_durable_candidate_snapshots(
3517        &mut self,
3518        candidate: &CandidateSchemaRevision,
3519    ) -> Result<(), InternalError> {
3520        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3521            let key = RawSchemaKey::from_entity_version(*entity_tag, snapshot.version());
3522            let value = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3523            match &mut self.backend {
3524                SchemaStoreBackend::Heap(map) => {
3525                    map.insert(key, value);
3526                }
3527                SchemaStoreBackend::Journaled {
3528                    canonical,
3529                    live,
3530                    tombstones,
3531                    ..
3532                } => {
3533                    live.remove(&key);
3534                    tombstones.remove(&key);
3535                    canonical.insert(key, value);
3536                }
3537            }
3538        }
3539        Ok(())
3540    }
3541
3542    fn candidate_entry_keys(
3543        candidate: &CandidateSchemaRevision,
3544        root_slot: usize,
3545    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3546        let mut keys = candidate
3547            .bundle()
3548            .entity_snapshots()
3549            .iter()
3550            .map(|(entity_tag, snapshot)| {
3551                RawSchemaKey::from_entity_version(*entity_tag, snapshot.version())
3552            })
3553            .collect::<BTreeSet<_>>();
3554        keys.insert(RawSchemaKey::from_accepted_bundle(
3555            candidate.root().bundle_key(),
3556        ));
3557        keys.insert(RawSchemaKey::from_accepted_root_slot(root_slot)?);
3558        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3559            for activation in snapshot
3560                .constraint_activations()
3561                .iter()
3562                .filter(|activation| activation.state() == ConstraintActivationState::Validating)
3563            {
3564                keys.insert(RawSchemaKey::from_constraint_validation_job(
3565                    *entity_tag,
3566                    activation.id(),
3567                ));
3568            }
3569        }
3570        Ok(keys)
3571    }
3572
3573    fn positioned_candidate_effect_keys(
3574        &self,
3575        incarnation: DatabaseIncarnationId,
3576        expected_revision: AcceptedSchemaRevision,
3577        candidate: &CandidateSchemaRevision,
3578        view: IdentityStateStorageView,
3579    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3580        let identity_transition =
3581            self.prepare_identity_state_transition(incarnation, candidate, view)?;
3582        let (first, second, candidate_is_current) = match view {
3583            IdentityStateStorageView::Effective => (
3584                self.accepted_root_slot_bytes(0)?,
3585                self.accepted_root_slot_bytes(1)?,
3586                self.current_root_matches_candidate(candidate)?,
3587            ),
3588            IdentityStateStorageView::Canonical => (
3589                self.canonical_root_slot_bytes(0)?,
3590                self.canonical_root_slot_bytes(1)?,
3591                self.canonical_root_matches_candidate(candidate)?,
3592            ),
3593        };
3594        let root_slot = if candidate_is_current {
3595            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3596                .ok_or_else(InternalError::store_corruption)?
3597                .slot()
3598        } else {
3599            prepare_accepted_schema_root_publication(
3600                [first.as_deref(), second.as_deref()],
3601                expected_revision,
3602                candidate,
3603            )
3604            .map_err(map_schema_publication_error)?
3605            .target_slot()
3606        };
3607        let mut keys = Self::candidate_entry_keys(candidate, root_slot)?;
3608        for state in identity_transition.into_updates() {
3609            keys.insert(RawSchemaKey::from_identity_state(
3610                state.owner().entity_tag(),
3611                state.owner().field_id(),
3612            ));
3613        }
3614
3615        let SchemaStoreBackend::Journaled {
3616            canonical,
3617            live,
3618            tombstones,
3619            positions,
3620        } = &self.backend
3621        else {
3622            return Err(InternalError::store_invariant());
3623        };
3624        for entry in canonical.iter() {
3625            let has_relevant_overlay = matches!(view, IdentityStateStorageView::Effective)
3626                || positions.is_positioned(entry.key())
3627                || live.contains_key(entry.key())
3628                || tombstones.contains(entry.key());
3629            if has_relevant_overlay
3630                && !keys.contains(entry.key())
3631                && !entry.key().is_identity_state()
3632            {
3633                keys.insert(*entry.key());
3634            }
3635        }
3636        if matches!(view, IdentityStateStorageView::Effective) {
3637            for key in live.keys() {
3638                if !keys.contains(key) && !key.is_identity_state() {
3639                    keys.insert(*key);
3640                }
3641            }
3642        }
3643        Ok(keys)
3644    }
3645
3646    fn positioned_journal_batch_keys(
3647        &self,
3648        incarnation: DatabaseIncarnationId,
3649        batch: &JournalBatch,
3650        view: IdentityStateStorageView,
3651    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3652        let mut keys = BTreeSet::new();
3653        for record in batch.records() {
3654            match record {
3655                JournalRecord::SchemaPut {
3656                    schema_snapshot_bytes,
3657                    ..
3658                } => {
3659                    let snapshot = decode_persisted_schema_snapshot(schema_snapshot_bytes)?;
3660                    let entity_tag = match view {
3661                        IdentityStateStorageView::Effective => self
3662                            .current_accepted_schema_bundle_ref()?
3663                            .ok_or_else(InternalError::store_corruption)?
3664                            .entity_snapshots()
3665                            .iter()
3666                            .find_map(|(entity_tag, accepted)| {
3667                                (accepted.entity_path() == snapshot.entity_path())
3668                                    .then_some(*entity_tag)
3669                            }),
3670                        IdentityStateStorageView::Canonical => self
3671                            .current_canonical_accepted_schema_bundle()?
3672                            .ok_or_else(InternalError::store_corruption)?
3673                            .entity_snapshots()
3674                            .iter()
3675                            .find_map(|(entity_tag, accepted)| {
3676                                (accepted.entity_path() == snapshot.entity_path())
3677                                    .then_some(*entity_tag)
3678                            }),
3679                    }
3680                    .ok_or_else(InternalError::store_corruption)?;
3681                    keys.insert(RawSchemaKey::from_entity_version(
3682                        entity_tag,
3683                        snapshot.version(),
3684                    ));
3685                }
3686                JournalRecord::AcceptedSchemaPublish {
3687                    expected_revision,
3688                    schema_bundle_bytes,
3689                    schema_root_bytes,
3690                    ..
3691                } => {
3692                    let candidate = CandidateSchemaRevision::from_encoded(
3693                        schema_bundle_bytes.clone(),
3694                        schema_root_bytes.clone(),
3695                    )?;
3696                    keys.extend(self.positioned_candidate_effect_keys(
3697                        incarnation,
3698                        *expected_revision,
3699                        &candidate,
3700                        view,
3701                    )?);
3702                }
3703                JournalRecord::ConstraintValidationJobPut {
3704                    entity_tag,
3705                    constraint_id,
3706                    ..
3707                }
3708                | JournalRecord::ConstraintValidationJobDelete {
3709                    entity_tag,
3710                    constraint_id,
3711                    ..
3712                } => {
3713                    keys.insert(RawSchemaKey::from_constraint_validation_job(
3714                        *entity_tag,
3715                        *constraint_id,
3716                    ));
3717                }
3718                JournalRecord::IdentityRangeAdvance { range } => {
3719                    keys.insert(RawSchemaKey::from_identity_state(
3720                        range.owner().entity_tag(),
3721                        range.owner().field_id(),
3722                    ));
3723                }
3724                JournalRecord::RowPut { .. }
3725                | JournalRecord::RowDelete { .. }
3726                | JournalRecord::AcceptedSchemaIndexDelete { .. }
3727                | JournalRecord::AcceptedSchemaIndexPut { .. }
3728                | JournalRecord::ConstraintValidationIndexPut { .. } => {}
3729                #[cfg(any(test, feature = "migration"))]
3730                JournalRecord::SchemaMigrationRowPut { .. }
3731                | JournalRecord::SchemaMigrationIndexPut { .. } => {}
3732            }
3733        }
3734        Ok(keys)
3735    }
3736
3737    // Keep only the current entity snapshots, immutable bundle, and selected
3738    // root. The inactive root is needed only during publication and is removed
3739    // after the new root has been verified.
3740    fn retain_durable_candidate_entries(
3741        &mut self,
3742        candidate: &CandidateSchemaRevision,
3743        root_slot: usize,
3744    ) -> Result<(), InternalError> {
3745        let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3746        self.accepted_bundle_cache.get_mut().take();
3747        match &mut self.backend {
3748            SchemaStoreBackend::Heap(map) => {
3749                map.retain(|key, _| keep.contains(key) || key.is_identity_state());
3750            }
3751            SchemaStoreBackend::Journaled {
3752                canonical,
3753                live,
3754                tombstones,
3755                ..
3756            } => {
3757                let stale = canonical
3758                    .iter()
3759                    .filter_map(|entry| {
3760                        (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3761                            .then_some(*entry.key())
3762                    })
3763                    .collect::<Vec<_>>();
3764                for key in stale {
3765                    canonical.remove(&key);
3766                }
3767                live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3768                tombstones.clear();
3769            }
3770        }
3771        Ok(())
3772    }
3773
3774    fn retain_materialized_candidate_entries(
3775        &mut self,
3776        candidate: &CandidateSchemaRevision,
3777        root_slot: usize,
3778    ) -> Result<(), InternalError> {
3779        let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3780        self.accepted_bundle_cache.get_mut().take();
3781        let SchemaStoreBackend::Journaled {
3782            canonical,
3783            live,
3784            tombstones,
3785            ..
3786        } = &mut self.backend
3787        else {
3788            return Err(InternalError::store_invariant());
3789        };
3790        live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3791        let canonical_keys = canonical
3792            .iter()
3793            .map(|entry| *entry.key())
3794            .collect::<Vec<_>>();
3795        for key in canonical_keys {
3796            if keep.contains(&key) || key.is_identity_state() {
3797                tombstones.remove(&key);
3798            } else {
3799                tombstones.insert(key);
3800            }
3801        }
3802        Ok(())
3803    }
3804
3805    fn retain_canonical_candidate_entries(
3806        &mut self,
3807        keep: &BTreeSet<RawSchemaKey>,
3808    ) -> Result<(), InternalError> {
3809        self.accepted_bundle_cache.get_mut().take();
3810        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3811            return Err(InternalError::store_invariant());
3812        };
3813        let stale = canonical
3814            .iter()
3815            .filter_map(|entry| {
3816                (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3817                    .then_some(*entry.key())
3818            })
3819            .collect::<Vec<_>>();
3820        for key in stale {
3821            canonical.remove(&key);
3822        }
3823        Ok(())
3824    }
3825
3826    /// Return whether one schema snapshot key is present.
3827    #[must_use]
3828    #[cfg(test)]
3829    fn contains_raw_snapshot(&self, key: &RawSchemaKey) -> bool {
3830        match &self.backend {
3831            SchemaStoreBackend::Heap(map) => map.contains_key(key),
3832            SchemaStoreBackend::Journaled { .. } => {
3833                self.get_raw_snapshot_for_backend(key).is_some()
3834            }
3835        }
3836    }
3837
3838    /// Return the number of schema snapshot entries in this store.
3839    #[must_use]
3840    #[cfg(test)]
3841    pub(in crate::db) fn len(&self) -> u64 {
3842        match &self.backend {
3843            SchemaStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
3844            SchemaStoreBackend::Journaled { .. } => {
3845                let mut count = 0_u64;
3846                let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3847                    count = count.saturating_add(1);
3848                    Ok(SchemaStoreVisit::Continue)
3849                });
3850                count
3851            }
3852        }
3853    }
3854
3855    /// Return whether this schema store currently has no persisted snapshots.
3856    #[must_use]
3857    #[cfg(test)]
3858    pub(in crate::db) fn is_empty(&self) -> bool {
3859        match &self.backend {
3860            SchemaStoreBackend::Heap(map) => map.is_empty(),
3861            SchemaStoreBackend::Journaled { .. } => {
3862                let mut empty = true;
3863                let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3864                    empty = false;
3865                    Ok(SchemaStoreVisit::Stop)
3866                });
3867                empty
3868            }
3869        }
3870    }
3871
3872    /// Clear all schema metadata entries from the store.
3873    #[cfg(test)]
3874    pub(in crate::db) fn clear(&mut self) {
3875        self.accepted_bundle_cache.get_mut().take();
3876        match &mut self.backend {
3877            SchemaStoreBackend::Heap(map) => map.clear(),
3878            SchemaStoreBackend::Journaled {
3879                canonical,
3880                live,
3881                tombstones,
3882                ..
3883            } => {
3884                live.clear();
3885                tombstones.clear();
3886                let keys = canonical
3887                    .iter()
3888                    .map(|entry| *entry.key())
3889                    .collect::<Vec<_>>();
3890                for key in keys {
3891                    if key.is_entity_snapshot() {
3892                        tombstones.insert(key);
3893                    } else {
3894                        canonical.remove(&key);
3895                    }
3896                }
3897            }
3898        }
3899    }
3900
3901    fn current_accepted_schema_bundle_ref(
3902        &self,
3903    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
3904        self.current_accepted_schema_authority_ref()
3905            .map(|authority| authority.map(|(_selection, bundle)| bundle))
3906    }
3907
3908    /// Borrow the effective accepted root and its cached, verified immutable bundle together.
3909    pub(in crate::db) fn current_accepted_schema_authority_ref(
3910        &self,
3911    ) -> Result<
3912        Option<(
3913            AcceptedSchemaRootSelection,
3914            Ref<'_, AcceptedSchemaRevisionBundle>,
3915        )>,
3916        InternalError,
3917    > {
3918        let selection = self.current_accepted_schema_root()?;
3919        self.accepted_schema_authority_ref_for_selection(selection)
3920    }
3921
3922    fn accepted_schema_authority_ref_for_selection(
3923        &self,
3924        selection: Option<AcceptedSchemaRootSelection>,
3925    ) -> Result<
3926        Option<(
3927            AcceptedSchemaRootSelection,
3928            Ref<'_, AcceptedSchemaRevisionBundle>,
3929        )>,
3930        InternalError,
3931    > {
3932        let Some(selection) = selection else {
3933            self.accepted_bundle_cache
3934                .try_borrow_mut()
3935                .map_err(|_| InternalError::store_invariant())?
3936                .take();
3937            return Ok(None);
3938        };
3939
3940        let cache_matches = self
3941            .accepted_bundle_cache
3942            .try_borrow()
3943            .map_err(|_| InternalError::store_invariant())?
3944            .as_ref()
3945            .is_some_and(|cached| cached.selection == selection);
3946        if !cache_matches {
3947            let key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3948            let raw = self
3949                .get_raw_snapshot(&key)
3950                .ok_or_else(InternalError::store_corruption)?;
3951            let bundle =
3952                decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
3953            self.validate_constraint_validation_job_closure(&bundle)?;
3954            #[cfg(test)]
3955            ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES
3956                .with(|misses| misses.set(misses.get().saturating_add(1)));
3957            let value_catalog = AcceptedValueCatalogHandle::new(
3958                bundle.enum_catalog().clone(),
3959                bundle.composite_catalog().clone(),
3960                self.accepted_catalog_scope
3961                    .get_or_init(AcceptedStoreCatalogScope::new)
3962                    .clone(),
3963                bundle.revision(),
3964                selection.root().fingerprint(),
3965            );
3966            let cardinality_domain = Rc::new(CardinalityAcceptedDomain::derive(&bundle)?);
3967            *self
3968                .accepted_bundle_cache
3969                .try_borrow_mut()
3970                .map_err(|_| InternalError::store_invariant())? = Some(AcceptedSchemaBundleCache {
3971                selection,
3972                bundle,
3973                cardinality_domain,
3974                value_catalog,
3975                entity_selections: RefCell::new(StdBTreeMap::new()),
3976            });
3977        }
3978
3979        let cache = self
3980            .accepted_bundle_cache
3981            .try_borrow()
3982            .map_err(|_| InternalError::store_invariant())?;
3983        let bundle = Ref::filter_map(cache, |cache| {
3984            cache
3985                .as_ref()
3986                .filter(|cached| cached.selection == selection)
3987                .map(|cached| &cached.bundle)
3988        })
3989        .map_err(|_| InternalError::store_invariant())?;
3990        self.validate_identity_state_closure(&bundle)?;
3991        Ok(Some((selection, bundle)))
3992    }
3993
3994    /// Reuse the accepted-domain projection for one already-selected effective root.
3995    pub(in crate::db) fn accepted_cardinality_domain_for_selection(
3996        &self,
3997        selection: Option<AcceptedSchemaRootSelection>,
3998    ) -> Result<Option<(AcceptedSchemaRootSelection, Rc<CardinalityAcceptedDomain>)>, InternalError>
3999    {
4000        let Some(selection) = selection else {
4001            self.accepted_bundle_cache
4002                .try_borrow_mut()
4003                .map_err(|_| InternalError::store_invariant())?
4004                .take();
4005            return Ok(None);
4006        };
4007        let cache_matches = self
4008            .accepted_bundle_cache
4009            .try_borrow()
4010            .map_err(|_| InternalError::store_invariant())?
4011            .as_ref()
4012            .is_some_and(|cached| cached.selection == selection);
4013        if !cache_matches {
4014            let authority = self
4015                .accepted_schema_authority_ref_for_selection(Some(selection))?
4016                .ok_or_else(InternalError::store_invariant)?;
4017            drop(authority);
4018        }
4019        let cache = self
4020            .accepted_bundle_cache
4021            .try_borrow()
4022            .map_err(|_| InternalError::store_invariant())?;
4023        let domain = cache
4024            .as_ref()
4025            .filter(|cached| cached.selection == selection)
4026            .map(|cached| Rc::clone(&cached.cardinality_domain))
4027            .ok_or_else(InternalError::store_invariant)?;
4028        Ok(Some((selection, domain)))
4029    }
4030
4031    /// Borrow the cached accepted-domain projection for an already-admitted root.
4032    pub(in crate::db) fn cached_cardinality_domain_for_root(
4033        &self,
4034        root: CardinalityAcceptedRootIdentity,
4035    ) -> Result<Option<Rc<CardinalityAcceptedDomain>>, InternalError> {
4036        let cache = self
4037            .accepted_bundle_cache
4038            .try_borrow()
4039            .map_err(|_| InternalError::store_invariant())?;
4040        Ok(cache
4041            .as_ref()
4042            .filter(|cached| root.matches(cached.selection.root()))
4043            .map(|cached| Rc::clone(&cached.cardinality_domain)))
4044    }
4045
4046    fn latest_raw_snapshots_by_entity(
4047        &self,
4048    ) -> StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)> {
4049        let mut latest_by_entity =
4050            StdBTreeMap::<EntityTag, (SchemaVersion, RawSchemaSnapshot)>::new();
4051
4052        let _: Result<(), std::convert::Infallible> = self.visit_raw_snapshots(|key, snapshot| {
4053            let version = SchemaVersion::new(key.version());
4054            match latest_by_entity.get_mut(&key.entity_tag()) {
4055                Some((latest_version, latest_snapshot)) if version > *latest_version => {
4056                    *latest_version = version;
4057                    *latest_snapshot = snapshot.clone();
4058                }
4059                None => {
4060                    latest_by_entity.insert(key.entity_tag(), (version, snapshot.clone()));
4061                }
4062                Some(_) => {}
4063            }
4064            Ok(SchemaStoreVisit::Continue)
4065        });
4066
4067        latest_by_entity
4068    }
4069
4070    /// Visit raw schema snapshots in canonical store order without exposing
4071    /// the backing stable-map iterator.
4072    fn visit_raw_snapshots<E>(
4073        &self,
4074        visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4075    ) -> Result<(), E> {
4076        let bounds = RawSchemaKey::all_entity_range_bounds();
4077        match &self.backend {
4078            SchemaStoreBackend::Heap(map) => {
4079                let mut visitor = visitor;
4080                for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4081                    if visitor(key, snapshot)?.should_stop() {
4082                        break;
4083                    }
4084                }
4085            }
4086            SchemaStoreBackend::Journaled {
4087                canonical,
4088                live,
4089                tombstones,
4090                ..
4091            } => Self::visit_journaled_raw_snapshot_range(
4092                canonical,
4093                live,
4094                tombstones,
4095                bounds,
4096                Direction::Asc,
4097                visitor,
4098            )?,
4099        }
4100
4101        Ok(())
4102    }
4103
4104    fn visit_constraint_validation_jobs_in_view<E>(
4105        &self,
4106        view: IdentityStateStorageView,
4107        visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4108    ) -> Result<(), E> {
4109        let bounds = RawSchemaKey::all_constraint_validation_job_range_bounds();
4110        match (&self.backend, view) {
4111            (SchemaStoreBackend::Heap(map), _) => {
4112                let mut visitor = visitor;
4113                for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4114                    if visitor(key, snapshot)?.should_stop() {
4115                        break;
4116                    }
4117                }
4118            }
4119            (
4120                SchemaStoreBackend::Journaled {
4121                    canonical,
4122                    live,
4123                    tombstones,
4124                    ..
4125                },
4126                IdentityStateStorageView::Effective,
4127            ) => Self::visit_journaled_raw_snapshot_range(
4128                canonical,
4129                live,
4130                tombstones,
4131                bounds,
4132                Direction::Asc,
4133                visitor,
4134            )?,
4135            (
4136                SchemaStoreBackend::Journaled { canonical, .. },
4137                IdentityStateStorageView::Canonical,
4138            ) => {
4139                let mut visitor = visitor;
4140                for entry in canonical.range((bounds.0, bounds.1)) {
4141                    if visitor(entry.key(), &entry.value())?.should_stop() {
4142                        break;
4143                    }
4144                }
4145            }
4146        }
4147        Ok(())
4148    }
4149
4150    #[cfg(test)]
4151    #[must_use]
4152    pub(in crate::db) fn canonical_len_for_tests(&self) -> u64 {
4153        match &self.backend {
4154            SchemaStoreBackend::Journaled { canonical: map, .. } => map.len(),
4155            SchemaStoreBackend::Heap(_) => 0,
4156        }
4157    }
4158
4159    fn get_raw_snapshot_for_backend(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
4160        let SchemaStoreBackend::Journaled {
4161            canonical,
4162            live,
4163            tombstones,
4164            ..
4165        } = &self.backend
4166        else {
4167            return None;
4168        };
4169
4170        if tombstones.contains(key) {
4171            return None;
4172        }
4173        live.get(key).cloned().or_else(|| canonical.get(key))
4174    }
4175
4176    fn visit_journaled_raw_snapshot_range<E>(
4177        canonical: &StableBTreeMap<
4178            RawSchemaKey,
4179            RawSchemaSnapshot,
4180            RuntimeMemory<DefaultMemoryImpl>,
4181        >,
4182        live: &StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
4183        tombstones: &BTreeSet<RawSchemaKey>,
4184        bounds: (RangeBound<RawSchemaKey>, RangeBound<RawSchemaKey>),
4185        direction: Direction,
4186        mut visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4187    ) -> Result<(), E> {
4188        match direction {
4189            Direction::Asc => {
4190                for entry in ordered_overlay_entries(
4191                    canonical.range((bounds.0, bounds.1)),
4192                    live.range((bounds.0, bounds.1)),
4193                    Direction::Asc,
4194                    |entry| entry.key(),
4195                    |entry| entry.0,
4196                    tombstones,
4197                ) {
4198                    let visit = match entry {
4199                        OrderedOverlayEntry::Canonical(canonical_entry) => {
4200                            visitor(canonical_entry.key(), &canonical_entry.value())?
4201                        }
4202                        OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4203                    };
4204                    if visit.should_stop() {
4205                        return Ok(());
4206                    }
4207                }
4208            }
4209            Direction::Desc => {
4210                for entry in ordered_overlay_entries(
4211                    canonical.range((bounds.0, bounds.1)).rev(),
4212                    live.range((bounds.0, bounds.1)).rev(),
4213                    Direction::Desc,
4214                    |entry| entry.key(),
4215                    |entry| entry.0,
4216                    tombstones,
4217                ) {
4218                    let visit = match entry {
4219                        OrderedOverlayEntry::Canonical(canonical_entry) => {
4220                            visitor(canonical_entry.key(), &canonical_entry.value())?
4221                        }
4222                        OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4223                    };
4224                    if visit.should_stop() {
4225                        return Ok(());
4226                    }
4227                }
4228            }
4229        }
4230
4231        Ok(())
4232    }
4233}
4234
4235fn map_schema_publication_error(error: AcceptedSchemaPublicationError) -> InternalError {
4236    match error {
4237        AcceptedSchemaPublicationError::StaleSchemaRevision { .. }
4238        | AcceptedSchemaPublicationError::RevisionExhausted => InternalError::store_unsupported(),
4239        AcceptedSchemaPublicationError::InvalidCandidate => InternalError::store_invariant(),
4240        AcceptedSchemaPublicationError::CorruptRootSlots => InternalError::store_corruption(),
4241    }
4242}
4243
4244fn derive_data_allocation_metadata(
4245    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4246) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4247    let mut max_version = SchemaVersion::initial();
4248    let mut hasher = new_hash_sha256();
4249    write_hash_tag_u8(&mut hasher, SCHEMA_STORE_DATA_ALLOCATION_FINGERPRINT_DOMAIN);
4250
4251    for (entity, (_, snapshot)) in latest_by_entity {
4252        let persisted = snapshot.decode_persisted_snapshot()?;
4253        if persisted.version() > max_version {
4254            max_version = persisted.version();
4255        }
4256
4257        let data_projection = PersistedSchemaSnapshot::new_with_primary_key_fields_and_indexes(
4258            persisted.version(),
4259            persisted.entity_path().to_string(),
4260            persisted.entity_name().to_string(),
4261            persisted.primary_key_field_ids().to_vec(),
4262            persisted.row_layout().clone(),
4263            persisted.fields().to_vec(),
4264            Vec::new(),
4265        );
4266        let constraint_catalog = crate::db::schema::AcceptedConstraintCatalog::initial(
4267            data_projection.fields(),
4268            data_projection.indexes(),
4269            data_projection.relations(),
4270        )
4271        .map_err(|_| InternalError::store_invariant())?;
4272        let data_projection = data_projection.with_constraint_catalog(constraint_catalog);
4273        let encoded = encode_persisted_schema_snapshot(&data_projection)?;
4274
4275        write_hash_u64(&mut hasher, entity.value());
4276        write_hash_u32(&mut hasher, persisted.version().get());
4277        write_hash_len_u32(&mut hasher, encoded.len());
4278        hasher.update(encoded);
4279    }
4280
4281    Ok(finalize_schema_metadata(
4282        max_version,
4283        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4284        hasher,
4285        latest_by_entity.len(),
4286    ))
4287}
4288
4289fn derive_index_allocation_metadata(
4290    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4291) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4292    let mut max_version = SchemaVersion::initial();
4293    let mut hasher = new_hash_sha256();
4294    write_hash_tag_u8(
4295        &mut hasher,
4296        SCHEMA_STORE_INDEX_ALLOCATION_FINGERPRINT_DOMAIN,
4297    );
4298
4299    for (entity, (_, snapshot)) in latest_by_entity {
4300        let persisted = snapshot.decode_persisted_snapshot()?;
4301        if persisted.version() > max_version {
4302            max_version = persisted.version();
4303        }
4304
4305        write_hash_u64(&mut hasher, entity.value());
4306        write_hash_u32(&mut hasher, persisted.version().get());
4307        write_hash_len_u32(&mut hasher, persisted.indexes().len());
4308        for index in persisted.indexes() {
4309            write_hash_u32(&mut hasher, u32::from(index.ordinal()));
4310            write_hash_str_u32(&mut hasher, index.name());
4311            write_hash_str_u32(&mut hasher, index.store());
4312            write_hash_tag_u8(&mut hasher, u8::from(index.unique()));
4313            write_hash_str_u32(&mut hasher, persisted_index_origin_name(index.origin()));
4314            match index.predicate_sql() {
4315                Some(predicate_sql) => {
4316                    write_hash_tag_u8(&mut hasher, 1);
4317                    write_hash_str_u32(&mut hasher, predicate_sql);
4318                }
4319                None => write_hash_tag_u8(&mut hasher, 0),
4320            }
4321            hash_persisted_index_key(&mut hasher, index.key());
4322        }
4323    }
4324
4325    Ok(finalize_schema_metadata(
4326        max_version,
4327        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4328        hasher,
4329        latest_by_entity.len(),
4330    ))
4331}
4332
4333fn derive_schema_catalog_metadata(
4334    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4335) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4336    let mut max_version = SchemaVersion::initial();
4337    let mut hasher = new_hash_sha256();
4338    write_hash_tag_u8(&mut hasher, SCHEMA_STORE_CATALOG_FINGERPRINT_DOMAIN);
4339
4340    for (entity, (version, snapshot)) in latest_by_entity {
4341        let persisted = snapshot.decode_persisted_snapshot()?;
4342        if persisted.version() > max_version {
4343            max_version = persisted.version();
4344        }
4345
4346        write_hash_u64(&mut hasher, entity.value());
4347        write_hash_u32(&mut hasher, version.get());
4348        write_hash_len_u32(&mut hasher, snapshot.as_bytes().len());
4349        hasher.update(snapshot.as_bytes());
4350    }
4351
4352    Ok(finalize_schema_metadata(
4353        max_version,
4354        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4355        hasher,
4356        latest_by_entity.len(),
4357    ))
4358}
4359
4360fn finalize_schema_metadata(
4361    schema_version: SchemaVersion,
4362    schema_fingerprint_method_version: u8,
4363    hasher: sha2::Sha256,
4364    entity_count: usize,
4365) -> SchemaStoreCatalogMetadata {
4366    let digest = finalize_hash_sha256(hasher);
4367    let mut schema_fingerprint = [0u8; 16];
4368    schema_fingerprint.copy_from_slice(&digest[..16]);
4369
4370    SchemaStoreCatalogMetadata::new(
4371        schema_version,
4372        schema_fingerprint_method_version,
4373        schema_fingerprint,
4374        u64::try_from(entity_count).unwrap_or(u64::MAX),
4375    )
4376}
4377
4378fn hash_persisted_index_key(hasher: &mut sha2::Sha256, key: &PersistedIndexKeySnapshot) {
4379    match key {
4380        PersistedIndexKeySnapshot::FieldPath(paths) => {
4381            write_hash_tag_u8(hasher, 1);
4382            write_hash_len_u32(hasher, paths.len());
4383            for path in paths {
4384                hash_persisted_index_field_path(hasher, path);
4385            }
4386        }
4387        PersistedIndexKeySnapshot::Items(items) => {
4388            write_hash_tag_u8(hasher, 2);
4389            write_hash_len_u32(hasher, items.len());
4390            for item in items {
4391                match item {
4392                    PersistedIndexKeyItemSnapshot::FieldPath(path) => {
4393                        write_hash_tag_u8(hasher, 1);
4394                        hash_persisted_index_field_path(hasher, path);
4395                    }
4396                    PersistedIndexKeyItemSnapshot::Expression(expression) => {
4397                        write_hash_tag_u8(hasher, 2);
4398                        write_hash_str_u32(hasher, persisted_expression_op_name(expression.op()));
4399                        hash_persisted_index_field_path(hasher, expression.source());
4400                        hash_accepted_field_kind(hasher, expression.input_kind());
4401                        hash_accepted_field_kind(hasher, expression.output_kind());
4402                        write_hash_str_u32(hasher, expression.canonical_text());
4403                    }
4404                }
4405            }
4406        }
4407    }
4408}
4409
4410fn hash_persisted_index_field_path(
4411    hasher: &mut sha2::Sha256,
4412    path: &crate::db::schema::PersistedIndexFieldPathSnapshot,
4413) {
4414    write_hash_u32(hasher, path.field_id().get());
4415    write_hash_u32(hasher, u32::from(path.slot().get()));
4416    write_hash_len_u32(hasher, path.path().len());
4417    for segment in path.path() {
4418        write_hash_str_u32(hasher, segment);
4419    }
4420    hash_accepted_field_kind(hasher, path.kind());
4421    write_hash_tag_u8(hasher, u8::from(path.nullable()));
4422}
4423
4424fn hash_accepted_field_kind(hasher: &mut sha2::Sha256, kind: &AcceptedFieldKind) {
4425    match kind {
4426        AcceptedFieldKind::Account => write_hash_tag_u8(hasher, 1),
4427        AcceptedFieldKind::Blob { max_len } => {
4428            write_hash_tag_u8(hasher, 2);
4429            hash_optional_u32(hasher, *max_len);
4430        }
4431        AcceptedFieldKind::Bool => {
4432            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_BOOL);
4433        }
4434        AcceptedFieldKind::Date => write_hash_tag_u8(hasher, 4),
4435        AcceptedFieldKind::Decimal { scale } => {
4436            write_hash_tag_u8(hasher, 5);
4437            write_hash_u32(hasher, *scale);
4438        }
4439        AcceptedFieldKind::Duration => write_hash_tag_u8(hasher, 6),
4440        AcceptedFieldKind::Enum { type_id } => {
4441            write_hash_tag_u8(hasher, 7);
4442            write_hash_u32(hasher, type_id.get());
4443        }
4444        AcceptedFieldKind::Float32 => write_hash_tag_u8(hasher, 8),
4445        AcceptedFieldKind::Float64 => write_hash_tag_u8(hasher, 9),
4446        AcceptedFieldKind::Int8 => write_hash_tag_u8(hasher, 10),
4447        AcceptedFieldKind::Int16 => write_hash_tag_u8(hasher, 11),
4448        AcceptedFieldKind::Int32 => write_hash_tag_u8(hasher, 12),
4449        AcceptedFieldKind::Int64 => write_hash_tag_u8(hasher, 13),
4450        AcceptedFieldKind::Int128 => write_hash_tag_u8(hasher, 14),
4451        AcceptedFieldKind::IntBig { max_bytes } => {
4452            write_hash_tag_u8(hasher, 15);
4453            write_hash_u32(hasher, *max_bytes);
4454        }
4455        AcceptedFieldKind::Principal => write_hash_tag_u8(hasher, 16),
4456        AcceptedFieldKind::Subaccount => write_hash_tag_u8(hasher, 17),
4457        AcceptedFieldKind::Text { max_len } => {
4458            write_hash_tag_u8(hasher, 18);
4459            hash_optional_u32(hasher, *max_len);
4460        }
4461        AcceptedFieldKind::Timestamp => write_hash_tag_u8(hasher, 19),
4462        AcceptedFieldKind::Nat8 => write_hash_tag_u8(hasher, 20),
4463        AcceptedFieldKind::Nat16 => write_hash_tag_u8(hasher, 21),
4464        AcceptedFieldKind::Nat32 => write_hash_tag_u8(hasher, 22),
4465        AcceptedFieldKind::Nat64 => write_hash_tag_u8(hasher, 23),
4466        AcceptedFieldKind::Nat128 => write_hash_tag_u8(hasher, 24),
4467        AcceptedFieldKind::NatBig { max_bytes } => {
4468            write_hash_tag_u8(hasher, 25);
4469            write_hash_u32(hasher, *max_bytes);
4470        }
4471        AcceptedFieldKind::Ulid => write_hash_tag_u8(hasher, 26),
4472        AcceptedFieldKind::Unit => write_hash_tag_u8(hasher, 27),
4473        AcceptedFieldKind::Relation {
4474            target_path,
4475            target_entity_name,
4476            target_entity_tag,
4477            target_store_path,
4478            key_kind,
4479        } => {
4480            write_hash_tag_u8(hasher, 28);
4481            write_hash_str_u32(hasher, target_path);
4482            write_hash_str_u32(hasher, target_entity_name);
4483            write_hash_u64(hasher, target_entity_tag.value());
4484            write_hash_str_u32(hasher, target_store_path);
4485            hash_accepted_field_kind(hasher, key_kind);
4486        }
4487        AcceptedFieldKind::List(inner) => {
4488            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_LIST);
4489            hash_accepted_field_kind(hasher, inner);
4490        }
4491        AcceptedFieldKind::Set(inner) => {
4492            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_SET);
4493            hash_accepted_field_kind(hasher, inner);
4494        }
4495        AcceptedFieldKind::Map { key, value } => {
4496            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_MAP);
4497            hash_accepted_field_kind(hasher, key);
4498            hash_accepted_field_kind(hasher, value);
4499        }
4500        AcceptedFieldKind::Composite { type_id } => {
4501            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_COMPOSITE);
4502            write_hash_u32(hasher, type_id.get());
4503        }
4504        AcceptedFieldKind::U256 => write_hash_tag_u8(hasher, 33),
4505    }
4506}
4507
4508fn hash_optional_u32(hasher: &mut sha2::Sha256, value: Option<u32>) {
4509    match value {
4510        Some(value) => {
4511            write_hash_tag_u8(hasher, 1);
4512            write_hash_u32(hasher, value);
4513        }
4514        None => write_hash_tag_u8(hasher, 0),
4515    }
4516}
4517
4518const fn persisted_index_origin_name(
4519    origin: crate::db::schema::PersistedIndexOrigin,
4520) -> &'static str {
4521    match origin {
4522        crate::db::schema::PersistedIndexOrigin::Generated => "generated",
4523        crate::db::schema::PersistedIndexOrigin::SqlDdl => "sql_ddl",
4524    }
4525}
4526
4527const fn persisted_expression_op_name(
4528    op: crate::db::schema::PersistedIndexExpressionOp,
4529) -> &'static str {
4530    match op {
4531        crate::db::schema::PersistedIndexExpressionOp::Lower => "lower",
4532        crate::db::schema::PersistedIndexExpressionOp::Upper => "upper",
4533        crate::db::schema::PersistedIndexExpressionOp::Trim => "trim",
4534        crate::db::schema::PersistedIndexExpressionOp::LowerTrim => "lower_trim",
4535        crate::db::schema::PersistedIndexExpressionOp::Date => "date",
4536        crate::db::schema::PersistedIndexExpressionOp::Year => "year",
4537        crate::db::schema::PersistedIndexExpressionOp::Month => "month",
4538        crate::db::schema::PersistedIndexExpressionOp::Day => "day",
4539    }
4540}
4541
4542///
4543/// TESTS
4544///
4545
4546#[cfg(test)]
4547mod tests;