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        snapshot.version(),
634        snapshot.primary_key_field_ids(),
635        snapshot.row_layout(),
636        snapshot.fields(),
637    )
638    .is_some()
639    {
640        return Err(InternalError::store_invariant());
641    }
642
643    Ok(())
644}
645
646///
647/// SchemaStoreCatalogMetadata
648///
649/// Accepted schema-store catalog metadata derived from latest persisted
650/// snapshots. This is diagnostic allocation metadata, not allocation identity.
651///
652
653#[derive(Clone, Copy, Debug, Eq, PartialEq)]
654pub(in crate::db) struct SchemaStoreCatalogMetadata {
655    schema_version: SchemaVersion,
656    schema_fingerprint_method_version: u8,
657    schema_fingerprint: CommitSchemaFingerprint,
658    entity_count: u64,
659}
660
661impl SchemaStoreCatalogMetadata {
662    /// Build catalog metadata from already-derived accepted schema facts.
663    #[must_use]
664    const fn new(
665        schema_version: SchemaVersion,
666        schema_fingerprint_method_version: u8,
667        schema_fingerprint: CommitSchemaFingerprint,
668        entity_count: u64,
669    ) -> Self {
670        Self {
671            schema_version,
672            schema_fingerprint_method_version,
673            schema_fingerprint,
674            entity_count,
675        }
676    }
677
678    /// Return the maximum latest schema version represented in the catalog.
679    #[must_use]
680    pub(in crate::db) const fn schema_version(self) -> SchemaVersion {
681        self.schema_version
682    }
683
684    /// Return the fingerprint method version for this diagnostic metadata row.
685    #[must_use]
686    pub(in crate::db) const fn schema_fingerprint_method_version(self) -> u8 {
687        self.schema_fingerprint_method_version
688    }
689
690    /// Return the deterministic catalog fingerprint for latest accepted
691    /// snapshots.
692    #[must_use]
693    pub(in crate::db) const fn schema_fingerprint(self) -> CommitSchemaFingerprint {
694        self.schema_fingerprint
695    }
696
697    /// Return number of entity schemas represented in this catalog metadata.
698    #[must_use]
699    pub(in crate::db) const fn entity_count(self) -> u64 {
700        self.entity_count
701    }
702}
703
704///
705/// SchemaStoreAllocationMetadata
706///
707/// Role-specific allocation metadata derived from latest accepted schema-store
708/// snapshots. These fingerprints describe the accepted contract that owns each
709/// allocation role; they are diagnostics, not allocation identity.
710///
711
712#[derive(Clone, Copy, Debug, Eq, PartialEq)]
713pub(in crate::db) struct SchemaStoreAllocationMetadata {
714    data: SchemaStoreCatalogMetadata,
715    index: SchemaStoreCatalogMetadata,
716    schema: SchemaStoreCatalogMetadata,
717}
718
719impl SchemaStoreAllocationMetadata {
720    /// Build one role-specific metadata set from already-derived accepted
721    /// schema facts.
722    #[must_use]
723    const fn new(
724        data: SchemaStoreCatalogMetadata,
725        index: SchemaStoreCatalogMetadata,
726        schema: SchemaStoreCatalogMetadata,
727    ) -> Self {
728        Self {
729            data,
730            index,
731            schema,
732        }
733    }
734
735    /// Return accepted row-layout allocation metadata for data memory.
736    #[must_use]
737    pub(in crate::db) const fn data(self) -> SchemaStoreCatalogMetadata {
738        self.data
739    }
740
741    /// Return accepted index-catalog allocation metadata for index memory.
742    #[must_use]
743    pub(in crate::db) const fn index(self) -> SchemaStoreCatalogMetadata {
744        self.index
745    }
746
747    /// Return accepted full schema-catalog allocation metadata for schema
748    /// memory.
749    #[must_use]
750    pub(in crate::db) const fn schema(self) -> SchemaStoreCatalogMetadata {
751        self.schema
752    }
753}
754
755///
756/// PendingRelationActivationDeleteBarrier
757///
758/// Accepted activation identity that blocks target deletion until a candidate
759/// reverse-relation generation is proven and promoted.
760///
761
762pub(in crate::db) struct PendingRelationActivationDeleteBarrier {
763    accepted_schema_fingerprint: CommitSchemaFingerprint,
764    source_entity_tag: EntityTag,
765    constraint_id: ConstraintId,
766}
767
768impl PendingRelationActivationDeleteBarrier {
769    #[must_use]
770    pub(in crate::db) const fn accepted_schema_fingerprint(&self) -> CommitSchemaFingerprint {
771        self.accepted_schema_fingerprint
772    }
773
774    #[must_use]
775    pub(in crate::db) const fn source_entity_tag(&self) -> EntityTag {
776        self.source_entity_tag
777    }
778
779    /// Return the stable accepted constraint identity.
780    #[must_use]
781    pub(in crate::db) const fn constraint_id(&self) -> ConstraintId {
782        self.constraint_id
783    }
784}
785
786///
787/// SchemaStore
788///
789/// Thin persistence wrapper over one journaled or heap schema metadata BTreeMap.
790/// Startup reconciliation writes and validates encoded schema snapshots here
791/// before row/index operations proceed.
792///
793
794pub struct SchemaStore {
795    backend: SchemaStoreBackend,
796    accepted_bundle_cache: RefCell<Option<AcceptedSchemaBundleCache>>,
797    cardinality_header_cache: RefCell<Option<(Vec<u8>, CardinalityGenerationHeader)>>,
798    accepted_catalog_scope: OnceCell<AcceptedStoreCatalogScope>,
799}
800
801struct AcceptedSchemaBundleCache {
802    selection: AcceptedSchemaRootSelection,
803    bundle: AcceptedSchemaRevisionBundle,
804    cardinality_domain: Rc<CardinalityAcceptedDomain>,
805    value_catalog: AcceptedValueCatalogHandle,
806    entity_selections: RefCell<StdBTreeMap<EntityTag, AcceptedCatalogSnapshotSelection>>,
807}
808
809enum SchemaStoreBackend {
810    Heap(StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>),
811    Journaled {
812        canonical:
813            StableBTreeMap<RawSchemaKey, RawSchemaSnapshot, RuntimeMemory<DefaultMemoryImpl>>,
814        live: StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
815        tombstones: BTreeSet<RawSchemaKey>,
816        positions: PositionedOverlayMetadata<RawSchemaKey>,
817    },
818}
819
820/// Control-flow result for schema-store traversal visitors.
821#[derive(Clone, Copy, Debug, Eq, PartialEq)]
822enum SchemaStoreVisit {
823    Continue,
824    #[cfg(test)]
825    Stop,
826}
827
828impl SchemaStoreVisit {
829    const fn should_stop(self) -> bool {
830        match self {
831            Self::Continue => false,
832            #[cfg(test)]
833            Self::Stop => true,
834        }
835    }
836}
837
838#[derive(Clone, Copy)]
839enum IdentityStateStorageView {
840    Effective,
841    Canonical,
842}
843
844/// Fully validated entity-snapshot bytes and identity, prepared before apply.
845/// Fields stay private so consumers cannot bypass the store's encoder.
846pub(in crate::db) struct PreparedSchemaSnapshot {
847    key: RawSchemaKey,
848    snapshot: RawSchemaSnapshot,
849}
850
851/// Candidate and canonical payloads retained from preflight until atomic fold.
852/// Not a reusable plan: preparation and application belong to one batch callback.
853pub(in crate::db) struct PreparedAcceptedSchemaFold {
854    candidate: CandidateSchemaRevision,
855    expected_revision: AcceptedSchemaRevision,
856    snapshots: Vec<PreparedSchemaSnapshot>,
857    identity_updates: Vec<(RawSchemaKey, Vec<u8>)>,
858    identity_removals: Vec<RawSchemaKey>,
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        let (updates, removals) = transition.into_effects();
1731        for owner in removals {
1732            self.remove_identity_state_key(
1733                RawSchemaKey::from_identity_state(owner.entity_tag(), owner.field_id()),
1734                target,
1735            )?;
1736        }
1737        for state in updates {
1738            let key = RawSchemaKey::from_identity_state(
1739                state.owner().entity_tag(),
1740                state.owner().field_id(),
1741            );
1742            let bytes = encode_identity_state(&state)?;
1743            match target {
1744                IdentityStateWriteTarget::Durable => {
1745                    self.insert_durable_raw_value(key, bytes);
1746                }
1747                IdentityStateWriteTarget::Materialized => {
1748                    self.insert_raw_snapshot(
1749                        key,
1750                        RawSchemaSnapshot::from_encoded_control_record(bytes),
1751                    );
1752                }
1753                IdentityStateWriteTarget::Canonical => {
1754                    self.insert_canonical_raw_value(key, bytes)?;
1755                }
1756            }
1757        }
1758        Ok(())
1759    }
1760
1761    // A retained allocator moves keys with its accepted field. These deletes
1762    // follow the same destination authority as schema publication; canonical
1763    // folding must not erase a newer positioned live state.
1764    fn remove_identity_state_key(
1765        &mut self,
1766        key: RawSchemaKey,
1767        target: IdentityStateWriteTarget,
1768    ) -> Result<(), InternalError> {
1769        match (&mut self.backend, target) {
1770            (
1771                SchemaStoreBackend::Heap(map),
1772                IdentityStateWriteTarget::Durable | IdentityStateWriteTarget::Materialized,
1773            ) => {
1774                map.remove(&key);
1775            }
1776            (
1777                SchemaStoreBackend::Journaled {
1778                    live, tombstones, ..
1779                },
1780                IdentityStateWriteTarget::Materialized,
1781            ) => {
1782                live.remove(&key);
1783                tombstones.insert(key);
1784            }
1785            (
1786                SchemaStoreBackend::Journaled {
1787                    canonical,
1788                    live,
1789                    tombstones,
1790                    ..
1791                },
1792                IdentityStateWriteTarget::Durable,
1793            ) => {
1794                live.remove(&key);
1795                tombstones.remove(&key);
1796                canonical.remove(&key);
1797            }
1798            (
1799                SchemaStoreBackend::Journaled { canonical, .. },
1800                IdentityStateWriteTarget::Canonical,
1801            ) => {
1802                canonical.remove(&key);
1803            }
1804            (SchemaStoreBackend::Heap(_), IdentityStateWriteTarget::Canonical) => {
1805                return Err(InternalError::store_invariant());
1806            }
1807        }
1808        Ok(())
1809    }
1810
1811    pub(in crate::db) fn current_canonical_accepted_schema_bundle(
1812        &self,
1813    ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
1814        self.current_canonical_accepted_schema_authority()
1815            .map(|authority| authority.map(|(_, bundle)| bundle))
1816    }
1817
1818    /// Load one canonical accepted root and its verified immutable bundle.
1819    pub(in crate::db) fn current_canonical_accepted_schema_authority(
1820        &self,
1821    ) -> Result<Option<(AcceptedSchemaRootSelection, AcceptedSchemaRevisionBundle)>, InternalError>
1822    {
1823        let Some(selection) = self.current_canonical_accepted_schema_root()? else {
1824            return Ok(None);
1825        };
1826        let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
1827        let raw = self
1828            .get_canonical_raw_value(&bundle_key)?
1829            .ok_or_else(InternalError::store_corruption)?;
1830        let bundle =
1831            decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
1832        Ok(Some((selection, bundle)))
1833    }
1834
1835    /// Return the accepted root selected only from canonical predecessor slots.
1836    pub(in crate::db) fn current_canonical_accepted_schema_root(
1837        &self,
1838    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
1839        let first = self.canonical_root_slot_bytes(0)?;
1840        let second = self.canonical_root_slot_bytes(1)?;
1841        select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
1842    }
1843
1844    /// Select effective and canonical roots from one canonical slot read.
1845    pub(in crate::db) fn current_effective_and_canonical_accepted_schema_roots(
1846        &self,
1847    ) -> Result<
1848        (
1849            Option<AcceptedSchemaRootSelection>,
1850            Option<AcceptedSchemaRootSelection>,
1851        ),
1852        InternalError,
1853    > {
1854        let SchemaStoreBackend::Journaled {
1855            canonical,
1856            live,
1857            tombstones,
1858            ..
1859        } = &self.backend
1860        else {
1861            return Err(InternalError::store_invariant());
1862        };
1863        let first_key = RawSchemaKey::from_accepted_root_slot(0)?;
1864        let second_key = RawSchemaKey::from_accepted_root_slot(1)?;
1865        let mut canonical_first = None;
1866        let mut canonical_second = None;
1867        for entry in canonical.range(first_key..=second_key) {
1868            if *entry.key() == first_key {
1869                canonical_first = Some(entry.value().clone());
1870            } else if *entry.key() == second_key {
1871                canonical_second = Some(entry.value().clone());
1872            } else {
1873                return Err(InternalError::store_corruption());
1874            }
1875        }
1876        let effective_first = if tombstones.contains(&first_key) {
1877            None
1878        } else {
1879            live.get(&first_key)
1880                .cloned()
1881                .or_else(|| canonical_first.clone())
1882        };
1883        let effective_second = if tombstones.contains(&second_key) {
1884            None
1885        } else {
1886            live.get(&second_key)
1887                .cloned()
1888                .or_else(|| canonical_second.clone())
1889        };
1890        let effective_first = effective_first.map(RawSchemaSnapshot::into_bytes);
1891        let effective_second = effective_second.map(RawSchemaSnapshot::into_bytes);
1892        let canonical_first = canonical_first.map(RawSchemaSnapshot::into_bytes);
1893        let canonical_second = canonical_second.map(RawSchemaSnapshot::into_bytes);
1894        Ok((
1895            select_current_accepted_schema_root([
1896                effective_first.as_deref(),
1897                effective_second.as_deref(),
1898            ])?,
1899            select_current_accepted_schema_root([
1900                canonical_first.as_deref(),
1901                canonical_second.as_deref(),
1902            ])?,
1903        ))
1904    }
1905
1906    /// Insert or replace one typed persisted schema snapshot.
1907    pub(in crate::db) fn insert_persisted_snapshot(
1908        &mut self,
1909        entity: EntityTag,
1910        snapshot: &PersistedSchemaSnapshot,
1911    ) -> Result<(), InternalError> {
1912        let prepared = Self::prepare_persisted_snapshot(entity, snapshot)?;
1913        self.apply_prepared_persisted_snapshot(prepared);
1914
1915        Ok(())
1916    }
1917
1918    /// Finish snapshot validation, fingerprinting and encoding before mutation.
1919    pub(in crate::db) fn prepare_persisted_snapshot(
1920        entity: EntityTag,
1921        snapshot: &PersistedSchemaSnapshot,
1922    ) -> Result<PreparedSchemaSnapshot, InternalError> {
1923        Ok(PreparedSchemaSnapshot {
1924            key: RawSchemaKey::from_entity_version(entity, snapshot.version()),
1925            snapshot: RawSchemaSnapshot::from_persisted_snapshot(snapshot)?,
1926        })
1927    }
1928
1929    /// Publish prepared bytes without repeating semantic construction.
1930    pub(in crate::db) fn apply_prepared_persisted_snapshot(
1931        &mut self,
1932        prepared: PreparedSchemaSnapshot,
1933    ) {
1934        let _ = self.insert_raw_snapshot(prepared.key, prepared.snapshot);
1935    }
1936
1937    /// Load one schema-owned constraint validation job.
1938    pub(in crate::db) fn constraint_validation_job(
1939        &self,
1940        entity: EntityTag,
1941        constraint_id: ConstraintId,
1942    ) -> Result<Option<ConstraintValidationJob>, InternalError> {
1943        let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1944        self.get_raw_snapshot(&key)
1945            .map(|raw| decode_constraint_validation_job(raw.as_bytes()))
1946            .transpose()
1947    }
1948
1949    /// Apply one marker-authorized validation job to the live schema projection.
1950    pub(in crate::db) fn apply_constraint_validation_job(
1951        &mut self,
1952        job: &ConstraintValidationJob,
1953    ) -> Result<(), InternalError> {
1954        let key =
1955            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1956        let bytes = encode_constraint_validation_job(job)?;
1957        let _ =
1958            self.insert_raw_snapshot(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1959        Ok(())
1960    }
1961
1962    /// Remove one marker-authorized validation job from the live projection.
1963    #[expect(
1964        clippy::unnecessary_wraps,
1965        reason = "marker apply operations share one fallible callback contract"
1966    )]
1967    pub(in crate::db) fn apply_constraint_validation_job_removal(
1968        &mut self,
1969        entity: EntityTag,
1970        constraint_id: ConstraintId,
1971    ) -> Result<(), InternalError> {
1972        let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1973        match &mut self.backend {
1974            SchemaStoreBackend::Heap(map) => {
1975                map.remove(&key);
1976            }
1977            SchemaStoreBackend::Journaled {
1978                live, tombstones, ..
1979            } => {
1980                live.remove(&key);
1981                tombstones.insert(key);
1982            }
1983        }
1984        Ok(())
1985    }
1986
1987    /// Fold one committed validation job into the canonical stable base.
1988    pub(in crate::db) fn fold_constraint_validation_job(
1989        &mut self,
1990        job: &ConstraintValidationJob,
1991    ) -> Result<(), InternalError> {
1992        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1993            return Err(InternalError::store_invariant());
1994        };
1995        let key =
1996            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1997        let bytes = encode_constraint_validation_job(job)?;
1998        canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1999        Ok(())
2000    }
2001
2002    /// Preflight one canonical validation-job fold without changing storage.
2003    pub(in crate::db) fn preflight_fold_constraint_validation_job(
2004        &self,
2005        job: &ConstraintValidationJob,
2006    ) -> Result<(), InternalError> {
2007        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2008            return Err(InternalError::store_invariant());
2009        }
2010        let _encoded = encode_constraint_validation_job(job)?;
2011        Ok(())
2012    }
2013
2014    /// Fold one committed validation-job removal into the canonical stable base.
2015    pub(in crate::db) fn fold_constraint_validation_job_removal(
2016        &mut self,
2017        entity: EntityTag,
2018        constraint_id: ConstraintId,
2019    ) -> Result<(), InternalError> {
2020        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
2021            return Err(InternalError::store_invariant());
2022        };
2023        canonical.remove(&RawSchemaKey::from_constraint_validation_job(
2024            entity,
2025            constraint_id,
2026        ));
2027        Ok(())
2028    }
2029
2030    /// Preflight one canonical validation-job removal without changing storage.
2031    pub(in crate::db) fn preflight_fold_constraint_validation_job_removal(
2032        &self,
2033    ) -> Result<(), InternalError> {
2034        match self.backend {
2035            SchemaStoreBackend::Journaled { .. } => Ok(()),
2036            SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
2037        }
2038    }
2039
2040    /// Reset the volatile projection for journaled recovery without mutating
2041    /// the canonical stable schema base.
2042    pub(in crate::db) fn reset_journaled_live_projection(&mut self) -> Result<(), InternalError> {
2043        let SchemaStoreBackend::Journaled {
2044            live,
2045            tombstones,
2046            positions,
2047            ..
2048        } = &mut self.backend
2049        else {
2050            return Err(InternalError::store_invariant());
2051        };
2052
2053        live.clear();
2054        tombstones.clear();
2055        positions.clear();
2056        self.accepted_bundle_cache.get_mut().take();
2057
2058        Ok(())
2059    }
2060
2061    /// Preflight every schema/control position represented by one online batch.
2062    pub(in crate::db) fn prepare_positioned_journal_batch_publication(
2063        &self,
2064        incarnation: DatabaseIncarnationId,
2065        batch: &JournalBatch,
2066        position: JournalOverlayPosition,
2067    ) -> Result<PreparedSchemaPositionPublication, InternalError> {
2068        let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2069            return Err(InternalError::store_invariant());
2070        };
2071        let keys = self.positioned_journal_batch_keys(
2072            incarnation,
2073            batch,
2074            IdentityStateStorageView::Effective,
2075        )?;
2076        for key in &keys {
2077            positions.preflight_publish(key, position)?;
2078        }
2079        Ok(PreparedSchemaPositionPublication {
2080            keys: keys.into_iter().collect(),
2081            position,
2082        })
2083    }
2084
2085    /// Preflight exact schema/control retirement before canonical mutation.
2086    pub(in crate::db) fn prepare_positioned_journal_batch_retirement(
2087        &self,
2088        incarnation: DatabaseIncarnationId,
2089        batch: &JournalBatch,
2090        position: JournalOverlayPosition,
2091    ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2092        let keys = self.positioned_journal_batch_keys(
2093            incarnation,
2094            batch,
2095            IdentityStateStorageView::Canonical,
2096        )?;
2097        self.prepare_positioned_key_retirements(keys, position)
2098    }
2099
2100    fn prepare_positioned_key_retirements(
2101        &self,
2102        keys: impl IntoIterator<Item = RawSchemaKey>,
2103        position: JournalOverlayPosition,
2104    ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2105        let SchemaStoreBackend::Journaled {
2106            live,
2107            tombstones,
2108            positions,
2109            ..
2110        } = &self.backend
2111        else {
2112            return Err(InternalError::store_invariant());
2113        };
2114        let mut entries = Vec::new();
2115        for key in keys {
2116            if !positions.is_positioned(&key) {
2117                // A prior row fold can create canonical-only derived metadata
2118                // after this older schema batch publishes. Its canonical fold
2119                // owns that key; there is no overlay for this batch to retire.
2120                if live.contains_key(&key) || tombstones.contains(&key) {
2121                    return Err(InternalError::store_invariant());
2122                }
2123                continue;
2124            }
2125            let retirement = positions.preflight_retirement(&key, position)?;
2126            entries.push((key, retirement));
2127        }
2128        Ok(PreparedSchemaPositionRetirement { entries })
2129    }
2130
2131    /// Publish schema positions after their values have been mechanically applied.
2132    pub(in crate::db) fn publish_prepared_journal_batch_positions(
2133        &mut self,
2134        prepared: PreparedSchemaPositionPublication,
2135    ) {
2136        let SchemaStoreBackend::Journaled { positions, .. } = &mut self.backend else {
2137            debug_assert!(
2138                false,
2139                "preflighted schema positions require a journaled store"
2140            );
2141            return;
2142        };
2143        for key in prepared.keys {
2144            positions.publish_preflighted(key, prepared.position);
2145        }
2146    }
2147
2148    /// Retire only exact schema/control overlays after canonical mutation.
2149    pub(in crate::db) fn apply_prepared_journal_batch_retirement(
2150        &mut self,
2151        prepared: PreparedSchemaPositionRetirement,
2152    ) {
2153        for (key, retirement) in prepared.entries {
2154            if retirement != PositionedOverlayRetirement::Exact {
2155                continue;
2156            }
2157            self.invalidate_accepted_bundle_cache_for_key(key);
2158            let SchemaStoreBackend::Journaled {
2159                live,
2160                tombstones,
2161                positions,
2162                ..
2163            } = &mut self.backend
2164            else {
2165                debug_assert!(
2166                    false,
2167                    "preflighted schema retirement requires a journaled store"
2168                );
2169                return;
2170            };
2171            live.remove(&key);
2172            tombstones.remove(&key);
2173            positions.retire_preflighted(&key, retirement);
2174        }
2175    }
2176
2177    #[cfg(test)]
2178    fn publish_positioned_journal_entry(
2179        &mut self,
2180        key: RawSchemaKey,
2181        snapshot: Option<RawSchemaSnapshot>,
2182        position: JournalOverlayPosition,
2183    ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
2184        let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2185            return Err(InternalError::store_invariant());
2186        };
2187        positions.preflight_publish(&key, position)?;
2188        self.invalidate_accepted_bundle_cache_for_key(key);
2189        let SchemaStoreBackend::Journaled {
2190            canonical,
2191            live,
2192            tombstones,
2193            positions,
2194        } = &mut self.backend
2195        else {
2196            return Err(InternalError::store_invariant());
2197        };
2198        let previous = if tombstones.contains(&key) {
2199            None
2200        } else {
2201            live.get(&key).cloned().or_else(|| canonical.get(&key))
2202        };
2203        if let Some(snapshot) = snapshot {
2204            tombstones.remove(&key);
2205            live.insert(key, snapshot);
2206        } else {
2207            live.remove(&key);
2208            tombstones.insert(key);
2209        }
2210        positions.publish_preflighted(key, position);
2211        Ok(previous)
2212    }
2213
2214    /// Seed a test's canonical snapshot through the maintained prepared handoff.
2215    #[cfg(test)]
2216    pub(in crate::db) fn fold_persisted_snapshot(
2217        &mut self,
2218        entity: EntityTag,
2219        snapshot: &PersistedSchemaSnapshot,
2220    ) -> Result<(), InternalError> {
2221        let prepared = self.prepare_fold_persisted_snapshot(entity, snapshot)?;
2222        self.apply_prepared_fold_persisted_snapshot(prepared)
2223    }
2224
2225    /// Apply only the encoded snapshot prepared for this journaled store.
2226    pub(in crate::db) fn apply_prepared_fold_persisted_snapshot(
2227        &mut self,
2228        prepared: PreparedSchemaSnapshot,
2229    ) -> Result<(), InternalError> {
2230        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
2231            return Err(InternalError::store_invariant());
2232        };
2233        canonical.insert(prepared.key, prepared.snapshot);
2234
2235        Ok(())
2236    }
2237
2238    /// Prepare one canonical fold, retaining encoded bytes without changing storage.
2239    pub(in crate::db) fn prepare_fold_persisted_snapshot(
2240        &self,
2241        entity: EntityTag,
2242        snapshot: &PersistedSchemaSnapshot,
2243    ) -> Result<PreparedSchemaSnapshot, InternalError> {
2244        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2245            return Err(InternalError::store_invariant());
2246        }
2247        Self::prepare_persisted_snapshot(entity, snapshot)
2248    }
2249
2250    /// Return the current accepted store root selected from its two checksummed slots.
2251    pub(in crate::db) fn current_accepted_schema_root(
2252        &self,
2253    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
2254        let first = self.accepted_root_slot_bytes(0)?;
2255        let second = self.accepted_root_slot_bytes(1)?;
2256        select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
2257    }
2258
2259    /// Load and verify the immutable bundle referenced by the current root.
2260    pub(in crate::db) fn current_accepted_schema_bundle(
2261        &self,
2262    ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
2263        self.borrow_current_accepted_schema_bundle()
2264            .map(|bundle| bundle.map(|bundle| bundle.clone()))
2265    }
2266
2267    /// Borrow the current verified bundle, rechecking durable job closure even
2268    /// on a cache hit. The borrow must end before mutating schema authority.
2269    pub(in crate::db) fn borrow_current_accepted_schema_bundle(
2270        &self,
2271    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
2272        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2273            return Ok(None);
2274        };
2275        self.validate_constraint_validation_job_closure(&bundle)?;
2276
2277        Ok(Some(bundle))
2278    }
2279
2280    /// Borrow the verified bundle selected by a current catalog operation.
2281    /// Root publication invalidates this store-owned cache before replacement;
2282    /// exact authority equality also binds store scope, revision and fingerprint.
2283    /// Mutable identity and validation-job records still require live checks.
2284    pub(in crate::db) fn borrow_accepted_schema_bundle_for_authority(
2285        &self,
2286        expected: &AcceptedSchemaAuthority,
2287    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
2288        let Some(selection) = self.current_accepted_selection_for_authority(expected)? else {
2289            return Ok(None);
2290        };
2291        let Some((_, bundle)) =
2292            self.accepted_schema_authority_ref_for_selection(Some(selection))?
2293        else {
2294            return Ok(None);
2295        };
2296        self.validate_constraint_validation_job_closure(&bundle)?;
2297
2298        Ok(Some(bundle))
2299    }
2300
2301    /// Project current accepted entity identity onto one registry-owned store path.
2302    pub(in crate::db) fn current_accepted_runtime_entities(
2303        &self,
2304        registered_store_path: &'static str,
2305    ) -> Result<Vec<AcceptedRuntimeEntity>, InternalError> {
2306        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2307            return Ok(Vec::new());
2308        };
2309        if bundle.store_path() != registered_store_path {
2310            return Err(InternalError::store_corruption());
2311        }
2312
2313        bundle
2314            .entity_snapshots()
2315            .iter()
2316            .map(|(entity_tag, snapshot)| {
2317                AcceptedRuntimeEntity::from_accepted_snapshot(
2318                    &bundle,
2319                    *entity_tag,
2320                    snapshot,
2321                    registered_store_path,
2322                )
2323            })
2324            .collect()
2325    }
2326
2327    /// Resolve one accepted entity tag without materializing the full store catalog.
2328    pub(in crate::db) fn current_accepted_runtime_entity_for_tag(
2329        &self,
2330        registered_store_path: &'static str,
2331        entity_tag: EntityTag,
2332    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2333        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2334            return Ok(None);
2335        };
2336        if bundle.store_path() != registered_store_path {
2337            return Err(InternalError::store_corruption());
2338        }
2339        let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2340            return Ok(None);
2341        };
2342
2343        AcceptedRuntimeEntity::from_accepted_snapshot(
2344            &bundle,
2345            entity_tag,
2346            snapshot,
2347            registered_store_path,
2348        )
2349        .map(Some)
2350    }
2351
2352    /// Resolve one entity tag from the canonical accepted predecessor.
2353    pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_tag(
2354        &self,
2355        registered_store_path: &'static str,
2356        entity_tag: EntityTag,
2357    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2358        let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2359            return Ok(None);
2360        };
2361        if bundle.store_path() != registered_store_path {
2362            return Err(InternalError::store_corruption());
2363        }
2364        let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2365            return Ok(None);
2366        };
2367
2368        AcceptedRuntimeEntity::from_accepted_snapshot(
2369            &bundle,
2370            entity_tag,
2371            snapshot,
2372            registered_store_path,
2373        )
2374        .map(Some)
2375    }
2376
2377    /// Resolve one accepted entity source path without materializing the full store catalog.
2378    pub(in crate::db) fn current_accepted_runtime_entity_for_path(
2379        &self,
2380        registered_store_path: &'static str,
2381        entity_path: &str,
2382    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2383        self.current_accepted_runtime_entity_matching(registered_store_path, |snapshot_path, _| {
2384            snapshot_path == entity_path
2385        })
2386    }
2387
2388    /// Resolve one entity path from the canonical accepted predecessor.
2389    pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_path(
2390        &self,
2391        registered_store_path: &'static str,
2392        entity_path: &str,
2393    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2394        let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2395            return Ok(None);
2396        };
2397        if bundle.store_path() != registered_store_path {
2398            return Err(InternalError::store_corruption());
2399        }
2400
2401        let mut matched = None;
2402        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2403            if snapshot.entity_path() != entity_path {
2404                continue;
2405            }
2406            let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2407                &bundle,
2408                *entity_tag,
2409                snapshot,
2410                registered_store_path,
2411            )?;
2412            if matched.replace(entity).is_some() {
2413                return Err(InternalError::store_corruption());
2414            }
2415        }
2416
2417        Ok(matched)
2418    }
2419
2420    /// Resolve one accepted entity display name without materializing the full store catalog.
2421    #[cfg(test)]
2422    pub(in crate::db) fn current_accepted_runtime_entity_for_name(
2423        &self,
2424        registered_store_path: &'static str,
2425        entity_name: &str,
2426    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2427        self.current_accepted_runtime_entity_matching(registered_store_path, |_, snapshot_name| {
2428            snapshot_name == entity_name
2429        })
2430    }
2431
2432    fn current_accepted_runtime_entity_matching(
2433        &self,
2434        registered_store_path: &'static str,
2435        mut predicate: impl FnMut(&str, &str) -> bool,
2436    ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2437        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2438            return Ok(None);
2439        };
2440        if bundle.store_path() != registered_store_path {
2441            return Err(InternalError::store_corruption());
2442        }
2443
2444        let mut matched = None;
2445        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2446            if !predicate(snapshot.entity_path(), snapshot.entity_name()) {
2447                continue;
2448            }
2449            let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2450                &bundle,
2451                *entity_tag,
2452                snapshot,
2453                registered_store_path,
2454            )?;
2455            if matched.replace(entity).is_some() {
2456                return Err(InternalError::store_corruption());
2457            }
2458        }
2459
2460        Ok(matched)
2461    }
2462
2463    /// Return the current accepted revision without decoding its bundle.
2464    pub(in crate::db) fn current_accepted_schema_revision(
2465        &self,
2466    ) -> Result<Option<AcceptedSchemaRevision>, InternalError> {
2467        Ok(self
2468            .current_accepted_schema_root()?
2469            .map(|selection| selection.root().revision()))
2470    }
2471
2472    /// Return the pending relation activation that blocks deletes from one target.
2473    ///
2474    /// This reads the immutable accepted-bundle cache directly so ordinary
2475    /// deletes do not decode and clone every store catalog merely to prove that
2476    /// no candidate reverse generation targets the deleted entity.
2477    pub(in crate::db) fn pending_relation_activation_for_target(
2478        &self,
2479        target_path: &str,
2480    ) -> Result<Option<PendingRelationActivationDeleteBarrier>, InternalError> {
2481        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2482            return Ok(None);
2483        };
2484        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2485            let Some(candidate) = snapshot
2486                .candidate_relations()
2487                .iter()
2488                .find(|candidate| candidate.target_path() == target_path)
2489            else {
2490                continue;
2491            };
2492            let activation = snapshot
2493                .constraint_activations()
2494                .iter()
2495                .find(|activation| {
2496                    matches!(
2497                        activation.kind(),
2498                        ConstraintActivationKind::Relation { relation_id }
2499                            if *relation_id == candidate.id()
2500                    )
2501                })
2502                .ok_or_else(InternalError::store_corruption)?;
2503            return Ok(Some(PendingRelationActivationDeleteBarrier {
2504                accepted_schema_fingerprint:
2505                    accepted_schema_cache_fingerprint_for_persisted_snapshot(snapshot)?,
2506                source_entity_tag: *entity_tag,
2507                constraint_id: activation.id(),
2508            }));
2509        }
2510
2511        Ok(None)
2512    }
2513
2514    /// Return whether one accepted source entity owns a live relation to a target.
2515    pub(in crate::db) fn entity_has_relation_to_target(
2516        &self,
2517        source_entity: EntityTag,
2518        target_path: &str,
2519    ) -> Result<bool, InternalError> {
2520        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2521            return Ok(false);
2522        };
2523        let Some(snapshot) = bundle.entity_snapshots().get(&source_entity) else {
2524            return Ok(false);
2525        };
2526
2527        Ok(snapshot
2528            .relations()
2529            .iter()
2530            .any(|relation| relation.target_path() == target_path))
2531    }
2532
2533    /// Reject any same-entity schema change beside one exact activation lifecycle step.
2534    pub(in crate::db) fn validate_live_activation_transition(
2535        &self,
2536        candidate: &AcceptedSchemaRevisionBundle,
2537    ) -> Result<(), InternalError> {
2538        let Some(current) = self.current_accepted_schema_bundle()? else {
2539            return Ok(());
2540        };
2541        Self::validate_activation_transition_from(&current, candidate)
2542    }
2543
2544    /// Validate one transition against the canonical accepted predecessor.
2545    pub(in crate::db) fn validate_canonical_activation_transition(
2546        &self,
2547        candidate: &AcceptedSchemaRevisionBundle,
2548    ) -> Result<(), InternalError> {
2549        let Some(current) = self.current_canonical_accepted_schema_bundle()? else {
2550            return Ok(());
2551        };
2552        Self::validate_activation_transition_from(&current, candidate)
2553    }
2554
2555    fn validate_activation_transition_from(
2556        current: &AcceptedSchemaRevisionBundle,
2557        candidate: &AcceptedSchemaRevisionBundle,
2558    ) -> Result<(), InternalError> {
2559        for (entity_tag, before) in current.entity_snapshots() {
2560            if before.constraint_activations().is_empty() {
2561                continue;
2562            }
2563            let after = candidate
2564                .entity_snapshots()
2565                .get(entity_tag)
2566                .ok_or_else(InternalError::store_invariant)?;
2567            if before == after {
2568                continue;
2569            }
2570            let expected_shape = before
2571                .clone()
2572                .with_constraint_catalog(after.constraint_catalog().clone());
2573            let catalog_only_transition = expected_shape == *after
2574                && before
2575                    .constraint_catalog()
2576                    .permits_live_activation_transition_to(after.constraint_catalog());
2577            let sql_row_local_abort_with_version =
2578                before.constraint_activations().iter().any(|activation| {
2579                    activation.origin() == ConstraintOrigin::SqlDdl
2580                        && matches!(
2581                            activation.kind(),
2582                            ConstraintActivationKind::Check { .. }
2583                                | ConstraintActivationKind::NotNull { .. }
2584                        )
2585                        && before.version().get().checked_add(1) == Some(after.version().get())
2586                        && before
2587                            .constraint_catalog()
2588                            .clone()
2589                            .with_aborted_activation(activation.id())
2590                            .is_ok_and(|catalog| catalog == *after.constraint_catalog())
2591                        && before
2592                            .clone()
2593                            .with_constraint_catalog(after.constraint_catalog().clone())
2594                            .with_schema_version(after.version())
2595                            == *after
2596                });
2597            let sql_unique_abort_with_version =
2598                before.constraint_activations().iter().any(|activation| {
2599                    activation.origin() == ConstraintOrigin::SqlDdl
2600                        && matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2601                        && before.version().get().checked_add(1) == Some(after.version().get())
2602                        && before
2603                            .with_aborted_unique_activation(activation.id(), after.version())
2604                            .is_ok_and(|expected| expected == *after)
2605                });
2606            let not_null_promotion = before.constraint_activations().iter().any(|activation| {
2607                matches!(activation.kind(), ConstraintActivationKind::NotNull { .. })
2608                    && before
2609                        .with_promoted_not_null_activation(activation.id(), after.version())
2610                        .is_ok_and(|expected| expected == *after)
2611            });
2612            let unique_promotion = before.constraint_activations().iter().any(|activation| {
2613                matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2614                    && before
2615                        .with_promoted_unique_activation(activation.id(), after.version())
2616                        .is_ok_and(|expected| expected == *after)
2617            });
2618            let relation_promotion = before.constraint_activations().iter().any(|activation| {
2619                matches!(activation.kind(), ConstraintActivationKind::Relation { .. })
2620                    && before
2621                        .with_promoted_relation_activation(activation.id(), after.version())
2622                        .is_ok_and(|expected| expected == *after)
2623            });
2624            if !catalog_only_transition
2625                && !sql_row_local_abort_with_version
2626                && !sql_unique_abort_with_version
2627                && !not_null_promotion
2628                && !unique_promotion
2629                && !relation_promotion
2630            {
2631                return Err(InternalError::store_invariant());
2632            }
2633        }
2634        Ok(())
2635    }
2636
2637    /// Prove exact pairing between live activations and durable validation jobs.
2638    pub(in crate::db) fn validate_constraint_validation_job_closure(
2639        &self,
2640        bundle: &AcceptedSchemaRevisionBundle,
2641    ) -> Result<(), InternalError> {
2642        self.validate_constraint_validation_job_closure_with_change(bundle, None, None)
2643    }
2644
2645    /// Prove the activation/job closure that would exist after one bounded
2646    /// marker-owned job replacement or removal.
2647    pub(in crate::db) fn validate_constraint_validation_job_closure_with_change(
2648        &self,
2649        bundle: &AcceptedSchemaRevisionBundle,
2650        replacement: Option<&ConstraintValidationJob>,
2651        removal: Option<(EntityTag, ConstraintId)>,
2652    ) -> Result<(), InternalError> {
2653        self.validate_constraint_validation_job_closure_with_change_in_view(
2654            bundle,
2655            replacement,
2656            removal,
2657            IdentityStateStorageView::Effective,
2658        )
2659    }
2660
2661    /// Prove activation/job closure against the canonical predecessor view.
2662    pub(in crate::db) fn validate_canonical_constraint_validation_job_closure_with_change(
2663        &self,
2664        bundle: &AcceptedSchemaRevisionBundle,
2665        replacement: Option<&ConstraintValidationJob>,
2666        removal: Option<(EntityTag, ConstraintId)>,
2667    ) -> Result<(), InternalError> {
2668        self.validate_constraint_validation_job_closure_with_change_in_view(
2669            bundle,
2670            replacement,
2671            removal,
2672            IdentityStateStorageView::Canonical,
2673        )
2674    }
2675
2676    fn validate_constraint_validation_job_closure_with_change_in_view(
2677        &self,
2678        bundle: &AcceptedSchemaRevisionBundle,
2679        replacement: Option<&ConstraintValidationJob>,
2680        removal: Option<(EntityTag, ConstraintId)>,
2681        view: IdentityStateStorageView,
2682    ) -> Result<(), InternalError> {
2683        if replacement.is_some() && removal.is_some() {
2684            return Err(InternalError::store_invariant());
2685        }
2686        let replacement_key = replacement.map(|job| {
2687            RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id())
2688        });
2689        let removal_key = removal.map(|(entity_tag, constraint_id)| {
2690            RawSchemaKey::from_constraint_validation_job(entity_tag, constraint_id)
2691        });
2692        let mut expected = BTreeSet::new();
2693        for (entity_tag, snapshot) in bundle.entity_snapshots() {
2694            for activation in snapshot.constraint_activations() {
2695                let key =
2696                    RawSchemaKey::from_constraint_validation_job(*entity_tag, activation.id());
2697                match activation.state() {
2698                    ConstraintActivationState::EnforcingNewWrites => {
2699                        if self
2700                            .constraint_validation_job_after_change(
2701                                key,
2702                                replacement,
2703                                replacement_key,
2704                                removal_key,
2705                                view,
2706                            )?
2707                            .is_some()
2708                        {
2709                            return Err(InternalError::store_corruption());
2710                        }
2711                    }
2712                    ConstraintActivationState::Validating => {
2713                        let job = self
2714                            .constraint_validation_job_after_change(
2715                                key,
2716                                replacement,
2717                                replacement_key,
2718                                removal_key,
2719                                view,
2720                            )?
2721                            .ok_or_else(InternalError::store_corruption)?;
2722                        if job.entity_tag() != *entity_tag
2723                            || job.entity_path() != snapshot.entity_path()
2724                        {
2725                            return Err(InternalError::store_corruption());
2726                        }
2727                        job.validate(Some(activation))?;
2728                        expected.insert(key);
2729                    }
2730                }
2731            }
2732        }
2733
2734        self.visit_constraint_validation_jobs_in_view(view, |key, raw| {
2735            if removal_key == Some(*key) || replacement_key == Some(*key) {
2736                return Ok(SchemaStoreVisit::Continue);
2737            }
2738            if !expected.contains(key) {
2739                return Err(InternalError::store_corruption());
2740            }
2741            let job = decode_constraint_validation_job(raw.as_bytes())?;
2742            if job.entity_tag() != key.entity_tag()
2743                || key.constraint_id() != Some(job.constraint_id())
2744            {
2745                return Err(InternalError::store_corruption());
2746            }
2747            Ok(SchemaStoreVisit::Continue)
2748        })?;
2749
2750        if let Some(key) = replacement_key
2751            && !expected.contains(&key)
2752        {
2753            return Err(InternalError::store_corruption());
2754        }
2755        if let Some(key) = removal_key
2756            && expected.contains(&key)
2757        {
2758            return Err(InternalError::store_corruption());
2759        }
2760
2761        Ok(())
2762    }
2763
2764    fn constraint_validation_job_after_change(
2765        &self,
2766        key: RawSchemaKey,
2767        replacement: Option<&ConstraintValidationJob>,
2768        replacement_key: Option<RawSchemaKey>,
2769        removal_key: Option<RawSchemaKey>,
2770        view: IdentityStateStorageView,
2771    ) -> Result<Option<ConstraintValidationJob>, InternalError> {
2772        if removal_key == Some(key) {
2773            return Ok(None);
2774        }
2775        if replacement_key == Some(key) {
2776            return Ok(replacement.cloned());
2777        }
2778        let raw = match view {
2779            IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
2780            IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
2781        };
2782        raw.map(|raw| decode_constraint_validation_job(raw.as_bytes()))
2783            .transpose()
2784    }
2785
2786    /// Return whether one retained schema authority still names this store's
2787    /// current immutable accepted root.
2788    pub(in crate::db) fn current_accepted_schema_authority_matches(
2789        &self,
2790        expected: &AcceptedSchemaAuthority,
2791    ) -> Result<bool, InternalError> {
2792        self.current_accepted_selection_for_authority(expected)
2793            .map(|selection| selection.is_some())
2794    }
2795
2796    // One store-scoped authority check for execution and borrowed metadata.
2797    // A dropped cache is not lost authority: reload its durable selection.
2798    fn current_accepted_selection_for_authority(
2799        &self,
2800        expected: &AcceptedSchemaAuthority,
2801    ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
2802        let Some(store_scope) = self.accepted_catalog_scope.get() else {
2803            return Ok(None);
2804        };
2805
2806        // Root-writing primitives invalidate this cache before publication,
2807        // so a retained selection is the current in-memory authority.
2808        let cached = self
2809            .accepted_bundle_cache
2810            .try_borrow()
2811            .map_err(|_| InternalError::store_invariant())?
2812            .as_ref()
2813            .map(|cached| cached.selection);
2814        let selection = match cached {
2815            Some(selection) => Some(selection),
2816            None => self.current_accepted_schema_root()?,
2817        };
2818        Ok(selection.filter(|selection| {
2819            let root = selection.root();
2820            expected.matches_store_root(store_scope, root.revision(), root.fingerprint())
2821        }))
2822    }
2823
2824    /// Publish a candidate directly into its canonical schema allocation.
2825    ///
2826    /// Journaled online revisions must use
2827    /// `apply_journaled_accepted_schema_candidate`; this path owns initial
2828    /// bootstrap and marker-owned live-projection updates.
2829    pub(in crate::db) fn publish_accepted_schema_candidate(
2830        &mut self,
2831        incarnation: DatabaseIncarnationId,
2832        expected_revision: AcceptedSchemaRevision,
2833        candidate: &CandidateSchemaRevision,
2834    ) -> Result<(), InternalError> {
2835        let identity_transition = self.prepare_identity_state_transition(
2836            incarnation,
2837            candidate,
2838            IdentityStateStorageView::Effective,
2839        )?;
2840        if self.current_root_matches_candidate(candidate)? {
2841            if !identity_transition.is_empty() {
2842                return Err(InternalError::identity_state_corruption());
2843            }
2844            let selection = self
2845                .current_accepted_schema_root()?
2846                .ok_or_else(InternalError::store_corruption)?;
2847            self.retain_durable_candidate_entries(candidate, selection.slot())?;
2848            return Ok(());
2849        }
2850        let first = self.accepted_root_slot_bytes(0)?;
2851        let second = self.accepted_root_slot_bytes(1)?;
2852        prepare_accepted_schema_root_publication(
2853            [first.as_deref(), second.as_deref()],
2854            expected_revision,
2855            candidate,
2856        )
2857        .map_err(map_schema_publication_error)?;
2858
2859        self.insert_durable_candidate_snapshots(candidate)?;
2860        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2861        self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2862        let persisted_bundle = self
2863            .get_raw_snapshot(&bundle_key)
2864            .ok_or_else(InternalError::store_corruption)?;
2865        let _verified = decode_verified_accepted_schema_revision_bundle(
2866            candidate.root(),
2867            persisted_bundle.as_bytes(),
2868        )?;
2869        self.apply_identity_state_transition(
2870            identity_transition,
2871            IdentityStateWriteTarget::Durable,
2872        )?;
2873
2874        // Re-read the root immediately before the inactive-slot write. This is
2875        // the compare-and-swap check after candidate persistence.
2876        let first = self.accepted_root_slot_bytes(0)?;
2877        let second = self.accepted_root_slot_bytes(1)?;
2878        let publication = prepare_accepted_schema_root_publication(
2879            [first.as_deref(), second.as_deref()],
2880            expected_revision,
2881            candidate,
2882        )
2883        .map_err(map_schema_publication_error)?;
2884        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
2885        self.insert_durable_raw_value(root_key, publication.encoded_root().to_vec());
2886
2887        let selected = self
2888            .current_accepted_schema_root()?
2889            .ok_or_else(InternalError::store_corruption)?;
2890        if selected.root() != candidate.root() {
2891            return Err(InternalError::store_corruption());
2892        }
2893        self.retain_durable_candidate_entries(candidate, selected.slot())?;
2894        Ok(())
2895    }
2896
2897    /// Restore one current accepted candidate into an empty live-only schema
2898    /// store from its durable database-control checkpoint.
2899    pub(in crate::db) fn restore_live_accepted_schema_checkpoint(
2900        &mut self,
2901        incarnation: DatabaseIncarnationId,
2902        candidate: &CandidateSchemaRevision,
2903        checkpoint_identity_states: &IdentityStateInventory,
2904    ) -> Result<(), InternalError> {
2905        if !matches!(self.backend, SchemaStoreBackend::Heap(_)) {
2906            return Err(InternalError::store_invariant());
2907        }
2908        let checkpoint_validation = prepare_identity_state_transition(
2909            incarnation,
2910            Some(candidate.bundle()),
2911            candidate.bundle(),
2912            checkpoint_identity_states.clone(),
2913        )?;
2914        if !checkpoint_validation.is_empty() {
2915            return Err(InternalError::identity_state_corruption());
2916        }
2917        if self.current_root_matches_candidate(candidate)? {
2918            for state in checkpoint_identity_states.values() {
2919                let key = RawSchemaKey::from_identity_state(
2920                    state.owner().entity_tag(),
2921                    state.owner().field_id(),
2922                );
2923                self.insert_durable_raw_value(key, encode_identity_state(state)?);
2924            }
2925            if self.identity_state_inventory(IdentityStateStorageView::Effective)?
2926                != *checkpoint_identity_states
2927            {
2928                return Err(InternalError::identity_state_corruption());
2929            }
2930            let selection = self
2931                .current_accepted_schema_root()?
2932                .ok_or_else(InternalError::store_corruption)?;
2933            self.retain_durable_candidate_entries(candidate, selection.slot())?;
2934            return Ok(());
2935        }
2936        if self.current_accepted_schema_root()?.is_some()
2937            || !self
2938                .identity_state_inventory(IdentityStateStorageView::Effective)?
2939                .is_empty()
2940        {
2941            return Err(InternalError::store_corruption());
2942        }
2943
2944        self.insert_durable_candidate_snapshots(candidate)?;
2945        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2946        self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2947        for state in checkpoint_identity_states.values() {
2948            let key = RawSchemaKey::from_identity_state(
2949                state.owner().entity_tag(),
2950                state.owner().field_id(),
2951            );
2952            self.insert_durable_raw_value(key, encode_identity_state(state)?);
2953        }
2954        let root_key = RawSchemaKey::from_accepted_root_slot(0)?;
2955        self.insert_durable_raw_value(root_key, candidate.encoded_root().to_vec());
2956
2957        let selected = self
2958            .current_accepted_schema_root()?
2959            .ok_or_else(InternalError::store_corruption)?;
2960        if selected.root() != candidate.root() {
2961            return Err(InternalError::store_corruption());
2962        }
2963        self.retain_durable_candidate_entries(candidate, selected.slot())?;
2964        Ok(())
2965    }
2966
2967    /// Preflight one accepted candidate without changing durable or live
2968    /// schema state.
2969    ///
2970    /// Returns `true` only when this exact candidate is already authoritative.
2971    /// Multi-store publication uses that distinction to reject partial replay
2972    /// before opening one marker-owned commit window.
2973    pub(in crate::db) fn preflight_accepted_schema_candidate(
2974        &self,
2975        incarnation: DatabaseIncarnationId,
2976        expected_revision: AcceptedSchemaRevision,
2977        candidate: &CandidateSchemaRevision,
2978    ) -> Result<bool, InternalError> {
2979        let identity_transition = self.prepare_identity_state_transition(
2980            incarnation,
2981            candidate,
2982            IdentityStateStorageView::Effective,
2983        )?;
2984        if self.current_root_matches_candidate(candidate)? {
2985            if !identity_transition.is_empty() {
2986                return Err(InternalError::identity_state_corruption());
2987            }
2988            return Ok(true);
2989        }
2990        let first = self.accepted_root_slot_bytes(0)?;
2991        let second = self.accepted_root_slot_bytes(1)?;
2992        prepare_accepted_schema_root_publication(
2993            [first.as_deref(), second.as_deref()],
2994            expected_revision,
2995            candidate,
2996        )
2997        .map_err(map_schema_publication_error)?;
2998
2999        Ok(false)
3000    }
3001
3002    /// Prepare one accepted candidate against canonical journaled authority.
3003    pub(in crate::db) fn prepare_fold_journaled_accepted_schema_candidate(
3004        &self,
3005        incarnation: DatabaseIncarnationId,
3006        expected_revision: AcceptedSchemaRevision,
3007        candidate: CandidateSchemaRevision,
3008    ) -> Result<PreparedAcceptedSchemaFold, InternalError> {
3009        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3010            return Err(InternalError::store_invariant());
3011        }
3012        let identity_transition = self.prepare_identity_state_transition(
3013            incarnation,
3014            &candidate,
3015            IdentityStateStorageView::Canonical,
3016        )?;
3017        let candidate_is_current = self.canonical_root_matches_candidate(&candidate)?;
3018        if candidate_is_current && !identity_transition.is_empty() {
3019            return Err(InternalError::identity_state_corruption());
3020        }
3021        let (updates, removals) = identity_transition.into_effects();
3022        let identity_removals = removals
3023            .into_iter()
3024            .map(|owner| RawSchemaKey::from_identity_state(owner.entity_tag(), owner.field_id()))
3025            .collect();
3026        let identity_updates = updates
3027            .into_iter()
3028            .map(|state| {
3029                Ok((
3030                    RawSchemaKey::from_identity_state(
3031                        state.owner().entity_tag(),
3032                        state.owner().field_id(),
3033                    ),
3034                    encode_identity_state(&state)?,
3035                ))
3036            })
3037            .collect::<Result<Vec<_>, InternalError>>()?;
3038        let snapshots = candidate
3039            .bundle()
3040            .entity_snapshots()
3041            .iter()
3042            .map(|(entity, snapshot)| Self::prepare_persisted_snapshot(*entity, snapshot))
3043            .collect::<Result<Vec<_>, _>>()?;
3044
3045        let first = self.canonical_root_slot_bytes(0)?;
3046        let second = self.canonical_root_slot_bytes(1)?;
3047        let root_slot = if candidate_is_current {
3048            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3049                .ok_or_else(InternalError::store_corruption)?
3050                .slot()
3051        } else {
3052            prepare_accepted_schema_root_publication(
3053                [first.as_deref(), second.as_deref()],
3054                expected_revision,
3055                &candidate,
3056            )
3057            .map_err(map_schema_publication_error)?
3058            .target_slot()
3059        };
3060        let retained = Self::candidate_entry_keys(&candidate, root_slot)?;
3061        Ok(PreparedAcceptedSchemaFold {
3062            candidate,
3063            expected_revision,
3064            snapshots,
3065            identity_updates,
3066            identity_removals,
3067            retained,
3068            root_slot,
3069        })
3070    }
3071
3072    /// Return the retained Identity owner count after admitting one candidate.
3073    pub(in crate::db) fn projected_identity_state_count(
3074        &self,
3075        incarnation: DatabaseIncarnationId,
3076        candidate: &CandidateSchemaRevision,
3077    ) -> Result<usize, InternalError> {
3078        Ok(self
3079            .prepare_identity_state_transition(
3080                incarnation,
3081                candidate,
3082                IdentityStateStorageView::Effective,
3083            )?
3084            .projected_inventory_len())
3085    }
3086
3087    /// Apply one marker-bound schema candidate to the journaled live projection.
3088    pub(in crate::db) fn apply_journaled_accepted_schema_candidate(
3089        &mut self,
3090        incarnation: DatabaseIncarnationId,
3091        expected_revision: AcceptedSchemaRevision,
3092        candidate: &CandidateSchemaRevision,
3093    ) -> Result<(), InternalError> {
3094        if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3095            return Err(InternalError::store_invariant());
3096        }
3097        let identity_transition = self.prepare_identity_state_transition(
3098            incarnation,
3099            candidate,
3100            IdentityStateStorageView::Effective,
3101        )?;
3102        if self.current_root_matches_candidate(candidate)? {
3103            if !identity_transition.is_empty() {
3104                return Err(InternalError::identity_state_corruption());
3105            }
3106            let selection = self
3107                .current_accepted_schema_root()?
3108                .ok_or_else(InternalError::store_corruption)?;
3109            self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3110            return Ok(());
3111        }
3112
3113        let first = self.accepted_root_slot_bytes(0)?;
3114        let second = self.accepted_root_slot_bytes(1)?;
3115        prepare_accepted_schema_root_publication(
3116            [first.as_deref(), second.as_deref()],
3117            expected_revision,
3118            candidate,
3119        )
3120        .map_err(map_schema_publication_error)?;
3121
3122        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3123            self.insert_persisted_snapshot(*entity_tag, snapshot)?;
3124        }
3125        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3126        self.insert_raw_snapshot(
3127            bundle_key,
3128            RawSchemaSnapshot::from_encoded_control_record(candidate.encoded_bundle().to_vec()),
3129        );
3130        let persisted_bundle = self
3131            .get_raw_snapshot(&bundle_key)
3132            .ok_or_else(InternalError::store_corruption)?;
3133        let _verified = decode_verified_accepted_schema_revision_bundle(
3134            candidate.root(),
3135            persisted_bundle.as_bytes(),
3136        )?;
3137        self.apply_identity_state_transition(
3138            identity_transition,
3139            IdentityStateWriteTarget::Materialized,
3140        )?;
3141
3142        let first = self.accepted_root_slot_bytes(0)?;
3143        let second = self.accepted_root_slot_bytes(1)?;
3144        let publication = prepare_accepted_schema_root_publication(
3145            [first.as_deref(), second.as_deref()],
3146            expected_revision,
3147            candidate,
3148        )
3149        .map_err(map_schema_publication_error)?;
3150        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3151        self.insert_raw_snapshot(
3152            root_key,
3153            RawSchemaSnapshot::from_encoded_control_record(publication.encoded_root().to_vec()),
3154        );
3155
3156        if !self.current_root_matches_candidate(candidate)? {
3157            return Err(InternalError::store_corruption());
3158        }
3159        let selection = self
3160            .current_accepted_schema_root()?
3161            .ok_or_else(InternalError::store_corruption)?;
3162        self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3163        Ok(())
3164    }
3165
3166    /// Consume one candidate's preflight payloads in the same atomic fold callback.
3167    pub(in crate::db) fn apply_prepared_accepted_schema_fold(
3168        &mut self,
3169        prepared: PreparedAcceptedSchemaFold,
3170    ) -> Result<(), InternalError> {
3171        let PreparedAcceptedSchemaFold {
3172            candidate,
3173            expected_revision,
3174            snapshots,
3175            identity_updates,
3176            identity_removals,
3177            retained,
3178            root_slot,
3179        } = prepared;
3180        if self.canonical_root_matches_candidate(&candidate)? {
3181            if !identity_updates.is_empty() || !identity_removals.is_empty() {
3182                return Err(InternalError::identity_state_corruption());
3183            }
3184            let first = self.canonical_root_slot_bytes(0)?;
3185            let second = self.canonical_root_slot_bytes(1)?;
3186            let selection =
3187                select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3188                    .ok_or_else(InternalError::store_corruption)?;
3189            if selection.slot() != root_slot {
3190                return Err(InternalError::store_invariant());
3191            }
3192            self.retain_canonical_candidate_entries(&retained)?;
3193            return Ok(());
3194        }
3195
3196        let first = self.canonical_root_slot_bytes(0)?;
3197        let second = self.canonical_root_slot_bytes(1)?;
3198        let publication = prepare_accepted_schema_root_publication(
3199            [first.as_deref(), second.as_deref()],
3200            expected_revision,
3201            &candidate,
3202        )
3203        .map_err(map_schema_publication_error)?;
3204        if publication.target_slot() != root_slot {
3205            return Err(InternalError::store_invariant());
3206        }
3207
3208        for snapshot in snapshots {
3209            self.apply_prepared_fold_persisted_snapshot(snapshot)?;
3210        }
3211        let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3212        self.insert_canonical_raw_value(bundle_key, candidate.encoded_bundle().to_vec())?;
3213        let persisted_bundle = self
3214            .get_canonical_raw_value(&bundle_key)?
3215            .ok_or_else(InternalError::store_corruption)?;
3216        let _verified = decode_verified_accepted_schema_revision_bundle(
3217            candidate.root(),
3218            persisted_bundle.as_bytes(),
3219        )?;
3220        for key in identity_removals {
3221            self.remove_identity_state_key(key, IdentityStateWriteTarget::Canonical)?;
3222        }
3223        for (key, bytes) in identity_updates {
3224            self.insert_canonical_raw_value(key, bytes)?;
3225        }
3226
3227        let first = self.canonical_root_slot_bytes(0)?;
3228        let second = self.canonical_root_slot_bytes(1)?;
3229        let publication = prepare_accepted_schema_root_publication(
3230            [first.as_deref(), second.as_deref()],
3231            expected_revision,
3232            &candidate,
3233        )
3234        .map_err(map_schema_publication_error)?;
3235        let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3236        self.insert_canonical_raw_value(root_key, publication.encoded_root().to_vec())?;
3237
3238        if !self.canonical_root_matches_candidate(&candidate)? {
3239            return Err(InternalError::store_corruption());
3240        }
3241        let first = self.canonical_root_slot_bytes(0)?;
3242        let second = self.canonical_root_slot_bytes(1)?;
3243        let selection = select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3244            .ok_or_else(InternalError::store_corruption)?;
3245        if selection.slot() != root_slot {
3246            return Err(InternalError::store_invariant());
3247        }
3248        self.retain_canonical_candidate_entries(&retained)?;
3249        Ok(())
3250    }
3251
3252    /// Load and decode one typed persisted schema snapshot.
3253    pub(in crate::db) fn get_persisted_snapshot(
3254        &self,
3255        entity: EntityTag,
3256        version: SchemaVersion,
3257    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3258        let key = RawSchemaKey::from_entity_version(entity, version);
3259        self.get_raw_snapshot(&key)
3260            .map(|snapshot| snapshot.decode_persisted_snapshot())
3261            .transpose()
3262    }
3263
3264    #[cfg(test)]
3265    fn latest_staged_persisted_snapshot(
3266        &self,
3267        entity: EntityTag,
3268    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3269        self.latest_raw_snapshots_by_entity()
3270            .remove(&entity)
3271            .map(|(_, snapshot)| snapshot.decode_persisted_snapshot())
3272            .transpose()
3273    }
3274
3275    /// Load one entity snapshot from the immutable bundle selected by the
3276    /// current accepted root.
3277    pub(in crate::db) fn current_accepted_persisted_snapshot(
3278        &self,
3279        entity: EntityTag,
3280    ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3281        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3282            return Ok(None);
3283        };
3284
3285        Ok(bundle.entity_snapshots().get(&entity).cloned())
3286    }
3287
3288    /// Return one accepted catalog selection from the current immutable root.
3289    pub(in crate::db) fn current_accepted_catalog_selection(
3290        &self,
3291        entity: EntityTag,
3292        entity_path: &str,
3293        store_path: &'static str,
3294    ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3295        let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3296            return Ok(None);
3297        };
3298        if bundle.store_path() != store_path {
3299            return Err(InternalError::store_corruption());
3300        }
3301        let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3302            return Ok(None);
3303        };
3304        if snapshot.entity_path() != entity_path {
3305            return Err(InternalError::store_corruption());
3306        }
3307
3308        let cache = self
3309            .accepted_bundle_cache
3310            .try_borrow()
3311            .map_err(|_| InternalError::store_invariant())?;
3312        let cached = cache.as_ref().ok_or_else(InternalError::store_invariant)?;
3313        if let Some(selection) = cached
3314            .entity_selections
3315            .try_borrow()
3316            .map_err(|_| InternalError::store_invariant())?
3317            .get(&entity)
3318            .cloned()
3319        {
3320            return Ok(Some(selection));
3321        }
3322
3323        let selected = AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3324            entity,
3325            store_path,
3326            snapshot,
3327            cached.value_catalog.clone(),
3328        )?;
3329        cached
3330            .entity_selections
3331            .try_borrow_mut()
3332            .map_err(|_| InternalError::store_invariant())?
3333            .insert(entity, selected.clone());
3334
3335        Ok(Some(selected))
3336    }
3337
3338    /// Return one accepted catalog selection from the canonical journal base.
3339    /// Recovery uses this while folding historical row batches whose schema
3340    /// revision can precede the current live accepted root.
3341    pub(in crate::db) fn current_canonical_accepted_catalog_selection(
3342        &self,
3343        entity: EntityTag,
3344        entity_path: &str,
3345        store_path: &'static str,
3346    ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3347        let first = self.canonical_root_slot_bytes(0)?;
3348        let second = self.canonical_root_slot_bytes(1)?;
3349        let Some(selection) =
3350            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3351        else {
3352            return Ok(None);
3353        };
3354        let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3355        let raw_bundle = self
3356            .get_canonical_raw_value(&bundle_key)?
3357            .ok_or_else(InternalError::store_corruption)?;
3358        let bundle = decode_verified_accepted_schema_revision_bundle(
3359            selection.root(),
3360            raw_bundle.as_bytes(),
3361        )?;
3362        if bundle.store_path() != store_path {
3363            return Err(InternalError::store_corruption());
3364        }
3365        let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3366            return Ok(None);
3367        };
3368        if snapshot.entity_path() != entity_path {
3369            return Err(InternalError::store_corruption());
3370        }
3371
3372        AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3373            entity,
3374            store_path,
3375            snapshot,
3376            AcceptedValueCatalogHandle::new(
3377                bundle.enum_catalog().clone(),
3378                bundle.composite_catalog().clone(),
3379                self.accepted_catalog_scope
3380                    .get_or_init(AcceptedStoreCatalogScope::new)
3381                    .clone(),
3382                bundle.revision(),
3383                selection.root().fingerprint(),
3384            ),
3385        )
3386        .map(Some)
3387    }
3388
3389    /// Derive accepted catalog metadata from latest persisted schema snapshots.
3390    ///
3391    /// This function intentionally reads only the persisted schema store. It
3392    /// does not reconstruct metadata from generated models when the store has
3393    /// no accepted snapshots.
3394    #[cfg(test)]
3395    pub(in crate::db) fn catalog_metadata(
3396        &self,
3397    ) -> Result<Option<SchemaStoreCatalogMetadata>, InternalError> {
3398        Ok(self
3399            .allocation_metadata()?
3400            .map(SchemaStoreAllocationMetadata::schema))
3401    }
3402
3403    /// Derive role-specific allocation metadata from latest persisted schema
3404    /// snapshots.
3405    ///
3406    /// This function intentionally reads only accepted schema-store payloads.
3407    /// It never reconstructs metadata from generated models when the store has
3408    /// no accepted snapshots.
3409    pub(in crate::db) fn allocation_metadata(
3410        &self,
3411    ) -> Result<Option<SchemaStoreAllocationMetadata>, InternalError> {
3412        let latest_by_entity = self.latest_raw_snapshots_by_entity();
3413        if latest_by_entity.is_empty() {
3414            return Ok(None);
3415        }
3416
3417        Ok(Some(SchemaStoreAllocationMetadata::new(
3418            derive_data_allocation_metadata(&latest_by_entity)?,
3419            derive_index_allocation_metadata(&latest_by_entity)?,
3420            derive_schema_catalog_metadata(&latest_by_entity)?,
3421        )))
3422    }
3423
3424    /// Insert or replace one raw schema snapshot.
3425    fn insert_raw_snapshot(
3426        &mut self,
3427        key: RawSchemaKey,
3428        snapshot: RawSchemaSnapshot,
3429    ) -> Option<RawSchemaSnapshot> {
3430        self.invalidate_accepted_bundle_cache_for_key(key);
3431        let previous_journaled = if matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3432            self.get_raw_snapshot_for_backend(&key)
3433        } else {
3434            None
3435        };
3436        match &mut self.backend {
3437            SchemaStoreBackend::Heap(map) => map.insert(key, snapshot),
3438            SchemaStoreBackend::Journaled {
3439                live, tombstones, ..
3440            } => {
3441                tombstones.remove(&key);
3442                live.insert(key, snapshot);
3443                previous_journaled
3444            }
3445        }
3446    }
3447
3448    /// Load one raw schema snapshot by key.
3449    #[must_use]
3450    fn get_raw_snapshot(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
3451        match &self.backend {
3452            SchemaStoreBackend::Heap(map) => map.get(key).cloned(),
3453            SchemaStoreBackend::Journaled { .. } => self.get_raw_snapshot_for_backend(key),
3454        }
3455    }
3456
3457    fn accepted_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3458        let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3459        Ok(self
3460            .get_raw_snapshot(&key)
3461            .map(RawSchemaSnapshot::into_bytes))
3462    }
3463
3464    fn canonical_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3465        let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3466        Ok(self
3467            .get_canonical_raw_value(&key)?
3468            .map(RawSchemaSnapshot::into_bytes))
3469    }
3470
3471    fn current_root_matches_candidate(
3472        &self,
3473        candidate: &CandidateSchemaRevision,
3474    ) -> Result<bool, InternalError> {
3475        let Some(selection) = self.current_accepted_schema_root()? else {
3476            return Ok(false);
3477        };
3478        if selection.root() != candidate.root() {
3479            return Ok(false);
3480        }
3481        let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3482        let bundle = self
3483            .get_raw_snapshot(&key)
3484            .ok_or_else(InternalError::store_corruption)?;
3485        let _verified =
3486            decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3487        Ok(true)
3488    }
3489
3490    fn canonical_root_matches_candidate(
3491        &self,
3492        candidate: &CandidateSchemaRevision,
3493    ) -> Result<bool, InternalError> {
3494        let first = self.canonical_root_slot_bytes(0)?;
3495        let second = self.canonical_root_slot_bytes(1)?;
3496        let Some(selection) =
3497            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3498        else {
3499            return Ok(false);
3500        };
3501        if selection.root() != candidate.root() {
3502            return Ok(false);
3503        }
3504        let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3505        let bundle = self
3506            .get_canonical_raw_value(&key)?
3507            .ok_or_else(InternalError::store_corruption)?;
3508        let _verified =
3509            decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3510        Ok(true)
3511    }
3512
3513    fn get_canonical_raw_value(
3514        &self,
3515        key: &RawSchemaKey,
3516    ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
3517        match &self.backend {
3518            SchemaStoreBackend::Journaled { canonical, .. } => Ok(canonical.get(key)),
3519            SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
3520        }
3521    }
3522
3523    fn insert_canonical_raw_value(
3524        &mut self,
3525        key: RawSchemaKey,
3526        bytes: Vec<u8>,
3527    ) -> Result<(), InternalError> {
3528        self.invalidate_accepted_bundle_cache_for_key(key);
3529        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3530            return Err(InternalError::store_invariant());
3531        };
3532        canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
3533        Ok(())
3534    }
3535
3536    // Initial accepted-catalog bootstrap persists immutable bundle/root values
3537    // directly in the schema allocation. Later online schema mutation will
3538    // carry the same values through the journal before calling this primitive.
3539    fn insert_durable_raw_value(&mut self, key: RawSchemaKey, bytes: Vec<u8>) {
3540        self.invalidate_accepted_bundle_cache_for_key(key);
3541        let value = RawSchemaSnapshot::from_encoded_control_record(bytes);
3542        match &mut self.backend {
3543            SchemaStoreBackend::Heap(map) => {
3544                map.insert(key, value);
3545            }
3546            SchemaStoreBackend::Journaled {
3547                canonical,
3548                live,
3549                tombstones,
3550                ..
3551            } => {
3552                live.remove(&key);
3553                tombstones.remove(&key);
3554                canonical.insert(key, value);
3555            }
3556        }
3557    }
3558
3559    fn invalidate_accepted_bundle_cache_for_key(&mut self, key: RawSchemaKey) {
3560        if key.is_accepted_root() {
3561            self.accepted_bundle_cache.get_mut().take();
3562        }
3563    }
3564
3565    fn insert_durable_candidate_snapshots(
3566        &mut self,
3567        candidate: &CandidateSchemaRevision,
3568    ) -> Result<(), InternalError> {
3569        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3570            let key = RawSchemaKey::from_entity_version(*entity_tag, snapshot.version());
3571            let value = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3572            match &mut self.backend {
3573                SchemaStoreBackend::Heap(map) => {
3574                    map.insert(key, value);
3575                }
3576                SchemaStoreBackend::Journaled {
3577                    canonical,
3578                    live,
3579                    tombstones,
3580                    ..
3581                } => {
3582                    live.remove(&key);
3583                    tombstones.remove(&key);
3584                    canonical.insert(key, value);
3585                }
3586            }
3587        }
3588        Ok(())
3589    }
3590
3591    fn candidate_entry_keys(
3592        candidate: &CandidateSchemaRevision,
3593        root_slot: usize,
3594    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3595        let mut keys = candidate
3596            .bundle()
3597            .entity_snapshots()
3598            .iter()
3599            .map(|(entity_tag, snapshot)| {
3600                RawSchemaKey::from_entity_version(*entity_tag, snapshot.version())
3601            })
3602            .collect::<BTreeSet<_>>();
3603        keys.insert(RawSchemaKey::from_accepted_bundle(
3604            candidate.root().bundle_key(),
3605        ));
3606        keys.insert(RawSchemaKey::from_accepted_root_slot(root_slot)?);
3607        for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3608            for activation in snapshot
3609                .constraint_activations()
3610                .iter()
3611                .filter(|activation| activation.state() == ConstraintActivationState::Validating)
3612            {
3613                keys.insert(RawSchemaKey::from_constraint_validation_job(
3614                    *entity_tag,
3615                    activation.id(),
3616                ));
3617            }
3618        }
3619        Ok(keys)
3620    }
3621
3622    fn positioned_candidate_effect_keys(
3623        &self,
3624        incarnation: DatabaseIncarnationId,
3625        expected_revision: AcceptedSchemaRevision,
3626        candidate: &CandidateSchemaRevision,
3627        view: IdentityStateStorageView,
3628    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3629        let identity_transition =
3630            self.prepare_identity_state_transition(incarnation, candidate, view)?;
3631        let (first, second, candidate_is_current) = match view {
3632            IdentityStateStorageView::Effective => (
3633                self.accepted_root_slot_bytes(0)?,
3634                self.accepted_root_slot_bytes(1)?,
3635                self.current_root_matches_candidate(candidate)?,
3636            ),
3637            IdentityStateStorageView::Canonical => (
3638                self.canonical_root_slot_bytes(0)?,
3639                self.canonical_root_slot_bytes(1)?,
3640                self.canonical_root_matches_candidate(candidate)?,
3641            ),
3642        };
3643        let root_slot = if candidate_is_current {
3644            select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3645                .ok_or_else(InternalError::store_corruption)?
3646                .slot()
3647        } else {
3648            prepare_accepted_schema_root_publication(
3649                [first.as_deref(), second.as_deref()],
3650                expected_revision,
3651                candidate,
3652            )
3653            .map_err(map_schema_publication_error)?
3654            .target_slot()
3655        };
3656        let mut keys = Self::candidate_entry_keys(candidate, root_slot)?;
3657        let (updates, removals) = identity_transition.into_effects();
3658        for owner in removals {
3659            keys.insert(RawSchemaKey::from_identity_state(
3660                owner.entity_tag(),
3661                owner.field_id(),
3662            ));
3663        }
3664        for state in updates {
3665            keys.insert(RawSchemaKey::from_identity_state(
3666                state.owner().entity_tag(),
3667                state.owner().field_id(),
3668            ));
3669        }
3670
3671        let SchemaStoreBackend::Journaled {
3672            canonical,
3673            live,
3674            tombstones,
3675            positions,
3676        } = &self.backend
3677        else {
3678            return Err(InternalError::store_invariant());
3679        };
3680        for entry in canonical.iter() {
3681            let has_relevant_overlay = matches!(view, IdentityStateStorageView::Effective)
3682                || positions.is_positioned(entry.key())
3683                || live.contains_key(entry.key())
3684                || tombstones.contains(entry.key());
3685            if has_relevant_overlay
3686                && !keys.contains(entry.key())
3687                && !entry.key().is_identity_state()
3688            {
3689                keys.insert(*entry.key());
3690            }
3691        }
3692        if matches!(view, IdentityStateStorageView::Effective) {
3693            for key in live.keys() {
3694                if !keys.contains(key) && !key.is_identity_state() {
3695                    keys.insert(*key);
3696                }
3697            }
3698        }
3699        Ok(keys)
3700    }
3701
3702    fn positioned_journal_batch_keys(
3703        &self,
3704        incarnation: DatabaseIncarnationId,
3705        batch: &JournalBatch,
3706        view: IdentityStateStorageView,
3707    ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3708        let mut keys = BTreeSet::new();
3709        for record in batch.records() {
3710            match record {
3711                JournalRecord::SchemaPut {
3712                    schema_snapshot_bytes,
3713                    ..
3714                } => {
3715                    let snapshot = decode_persisted_schema_snapshot(schema_snapshot_bytes)?;
3716                    let entity_tag = match view {
3717                        IdentityStateStorageView::Effective => self
3718                            .current_accepted_schema_bundle_ref()?
3719                            .ok_or_else(InternalError::store_corruption)?
3720                            .entity_snapshots()
3721                            .iter()
3722                            .find_map(|(entity_tag, accepted)| {
3723                                (accepted.entity_path() == snapshot.entity_path())
3724                                    .then_some(*entity_tag)
3725                            }),
3726                        IdentityStateStorageView::Canonical => self
3727                            .current_canonical_accepted_schema_bundle()?
3728                            .ok_or_else(InternalError::store_corruption)?
3729                            .entity_snapshots()
3730                            .iter()
3731                            .find_map(|(entity_tag, accepted)| {
3732                                (accepted.entity_path() == snapshot.entity_path())
3733                                    .then_some(*entity_tag)
3734                            }),
3735                    }
3736                    .ok_or_else(InternalError::store_corruption)?;
3737                    keys.insert(RawSchemaKey::from_entity_version(
3738                        entity_tag,
3739                        snapshot.version(),
3740                    ));
3741                }
3742                JournalRecord::AcceptedSchemaPublish {
3743                    expected_revision,
3744                    schema_bundle_bytes,
3745                    schema_root_bytes,
3746                    ..
3747                } => {
3748                    let candidate = CandidateSchemaRevision::from_encoded(
3749                        schema_bundle_bytes.clone(),
3750                        schema_root_bytes.clone(),
3751                    )?;
3752                    keys.extend(self.positioned_candidate_effect_keys(
3753                        incarnation,
3754                        *expected_revision,
3755                        &candidate,
3756                        view,
3757                    )?);
3758                }
3759                JournalRecord::ConstraintValidationJobPut {
3760                    entity_tag,
3761                    constraint_id,
3762                    ..
3763                }
3764                | JournalRecord::ConstraintValidationJobDelete {
3765                    entity_tag,
3766                    constraint_id,
3767                    ..
3768                } => {
3769                    keys.insert(RawSchemaKey::from_constraint_validation_job(
3770                        *entity_tag,
3771                        *constraint_id,
3772                    ));
3773                }
3774                JournalRecord::IdentityRangeAdvance { range } => {
3775                    keys.insert(RawSchemaKey::from_identity_state(
3776                        range.owner().entity_tag(),
3777                        range.owner().field_id(),
3778                    ));
3779                }
3780                JournalRecord::RowPut { .. }
3781                | JournalRecord::RowDelete { .. }
3782                | JournalRecord::AcceptedSchemaIndexDelete { .. }
3783                | JournalRecord::AcceptedSchemaIndexPut { .. }
3784                | JournalRecord::ConstraintValidationIndexPut { .. } => {}
3785                #[cfg(any(test, feature = "migration"))]
3786                JournalRecord::SchemaMigrationRowPut { .. }
3787                | JournalRecord::SchemaMigrationIndexPut { .. } => {}
3788            }
3789        }
3790        Ok(keys)
3791    }
3792
3793    // Keep only the current entity snapshots, immutable bundle, and selected
3794    // root. The inactive root is needed only during publication and is removed
3795    // after the new root has been verified.
3796    fn retain_durable_candidate_entries(
3797        &mut self,
3798        candidate: &CandidateSchemaRevision,
3799        root_slot: usize,
3800    ) -> Result<(), InternalError> {
3801        let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3802        self.accepted_bundle_cache.get_mut().take();
3803        match &mut self.backend {
3804            SchemaStoreBackend::Heap(map) => {
3805                map.retain(|key, _| keep.contains(key) || key.is_identity_state());
3806            }
3807            SchemaStoreBackend::Journaled {
3808                canonical,
3809                live,
3810                tombstones,
3811                ..
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                live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3824                tombstones.clear();
3825            }
3826        }
3827        Ok(())
3828    }
3829
3830    fn retain_materialized_candidate_entries(
3831        &mut self,
3832        candidate: &CandidateSchemaRevision,
3833        root_slot: usize,
3834    ) -> Result<(), InternalError> {
3835        let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3836        self.accepted_bundle_cache.get_mut().take();
3837        let SchemaStoreBackend::Journaled {
3838            canonical,
3839            live,
3840            tombstones,
3841            ..
3842        } = &mut self.backend
3843        else {
3844            return Err(InternalError::store_invariant());
3845        };
3846        live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3847        let canonical_keys = canonical
3848            .iter()
3849            .map(|entry| *entry.key())
3850            .collect::<Vec<_>>();
3851        for key in canonical_keys {
3852            // Identity moves and retirements own their exact live effects;
3853            // generic catalog cleanup must preserve a moved owner's tombstone.
3854            if key.is_identity_state() {
3855                continue;
3856            }
3857            if keep.contains(&key) {
3858                tombstones.remove(&key);
3859            } else {
3860                tombstones.insert(key);
3861            }
3862        }
3863        Ok(())
3864    }
3865
3866    fn retain_canonical_candidate_entries(
3867        &mut self,
3868        keep: &BTreeSet<RawSchemaKey>,
3869    ) -> Result<(), InternalError> {
3870        self.accepted_bundle_cache.get_mut().take();
3871        let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3872            return Err(InternalError::store_invariant());
3873        };
3874        let stale = canonical
3875            .iter()
3876            .filter_map(|entry| {
3877                (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3878                    .then_some(*entry.key())
3879            })
3880            .collect::<Vec<_>>();
3881        for key in stale {
3882            canonical.remove(&key);
3883        }
3884        Ok(())
3885    }
3886
3887    /// Return whether one schema snapshot key is present.
3888    #[must_use]
3889    #[cfg(test)]
3890    fn contains_raw_snapshot(&self, key: &RawSchemaKey) -> bool {
3891        match &self.backend {
3892            SchemaStoreBackend::Heap(map) => map.contains_key(key),
3893            SchemaStoreBackend::Journaled { .. } => {
3894                self.get_raw_snapshot_for_backend(key).is_some()
3895            }
3896        }
3897    }
3898
3899    /// Return the number of schema snapshot entries in this store.
3900    #[must_use]
3901    #[cfg(test)]
3902    pub(in crate::db) fn len(&self) -> u64 {
3903        match &self.backend {
3904            SchemaStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
3905            SchemaStoreBackend::Journaled { .. } => {
3906                let mut count = 0_u64;
3907                let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3908                    count = count.saturating_add(1);
3909                    Ok(SchemaStoreVisit::Continue)
3910                });
3911                count
3912            }
3913        }
3914    }
3915
3916    /// Return whether this schema store currently has no persisted snapshots.
3917    #[must_use]
3918    #[cfg(test)]
3919    pub(in crate::db) fn is_empty(&self) -> bool {
3920        match &self.backend {
3921            SchemaStoreBackend::Heap(map) => map.is_empty(),
3922            SchemaStoreBackend::Journaled { .. } => {
3923                let mut empty = true;
3924                let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3925                    empty = false;
3926                    Ok(SchemaStoreVisit::Stop)
3927                });
3928                empty
3929            }
3930        }
3931    }
3932
3933    /// Clear all schema metadata entries from the store.
3934    #[cfg(test)]
3935    pub(in crate::db) fn clear(&mut self) {
3936        self.accepted_bundle_cache.get_mut().take();
3937        match &mut self.backend {
3938            SchemaStoreBackend::Heap(map) => map.clear(),
3939            SchemaStoreBackend::Journaled {
3940                canonical,
3941                live,
3942                tombstones,
3943                ..
3944            } => {
3945                live.clear();
3946                tombstones.clear();
3947                let keys = canonical
3948                    .iter()
3949                    .map(|entry| *entry.key())
3950                    .collect::<Vec<_>>();
3951                for key in keys {
3952                    if key.is_entity_snapshot() {
3953                        tombstones.insert(key);
3954                    } else {
3955                        canonical.remove(&key);
3956                    }
3957                }
3958            }
3959        }
3960    }
3961
3962    fn current_accepted_schema_bundle_ref(
3963        &self,
3964    ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
3965        self.current_accepted_schema_authority_ref()
3966            .map(|authority| authority.map(|(_selection, bundle)| bundle))
3967    }
3968
3969    /// Borrow the effective accepted root and its cached, verified immutable bundle together.
3970    pub(in crate::db) fn current_accepted_schema_authority_ref(
3971        &self,
3972    ) -> Result<
3973        Option<(
3974            AcceptedSchemaRootSelection,
3975            Ref<'_, AcceptedSchemaRevisionBundle>,
3976        )>,
3977        InternalError,
3978    > {
3979        let selection = self.current_accepted_schema_root()?;
3980        self.accepted_schema_authority_ref_for_selection(selection)
3981    }
3982
3983    fn accepted_schema_authority_ref_for_selection(
3984        &self,
3985        selection: Option<AcceptedSchemaRootSelection>,
3986    ) -> Result<
3987        Option<(
3988            AcceptedSchemaRootSelection,
3989            Ref<'_, AcceptedSchemaRevisionBundle>,
3990        )>,
3991        InternalError,
3992    > {
3993        let Some(selection) = selection else {
3994            self.accepted_bundle_cache
3995                .try_borrow_mut()
3996                .map_err(|_| InternalError::store_invariant())?
3997                .take();
3998            return Ok(None);
3999        };
4000
4001        let cache_matches = self
4002            .accepted_bundle_cache
4003            .try_borrow()
4004            .map_err(|_| InternalError::store_invariant())?
4005            .as_ref()
4006            .is_some_and(|cached| cached.selection == selection);
4007        if !cache_matches {
4008            let key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
4009            let raw = self
4010                .get_raw_snapshot(&key)
4011                .ok_or_else(InternalError::store_corruption)?;
4012            let bundle =
4013                decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
4014            self.validate_constraint_validation_job_closure(&bundle)?;
4015            #[cfg(test)]
4016            ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES
4017                .with(|misses| misses.set(misses.get().saturating_add(1)));
4018            let value_catalog = AcceptedValueCatalogHandle::new(
4019                bundle.enum_catalog().clone(),
4020                bundle.composite_catalog().clone(),
4021                self.accepted_catalog_scope
4022                    .get_or_init(AcceptedStoreCatalogScope::new)
4023                    .clone(),
4024                bundle.revision(),
4025                selection.root().fingerprint(),
4026            );
4027            let cardinality_domain = Rc::new(CardinalityAcceptedDomain::derive(&bundle)?);
4028            *self
4029                .accepted_bundle_cache
4030                .try_borrow_mut()
4031                .map_err(|_| InternalError::store_invariant())? = Some(AcceptedSchemaBundleCache {
4032                selection,
4033                bundle,
4034                cardinality_domain,
4035                value_catalog,
4036                entity_selections: RefCell::new(StdBTreeMap::new()),
4037            });
4038        }
4039
4040        let cache = self
4041            .accepted_bundle_cache
4042            .try_borrow()
4043            .map_err(|_| InternalError::store_invariant())?;
4044        let bundle = Ref::filter_map(cache, |cache| {
4045            cache
4046                .as_ref()
4047                .filter(|cached| cached.selection == selection)
4048                .map(|cached| &cached.bundle)
4049        })
4050        .map_err(|_| InternalError::store_invariant())?;
4051        self.validate_identity_state_closure(&bundle)?;
4052        Ok(Some((selection, bundle)))
4053    }
4054
4055    /// Reuse the accepted-domain projection for one already-selected effective root.
4056    pub(in crate::db) fn accepted_cardinality_domain_for_selection(
4057        &self,
4058        selection: Option<AcceptedSchemaRootSelection>,
4059    ) -> Result<Option<(AcceptedSchemaRootSelection, Rc<CardinalityAcceptedDomain>)>, InternalError>
4060    {
4061        let Some(selection) = selection else {
4062            self.accepted_bundle_cache
4063                .try_borrow_mut()
4064                .map_err(|_| InternalError::store_invariant())?
4065                .take();
4066            return Ok(None);
4067        };
4068        let cache_matches = self
4069            .accepted_bundle_cache
4070            .try_borrow()
4071            .map_err(|_| InternalError::store_invariant())?
4072            .as_ref()
4073            .is_some_and(|cached| cached.selection == selection);
4074        if !cache_matches {
4075            let authority = self
4076                .accepted_schema_authority_ref_for_selection(Some(selection))?
4077                .ok_or_else(InternalError::store_invariant)?;
4078            drop(authority);
4079        }
4080        let cache = self
4081            .accepted_bundle_cache
4082            .try_borrow()
4083            .map_err(|_| InternalError::store_invariant())?;
4084        let domain = cache
4085            .as_ref()
4086            .filter(|cached| cached.selection == selection)
4087            .map(|cached| Rc::clone(&cached.cardinality_domain))
4088            .ok_or_else(InternalError::store_invariant)?;
4089        Ok(Some((selection, domain)))
4090    }
4091
4092    /// Borrow the cached accepted-domain projection for an already-admitted root.
4093    pub(in crate::db) fn cached_cardinality_domain_for_root(
4094        &self,
4095        root: CardinalityAcceptedRootIdentity,
4096    ) -> Result<Option<Rc<CardinalityAcceptedDomain>>, InternalError> {
4097        let cache = self
4098            .accepted_bundle_cache
4099            .try_borrow()
4100            .map_err(|_| InternalError::store_invariant())?;
4101        Ok(cache
4102            .as_ref()
4103            .filter(|cached| root.matches(cached.selection.root()))
4104            .map(|cached| Rc::clone(&cached.cardinality_domain)))
4105    }
4106
4107    fn latest_raw_snapshots_by_entity(
4108        &self,
4109    ) -> StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)> {
4110        let mut latest_by_entity =
4111            StdBTreeMap::<EntityTag, (SchemaVersion, RawSchemaSnapshot)>::new();
4112
4113        let _: Result<(), std::convert::Infallible> = self.visit_raw_snapshots(|key, snapshot| {
4114            let version = SchemaVersion::new(key.version());
4115            match latest_by_entity.get_mut(&key.entity_tag()) {
4116                Some((latest_version, latest_snapshot)) if version > *latest_version => {
4117                    *latest_version = version;
4118                    *latest_snapshot = snapshot.clone();
4119                }
4120                None => {
4121                    latest_by_entity.insert(key.entity_tag(), (version, snapshot.clone()));
4122                }
4123                Some(_) => {}
4124            }
4125            Ok(SchemaStoreVisit::Continue)
4126        });
4127
4128        latest_by_entity
4129    }
4130
4131    /// Visit raw schema snapshots in canonical store order without exposing
4132    /// the backing stable-map iterator.
4133    fn visit_raw_snapshots<E>(
4134        &self,
4135        visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4136    ) -> Result<(), E> {
4137        let bounds = RawSchemaKey::all_entity_range_bounds();
4138        match &self.backend {
4139            SchemaStoreBackend::Heap(map) => {
4140                let mut visitor = visitor;
4141                for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4142                    if visitor(key, snapshot)?.should_stop() {
4143                        break;
4144                    }
4145                }
4146            }
4147            SchemaStoreBackend::Journaled {
4148                canonical,
4149                live,
4150                tombstones,
4151                ..
4152            } => Self::visit_journaled_raw_snapshot_range(
4153                canonical,
4154                live,
4155                tombstones,
4156                bounds,
4157                Direction::Asc,
4158                visitor,
4159            )?,
4160        }
4161
4162        Ok(())
4163    }
4164
4165    fn visit_constraint_validation_jobs_in_view<E>(
4166        &self,
4167        view: IdentityStateStorageView,
4168        visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4169    ) -> Result<(), E> {
4170        let bounds = RawSchemaKey::all_constraint_validation_job_range_bounds();
4171        match (&self.backend, view) {
4172            (SchemaStoreBackend::Heap(map), _) => {
4173                let mut visitor = visitor;
4174                for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4175                    if visitor(key, snapshot)?.should_stop() {
4176                        break;
4177                    }
4178                }
4179            }
4180            (
4181                SchemaStoreBackend::Journaled {
4182                    canonical,
4183                    live,
4184                    tombstones,
4185                    ..
4186                },
4187                IdentityStateStorageView::Effective,
4188            ) => Self::visit_journaled_raw_snapshot_range(
4189                canonical,
4190                live,
4191                tombstones,
4192                bounds,
4193                Direction::Asc,
4194                visitor,
4195            )?,
4196            (
4197                SchemaStoreBackend::Journaled { canonical, .. },
4198                IdentityStateStorageView::Canonical,
4199            ) => {
4200                let mut visitor = visitor;
4201                for entry in canonical.range((bounds.0, bounds.1)) {
4202                    if visitor(entry.key(), &entry.value())?.should_stop() {
4203                        break;
4204                    }
4205                }
4206            }
4207        }
4208        Ok(())
4209    }
4210
4211    #[cfg(test)]
4212    #[must_use]
4213    pub(in crate::db) fn canonical_len_for_tests(&self) -> u64 {
4214        match &self.backend {
4215            SchemaStoreBackend::Journaled { canonical: map, .. } => map.len(),
4216            SchemaStoreBackend::Heap(_) => 0,
4217        }
4218    }
4219
4220    fn get_raw_snapshot_for_backend(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
4221        let SchemaStoreBackend::Journaled {
4222            canonical,
4223            live,
4224            tombstones,
4225            ..
4226        } = &self.backend
4227        else {
4228            return None;
4229        };
4230
4231        if tombstones.contains(key) {
4232            return None;
4233        }
4234        live.get(key).cloned().or_else(|| canonical.get(key))
4235    }
4236
4237    fn visit_journaled_raw_snapshot_range<E>(
4238        canonical: &StableBTreeMap<
4239            RawSchemaKey,
4240            RawSchemaSnapshot,
4241            RuntimeMemory<DefaultMemoryImpl>,
4242        >,
4243        live: &StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
4244        tombstones: &BTreeSet<RawSchemaKey>,
4245        bounds: (RangeBound<RawSchemaKey>, RangeBound<RawSchemaKey>),
4246        direction: Direction,
4247        mut visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4248    ) -> Result<(), E> {
4249        match direction {
4250            Direction::Asc => {
4251                for entry in ordered_overlay_entries(
4252                    canonical.range((bounds.0, bounds.1)),
4253                    live.range((bounds.0, bounds.1)),
4254                    Direction::Asc,
4255                    |entry| entry.key(),
4256                    |entry| entry.0,
4257                    tombstones,
4258                ) {
4259                    let visit = match entry {
4260                        OrderedOverlayEntry::Canonical(canonical_entry) => {
4261                            visitor(canonical_entry.key(), &canonical_entry.value())?
4262                        }
4263                        OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4264                    };
4265                    if visit.should_stop() {
4266                        return Ok(());
4267                    }
4268                }
4269            }
4270            Direction::Desc => {
4271                for entry in ordered_overlay_entries(
4272                    canonical.range((bounds.0, bounds.1)).rev(),
4273                    live.range((bounds.0, bounds.1)).rev(),
4274                    Direction::Desc,
4275                    |entry| entry.key(),
4276                    |entry| entry.0,
4277                    tombstones,
4278                ) {
4279                    let visit = match entry {
4280                        OrderedOverlayEntry::Canonical(canonical_entry) => {
4281                            visitor(canonical_entry.key(), &canonical_entry.value())?
4282                        }
4283                        OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4284                    };
4285                    if visit.should_stop() {
4286                        return Ok(());
4287                    }
4288                }
4289            }
4290        }
4291
4292        Ok(())
4293    }
4294}
4295
4296fn map_schema_publication_error(error: AcceptedSchemaPublicationError) -> InternalError {
4297    match error {
4298        AcceptedSchemaPublicationError::StaleSchemaRevision { .. }
4299        | AcceptedSchemaPublicationError::RevisionExhausted => InternalError::store_unsupported(),
4300        AcceptedSchemaPublicationError::InvalidCandidate => InternalError::store_invariant(),
4301        AcceptedSchemaPublicationError::CorruptRootSlots => InternalError::store_corruption(),
4302    }
4303}
4304
4305fn derive_data_allocation_metadata(
4306    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4307) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4308    let mut max_version = SchemaVersion::initial();
4309    let mut hasher = new_hash_sha256();
4310    write_hash_tag_u8(&mut hasher, SCHEMA_STORE_DATA_ALLOCATION_FINGERPRINT_DOMAIN);
4311
4312    for (entity, (_, snapshot)) in latest_by_entity {
4313        let persisted = snapshot.decode_persisted_snapshot()?;
4314        if persisted.version() > max_version {
4315            max_version = persisted.version();
4316        }
4317
4318        let data_projection = PersistedSchemaSnapshot::new_with_primary_key_fields_and_indexes(
4319            persisted.version(),
4320            persisted.entity_path().to_string(),
4321            persisted.entity_name().to_string(),
4322            persisted.primary_key_field_ids().to_vec(),
4323            persisted.row_layout().clone(),
4324            persisted.fields().to_vec(),
4325            Vec::new(),
4326        );
4327        let constraint_catalog = crate::db::schema::AcceptedConstraintCatalog::initial(
4328            data_projection.fields(),
4329            data_projection.indexes(),
4330            data_projection.relations(),
4331        )
4332        .map_err(|_| InternalError::store_invariant())?;
4333        let data_projection = data_projection.with_constraint_catalog(constraint_catalog);
4334        let encoded = encode_persisted_schema_snapshot(&data_projection)?;
4335
4336        write_hash_u64(&mut hasher, entity.value());
4337        write_hash_u32(&mut hasher, persisted.version().get());
4338        write_hash_len_u32(&mut hasher, encoded.len());
4339        hasher.update(encoded);
4340    }
4341
4342    Ok(finalize_schema_metadata(
4343        max_version,
4344        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4345        hasher,
4346        latest_by_entity.len(),
4347    ))
4348}
4349
4350fn derive_index_allocation_metadata(
4351    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4352) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4353    let mut max_version = SchemaVersion::initial();
4354    let mut hasher = new_hash_sha256();
4355    write_hash_tag_u8(
4356        &mut hasher,
4357        SCHEMA_STORE_INDEX_ALLOCATION_FINGERPRINT_DOMAIN,
4358    );
4359
4360    for (entity, (_, snapshot)) in latest_by_entity {
4361        let persisted = snapshot.decode_persisted_snapshot()?;
4362        if persisted.version() > max_version {
4363            max_version = persisted.version();
4364        }
4365
4366        write_hash_u64(&mut hasher, entity.value());
4367        write_hash_u32(&mut hasher, persisted.version().get());
4368        write_hash_len_u32(&mut hasher, persisted.indexes().len());
4369        for index in persisted.indexes() {
4370            write_hash_u32(&mut hasher, u32::from(index.ordinal()));
4371            write_hash_str_u32(&mut hasher, index.name());
4372            write_hash_str_u32(&mut hasher, index.store());
4373            write_hash_tag_u8(&mut hasher, u8::from(index.unique()));
4374            write_hash_str_u32(&mut hasher, persisted_index_origin_name(index.origin()));
4375            match index.predicate_sql() {
4376                Some(predicate_sql) => {
4377                    write_hash_tag_u8(&mut hasher, 1);
4378                    write_hash_str_u32(&mut hasher, predicate_sql);
4379                }
4380                None => write_hash_tag_u8(&mut hasher, 0),
4381            }
4382            hash_persisted_index_key(&mut hasher, index.key());
4383        }
4384    }
4385
4386    Ok(finalize_schema_metadata(
4387        max_version,
4388        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4389        hasher,
4390        latest_by_entity.len(),
4391    ))
4392}
4393
4394fn derive_schema_catalog_metadata(
4395    latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4396) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4397    let mut max_version = SchemaVersion::initial();
4398    let mut hasher = new_hash_sha256();
4399    write_hash_tag_u8(&mut hasher, SCHEMA_STORE_CATALOG_FINGERPRINT_DOMAIN);
4400
4401    for (entity, (version, snapshot)) in latest_by_entity {
4402        let persisted = snapshot.decode_persisted_snapshot()?;
4403        if persisted.version() > max_version {
4404            max_version = persisted.version();
4405        }
4406
4407        write_hash_u64(&mut hasher, entity.value());
4408        write_hash_u32(&mut hasher, version.get());
4409        write_hash_len_u32(&mut hasher, snapshot.as_bytes().len());
4410        hasher.update(snapshot.as_bytes());
4411    }
4412
4413    Ok(finalize_schema_metadata(
4414        max_version,
4415        SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4416        hasher,
4417        latest_by_entity.len(),
4418    ))
4419}
4420
4421fn finalize_schema_metadata(
4422    schema_version: SchemaVersion,
4423    schema_fingerprint_method_version: u8,
4424    hasher: sha2::Sha256,
4425    entity_count: usize,
4426) -> SchemaStoreCatalogMetadata {
4427    let digest = finalize_hash_sha256(hasher);
4428    let mut schema_fingerprint = [0u8; 16];
4429    schema_fingerprint.copy_from_slice(&digest[..16]);
4430
4431    SchemaStoreCatalogMetadata::new(
4432        schema_version,
4433        schema_fingerprint_method_version,
4434        schema_fingerprint,
4435        u64::try_from(entity_count).unwrap_or(u64::MAX),
4436    )
4437}
4438
4439fn hash_persisted_index_key(hasher: &mut sha2::Sha256, key: &PersistedIndexKeySnapshot) {
4440    match key {
4441        PersistedIndexKeySnapshot::FieldPath(paths) => {
4442            write_hash_tag_u8(hasher, 1);
4443            write_hash_len_u32(hasher, paths.len());
4444            for path in paths {
4445                hash_persisted_index_field_path(hasher, path);
4446            }
4447        }
4448        PersistedIndexKeySnapshot::Items(items) => {
4449            write_hash_tag_u8(hasher, 2);
4450            write_hash_len_u32(hasher, items.len());
4451            for item in items {
4452                match item {
4453                    PersistedIndexKeyItemSnapshot::FieldPath(path) => {
4454                        write_hash_tag_u8(hasher, 1);
4455                        hash_persisted_index_field_path(hasher, path);
4456                    }
4457                    PersistedIndexKeyItemSnapshot::Expression(expression) => {
4458                        write_hash_tag_u8(hasher, 2);
4459                        write_hash_str_u32(hasher, persisted_expression_op_name(expression.op()));
4460                        hash_persisted_index_field_path(hasher, expression.source());
4461                        hash_accepted_field_kind(hasher, expression.input_kind());
4462                        hash_accepted_field_kind(hasher, expression.output_kind());
4463                        write_hash_str_u32(hasher, expression.canonical_text());
4464                    }
4465                }
4466            }
4467        }
4468    }
4469}
4470
4471fn hash_persisted_index_field_path(
4472    hasher: &mut sha2::Sha256,
4473    path: &crate::db::schema::PersistedIndexFieldPathSnapshot,
4474) {
4475    write_hash_u32(hasher, path.field_id().get());
4476    write_hash_u32(hasher, u32::from(path.slot().get()));
4477    write_hash_len_u32(hasher, path.path().len());
4478    for segment in path.path() {
4479        write_hash_str_u32(hasher, segment);
4480    }
4481    hash_accepted_field_kind(hasher, path.kind());
4482    write_hash_tag_u8(hasher, u8::from(path.nullable()));
4483}
4484
4485fn hash_accepted_field_kind(hasher: &mut sha2::Sha256, kind: &AcceptedFieldKind) {
4486    match kind {
4487        AcceptedFieldKind::Account => write_hash_tag_u8(hasher, 1),
4488        AcceptedFieldKind::Blob { max_len } => {
4489            write_hash_tag_u8(hasher, 2);
4490            hash_optional_u32(hasher, *max_len);
4491        }
4492        AcceptedFieldKind::Bool => {
4493            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_BOOL);
4494        }
4495        AcceptedFieldKind::Date => write_hash_tag_u8(hasher, 4),
4496        AcceptedFieldKind::Decimal { scale } => {
4497            write_hash_tag_u8(hasher, 5);
4498            write_hash_u32(hasher, *scale);
4499        }
4500        AcceptedFieldKind::Duration => write_hash_tag_u8(hasher, 6),
4501        AcceptedFieldKind::Enum { type_id } => {
4502            write_hash_tag_u8(hasher, 7);
4503            write_hash_u32(hasher, type_id.get());
4504        }
4505        AcceptedFieldKind::Float32 => write_hash_tag_u8(hasher, 8),
4506        AcceptedFieldKind::Float64 => write_hash_tag_u8(hasher, 9),
4507        AcceptedFieldKind::Int8 => write_hash_tag_u8(hasher, 10),
4508        AcceptedFieldKind::Int16 => write_hash_tag_u8(hasher, 11),
4509        AcceptedFieldKind::Int32 => write_hash_tag_u8(hasher, 12),
4510        AcceptedFieldKind::Int64 => write_hash_tag_u8(hasher, 13),
4511        AcceptedFieldKind::Int128 => write_hash_tag_u8(hasher, 14),
4512        AcceptedFieldKind::IntBig { max_bytes } => {
4513            write_hash_tag_u8(hasher, 15);
4514            write_hash_u32(hasher, *max_bytes);
4515        }
4516        AcceptedFieldKind::Principal => write_hash_tag_u8(hasher, 16),
4517        AcceptedFieldKind::Subaccount => write_hash_tag_u8(hasher, 17),
4518        AcceptedFieldKind::Text { max_len } => {
4519            write_hash_tag_u8(hasher, 18);
4520            hash_optional_u32(hasher, *max_len);
4521        }
4522        AcceptedFieldKind::Timestamp => write_hash_tag_u8(hasher, 19),
4523        AcceptedFieldKind::Nat8 => write_hash_tag_u8(hasher, 20),
4524        AcceptedFieldKind::Nat16 => write_hash_tag_u8(hasher, 21),
4525        AcceptedFieldKind::Nat32 => write_hash_tag_u8(hasher, 22),
4526        AcceptedFieldKind::Nat64 => write_hash_tag_u8(hasher, 23),
4527        AcceptedFieldKind::Nat128 => write_hash_tag_u8(hasher, 24),
4528        AcceptedFieldKind::NatBig { max_bytes } => {
4529            write_hash_tag_u8(hasher, 25);
4530            write_hash_u32(hasher, *max_bytes);
4531        }
4532        AcceptedFieldKind::Ulid => write_hash_tag_u8(hasher, 26),
4533        AcceptedFieldKind::Unit => write_hash_tag_u8(hasher, 27),
4534        AcceptedFieldKind::Relation {
4535            target_path,
4536            target_entity_name,
4537            target_entity_tag,
4538            target_store_path,
4539            key_kind,
4540        } => {
4541            write_hash_tag_u8(hasher, 28);
4542            write_hash_str_u32(hasher, target_path);
4543            write_hash_str_u32(hasher, target_entity_name);
4544            write_hash_u64(hasher, target_entity_tag.value());
4545            write_hash_str_u32(hasher, target_store_path);
4546            hash_accepted_field_kind(hasher, key_kind);
4547        }
4548        AcceptedFieldKind::List(inner) => {
4549            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_LIST);
4550            hash_accepted_field_kind(hasher, inner);
4551        }
4552        AcceptedFieldKind::Set(inner) => {
4553            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_SET);
4554            hash_accepted_field_kind(hasher, inner);
4555        }
4556        AcceptedFieldKind::Map { key, value } => {
4557            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_MAP);
4558            hash_accepted_field_kind(hasher, key);
4559            hash_accepted_field_kind(hasher, value);
4560        }
4561        AcceptedFieldKind::Composite { type_id } => {
4562            write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_COMPOSITE);
4563            write_hash_u32(hasher, type_id.get());
4564        }
4565        AcceptedFieldKind::U256 => write_hash_tag_u8(hasher, 33),
4566    }
4567}
4568
4569fn hash_optional_u32(hasher: &mut sha2::Sha256, value: Option<u32>) {
4570    match value {
4571        Some(value) => {
4572            write_hash_tag_u8(hasher, 1);
4573            write_hash_u32(hasher, value);
4574        }
4575        None => write_hash_tag_u8(hasher, 0),
4576    }
4577}
4578
4579const fn persisted_index_origin_name(
4580    origin: crate::db::schema::PersistedIndexOrigin,
4581) -> &'static str {
4582    match origin {
4583        crate::db::schema::PersistedIndexOrigin::Generated => "generated",
4584        crate::db::schema::PersistedIndexOrigin::SqlDdl => "sql_ddl",
4585    }
4586}
4587
4588const fn persisted_expression_op_name(
4589    op: crate::db::schema::PersistedIndexExpressionOp,
4590) -> &'static str {
4591    match op {
4592        crate::db::schema::PersistedIndexExpressionOp::Lower => "lower",
4593        crate::db::schema::PersistedIndexExpressionOp::Upper => "upper",
4594        crate::db::schema::PersistedIndexExpressionOp::Trim => "trim",
4595        crate::db::schema::PersistedIndexExpressionOp::LowerTrim => "lower_trim",
4596        crate::db::schema::PersistedIndexExpressionOp::Date => "date",
4597        crate::db::schema::PersistedIndexExpressionOp::Year => "year",
4598        crate::db::schema::PersistedIndexExpressionOp::Month => "month",
4599        crate::db::schema::PersistedIndexExpressionOp::Day => "day",
4600    }
4601}
4602
4603///
4604/// TESTS
4605///
4606
4607#[cfg(test)]
4608mod tests;