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