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