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 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2224 return Ok(None);
2225 };
2226 self.validate_constraint_validation_job_closure(&bundle)?;
2227 Ok(Some(bundle.clone()))
2228 }
2229
2230 pub(in crate::db) fn current_accepted_runtime_entities(
2232 &self,
2233 registered_store_path: &'static str,
2234 ) -> Result<Vec<AcceptedRuntimeEntity>, InternalError> {
2235 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2236 return Ok(Vec::new());
2237 };
2238 if bundle.store_path() != registered_store_path {
2239 return Err(InternalError::store_corruption());
2240 }
2241
2242 bundle
2243 .entity_snapshots()
2244 .iter()
2245 .map(|(entity_tag, snapshot)| {
2246 AcceptedRuntimeEntity::from_accepted_snapshot(
2247 &bundle,
2248 *entity_tag,
2249 snapshot,
2250 registered_store_path,
2251 )
2252 })
2253 .collect()
2254 }
2255
2256 pub(in crate::db) fn current_accepted_runtime_entity_for_tag(
2258 &self,
2259 registered_store_path: &'static str,
2260 entity_tag: EntityTag,
2261 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2262 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2263 return Ok(None);
2264 };
2265 if bundle.store_path() != registered_store_path {
2266 return Err(InternalError::store_corruption());
2267 }
2268 let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2269 return Ok(None);
2270 };
2271
2272 AcceptedRuntimeEntity::from_accepted_snapshot(
2273 &bundle,
2274 entity_tag,
2275 snapshot,
2276 registered_store_path,
2277 )
2278 .map(Some)
2279 }
2280
2281 pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_tag(
2283 &self,
2284 registered_store_path: &'static str,
2285 entity_tag: EntityTag,
2286 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2287 let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2288 return Ok(None);
2289 };
2290 if bundle.store_path() != registered_store_path {
2291 return Err(InternalError::store_corruption());
2292 }
2293 let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2294 return Ok(None);
2295 };
2296
2297 AcceptedRuntimeEntity::from_accepted_snapshot(
2298 &bundle,
2299 entity_tag,
2300 snapshot,
2301 registered_store_path,
2302 )
2303 .map(Some)
2304 }
2305
2306 pub(in crate::db) fn current_accepted_runtime_entity_for_path(
2308 &self,
2309 registered_store_path: &'static str,
2310 entity_path: &str,
2311 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2312 self.current_accepted_runtime_entity_matching(registered_store_path, |snapshot_path, _| {
2313 snapshot_path == entity_path
2314 })
2315 }
2316
2317 pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_path(
2319 &self,
2320 registered_store_path: &'static str,
2321 entity_path: &str,
2322 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2323 let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2324 return Ok(None);
2325 };
2326 if bundle.store_path() != registered_store_path {
2327 return Err(InternalError::store_corruption());
2328 }
2329
2330 let mut matched = None;
2331 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2332 if snapshot.entity_path() != entity_path {
2333 continue;
2334 }
2335 let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2336 &bundle,
2337 *entity_tag,
2338 snapshot,
2339 registered_store_path,
2340 )?;
2341 if matched.replace(entity).is_some() {
2342 return Err(InternalError::store_corruption());
2343 }
2344 }
2345
2346 Ok(matched)
2347 }
2348
2349 #[cfg(test)]
2351 pub(in crate::db) fn current_accepted_runtime_entity_for_name(
2352 &self,
2353 registered_store_path: &'static str,
2354 entity_name: &str,
2355 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2356 self.current_accepted_runtime_entity_matching(registered_store_path, |_, snapshot_name| {
2357 snapshot_name == entity_name
2358 })
2359 }
2360
2361 fn current_accepted_runtime_entity_matching(
2362 &self,
2363 registered_store_path: &'static str,
2364 mut predicate: impl FnMut(&str, &str) -> bool,
2365 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2366 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2367 return Ok(None);
2368 };
2369 if bundle.store_path() != registered_store_path {
2370 return Err(InternalError::store_corruption());
2371 }
2372
2373 let mut matched = None;
2374 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2375 if !predicate(snapshot.entity_path(), snapshot.entity_name()) {
2376 continue;
2377 }
2378 let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2379 &bundle,
2380 *entity_tag,
2381 snapshot,
2382 registered_store_path,
2383 )?;
2384 if matched.replace(entity).is_some() {
2385 return Err(InternalError::store_corruption());
2386 }
2387 }
2388
2389 Ok(matched)
2390 }
2391
2392 pub(in crate::db) fn current_accepted_schema_revision(
2394 &self,
2395 ) -> Result<Option<AcceptedSchemaRevision>, InternalError> {
2396 Ok(self
2397 .current_accepted_schema_root()?
2398 .map(|selection| selection.root().revision()))
2399 }
2400
2401 pub(in crate::db) fn pending_relation_activation_for_target(
2407 &self,
2408 target_path: &str,
2409 ) -> Result<Option<PendingRelationActivationDeleteBarrier>, InternalError> {
2410 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2411 return Ok(None);
2412 };
2413 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2414 let Some(candidate) = snapshot
2415 .candidate_relations()
2416 .iter()
2417 .find(|candidate| candidate.target_path() == target_path)
2418 else {
2419 continue;
2420 };
2421 let activation = snapshot
2422 .constraint_activations()
2423 .iter()
2424 .find(|activation| {
2425 matches!(
2426 activation.kind(),
2427 ConstraintActivationKind::Relation { relation_id }
2428 if *relation_id == candidate.id()
2429 )
2430 })
2431 .ok_or_else(InternalError::store_corruption)?;
2432 return Ok(Some(PendingRelationActivationDeleteBarrier {
2433 accepted_schema_fingerprint:
2434 accepted_schema_cache_fingerprint_for_persisted_snapshot(snapshot)?,
2435 source_entity_tag: *entity_tag,
2436 constraint_id: activation.id(),
2437 }));
2438 }
2439
2440 Ok(None)
2441 }
2442
2443 pub(in crate::db) fn entity_has_relation_to_target(
2445 &self,
2446 source_entity: EntityTag,
2447 target_path: &str,
2448 ) -> Result<bool, InternalError> {
2449 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2450 return Ok(false);
2451 };
2452 let Some(snapshot) = bundle.entity_snapshots().get(&source_entity) else {
2453 return Ok(false);
2454 };
2455
2456 Ok(snapshot
2457 .relations()
2458 .iter()
2459 .any(|relation| relation.target_path() == target_path))
2460 }
2461
2462 pub(in crate::db) fn validate_live_activation_transition(
2464 &self,
2465 candidate: &AcceptedSchemaRevisionBundle,
2466 ) -> Result<(), InternalError> {
2467 let Some(current) = self.current_accepted_schema_bundle()? else {
2468 return Ok(());
2469 };
2470 Self::validate_activation_transition_from(¤t, candidate)
2471 }
2472
2473 pub(in crate::db) fn validate_canonical_activation_transition(
2475 &self,
2476 candidate: &AcceptedSchemaRevisionBundle,
2477 ) -> Result<(), InternalError> {
2478 let Some(current) = self.current_canonical_accepted_schema_bundle()? else {
2479 return Ok(());
2480 };
2481 Self::validate_activation_transition_from(¤t, candidate)
2482 }
2483
2484 fn validate_activation_transition_from(
2485 current: &AcceptedSchemaRevisionBundle,
2486 candidate: &AcceptedSchemaRevisionBundle,
2487 ) -> Result<(), InternalError> {
2488 for (entity_tag, before) in current.entity_snapshots() {
2489 if before.constraint_activations().is_empty() {
2490 continue;
2491 }
2492 let after = candidate
2493 .entity_snapshots()
2494 .get(entity_tag)
2495 .ok_or_else(InternalError::store_invariant)?;
2496 if before == after {
2497 continue;
2498 }
2499 let expected_shape = before
2500 .clone()
2501 .with_constraint_catalog(after.constraint_catalog().clone());
2502 let catalog_only_transition = expected_shape == *after
2503 && before
2504 .constraint_catalog()
2505 .permits_live_activation_transition_to(after.constraint_catalog());
2506 let sql_row_local_abort_with_version =
2507 before.constraint_activations().iter().any(|activation| {
2508 activation.origin() == ConstraintOrigin::SqlDdl
2509 && matches!(
2510 activation.kind(),
2511 ConstraintActivationKind::Check { .. }
2512 | ConstraintActivationKind::NotNull { .. }
2513 )
2514 && before.version().get().checked_add(1) == Some(after.version().get())
2515 && before
2516 .constraint_catalog()
2517 .clone()
2518 .with_aborted_activation(activation.id())
2519 .is_ok_and(|catalog| catalog == *after.constraint_catalog())
2520 && before
2521 .clone()
2522 .with_constraint_catalog(after.constraint_catalog().clone())
2523 .with_schema_version(after.version())
2524 == *after
2525 });
2526 let sql_unique_abort_with_version =
2527 before.constraint_activations().iter().any(|activation| {
2528 activation.origin() == ConstraintOrigin::SqlDdl
2529 && matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2530 && before.version().get().checked_add(1) == Some(after.version().get())
2531 && before
2532 .with_aborted_unique_activation(activation.id(), after.version())
2533 .is_ok_and(|expected| expected == *after)
2534 });
2535 let not_null_promotion = before.constraint_activations().iter().any(|activation| {
2536 matches!(activation.kind(), ConstraintActivationKind::NotNull { .. })
2537 && before
2538 .with_promoted_not_null_activation(activation.id(), after.version())
2539 .is_ok_and(|expected| expected == *after)
2540 });
2541 let unique_promotion = before.constraint_activations().iter().any(|activation| {
2542 matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2543 && before
2544 .with_promoted_unique_activation(activation.id(), after.version())
2545 .is_ok_and(|expected| expected == *after)
2546 });
2547 let relation_promotion = before.constraint_activations().iter().any(|activation| {
2548 matches!(activation.kind(), ConstraintActivationKind::Relation { .. })
2549 && before
2550 .with_promoted_relation_activation(activation.id(), after.version())
2551 .is_ok_and(|expected| expected == *after)
2552 });
2553 if !catalog_only_transition
2554 && !sql_row_local_abort_with_version
2555 && !sql_unique_abort_with_version
2556 && !not_null_promotion
2557 && !unique_promotion
2558 && !relation_promotion
2559 {
2560 return Err(InternalError::store_invariant());
2561 }
2562 }
2563 Ok(())
2564 }
2565
2566 pub(in crate::db) fn validate_constraint_validation_job_closure(
2568 &self,
2569 bundle: &AcceptedSchemaRevisionBundle,
2570 ) -> Result<(), InternalError> {
2571 self.validate_constraint_validation_job_closure_with_change(bundle, None, None)
2572 }
2573
2574 pub(in crate::db) fn validate_constraint_validation_job_closure_with_change(
2577 &self,
2578 bundle: &AcceptedSchemaRevisionBundle,
2579 replacement: Option<&ConstraintValidationJob>,
2580 removal: Option<(EntityTag, ConstraintId)>,
2581 ) -> Result<(), InternalError> {
2582 self.validate_constraint_validation_job_closure_with_change_in_view(
2583 bundle,
2584 replacement,
2585 removal,
2586 IdentityStateStorageView::Effective,
2587 )
2588 }
2589
2590 pub(in crate::db) fn validate_canonical_constraint_validation_job_closure_with_change(
2592 &self,
2593 bundle: &AcceptedSchemaRevisionBundle,
2594 replacement: Option<&ConstraintValidationJob>,
2595 removal: Option<(EntityTag, ConstraintId)>,
2596 ) -> Result<(), InternalError> {
2597 self.validate_constraint_validation_job_closure_with_change_in_view(
2598 bundle,
2599 replacement,
2600 removal,
2601 IdentityStateStorageView::Canonical,
2602 )
2603 }
2604
2605 fn validate_constraint_validation_job_closure_with_change_in_view(
2606 &self,
2607 bundle: &AcceptedSchemaRevisionBundle,
2608 replacement: Option<&ConstraintValidationJob>,
2609 removal: Option<(EntityTag, ConstraintId)>,
2610 view: IdentityStateStorageView,
2611 ) -> Result<(), InternalError> {
2612 if replacement.is_some() && removal.is_some() {
2613 return Err(InternalError::store_invariant());
2614 }
2615 let replacement_key = replacement.map(|job| {
2616 RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id())
2617 });
2618 let removal_key = removal.map(|(entity_tag, constraint_id)| {
2619 RawSchemaKey::from_constraint_validation_job(entity_tag, constraint_id)
2620 });
2621 let mut expected = BTreeSet::new();
2622 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2623 for activation in snapshot.constraint_activations() {
2624 let key =
2625 RawSchemaKey::from_constraint_validation_job(*entity_tag, activation.id());
2626 match activation.state() {
2627 ConstraintActivationState::EnforcingNewWrites => {
2628 if self
2629 .constraint_validation_job_after_change(
2630 key,
2631 replacement,
2632 replacement_key,
2633 removal_key,
2634 view,
2635 )?
2636 .is_some()
2637 {
2638 return Err(InternalError::store_corruption());
2639 }
2640 }
2641 ConstraintActivationState::Validating => {
2642 let job = self
2643 .constraint_validation_job_after_change(
2644 key,
2645 replacement,
2646 replacement_key,
2647 removal_key,
2648 view,
2649 )?
2650 .ok_or_else(InternalError::store_corruption)?;
2651 if job.entity_tag() != *entity_tag
2652 || job.entity_path() != snapshot.entity_path()
2653 {
2654 return Err(InternalError::store_corruption());
2655 }
2656 job.validate(Some(activation))?;
2657 expected.insert(key);
2658 }
2659 }
2660 }
2661 }
2662
2663 self.visit_constraint_validation_jobs_in_view(view, |key, raw| {
2664 if removal_key == Some(*key) || replacement_key == Some(*key) {
2665 return Ok(SchemaStoreVisit::Continue);
2666 }
2667 if !expected.contains(key) {
2668 return Err(InternalError::store_corruption());
2669 }
2670 let job = decode_constraint_validation_job(raw.as_bytes())?;
2671 if job.entity_tag() != key.entity_tag()
2672 || key.constraint_id() != Some(job.constraint_id())
2673 {
2674 return Err(InternalError::store_corruption());
2675 }
2676 Ok(SchemaStoreVisit::Continue)
2677 })?;
2678
2679 if let Some(key) = replacement_key
2680 && !expected.contains(&key)
2681 {
2682 return Err(InternalError::store_corruption());
2683 }
2684 if let Some(key) = removal_key
2685 && expected.contains(&key)
2686 {
2687 return Err(InternalError::store_corruption());
2688 }
2689
2690 Ok(())
2691 }
2692
2693 fn constraint_validation_job_after_change(
2694 &self,
2695 key: RawSchemaKey,
2696 replacement: Option<&ConstraintValidationJob>,
2697 replacement_key: Option<RawSchemaKey>,
2698 removal_key: Option<RawSchemaKey>,
2699 view: IdentityStateStorageView,
2700 ) -> Result<Option<ConstraintValidationJob>, InternalError> {
2701 if removal_key == Some(key) {
2702 return Ok(None);
2703 }
2704 if replacement_key == Some(key) {
2705 return Ok(replacement.cloned());
2706 }
2707 let raw = match view {
2708 IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
2709 IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
2710 };
2711 raw.map(|raw| decode_constraint_validation_job(raw.as_bytes()))
2712 .transpose()
2713 }
2714
2715 pub(in crate::db) fn current_accepted_schema_authority_matches(
2718 &self,
2719 expected: &AcceptedSchemaAuthority,
2720 ) -> Result<bool, InternalError> {
2721 let Some(store_scope) = self.accepted_catalog_scope.get() else {
2722 return Ok(false);
2723 };
2724
2725 if let Some(cached) = self
2728 .accepted_bundle_cache
2729 .try_borrow()
2730 .map_err(|_| InternalError::store_invariant())?
2731 .as_ref()
2732 {
2733 let root = cached.selection.root();
2734 return Ok(expected.matches_store_root(
2735 store_scope,
2736 root.revision(),
2737 root.fingerprint(),
2738 ));
2739 }
2740
2741 let Some(selection) = self.current_accepted_schema_root()? else {
2742 return Ok(false);
2743 };
2744 let root = selection.root();
2745
2746 Ok(expected.matches_store_root(store_scope, root.revision(), root.fingerprint()))
2747 }
2748
2749 pub(in crate::db) fn publish_accepted_schema_candidate(
2755 &mut self,
2756 incarnation: DatabaseIncarnationId,
2757 expected_revision: AcceptedSchemaRevision,
2758 candidate: &CandidateSchemaRevision,
2759 ) -> Result<(), InternalError> {
2760 let identity_transition = self.prepare_identity_state_transition(
2761 incarnation,
2762 candidate,
2763 IdentityStateStorageView::Effective,
2764 )?;
2765 if self.current_root_matches_candidate(candidate)? {
2766 if !identity_transition.is_empty() {
2767 return Err(InternalError::identity_state_corruption());
2768 }
2769 let selection = self
2770 .current_accepted_schema_root()?
2771 .ok_or_else(InternalError::store_corruption)?;
2772 self.retain_durable_candidate_entries(candidate, selection.slot())?;
2773 return Ok(());
2774 }
2775 let first = self.accepted_root_slot_bytes(0)?;
2776 let second = self.accepted_root_slot_bytes(1)?;
2777 prepare_accepted_schema_root_publication(
2778 [first.as_deref(), second.as_deref()],
2779 expected_revision,
2780 candidate,
2781 )
2782 .map_err(map_schema_publication_error)?;
2783
2784 self.insert_durable_candidate_snapshots(candidate)?;
2785 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2786 self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2787 let persisted_bundle = self
2788 .get_raw_snapshot(&bundle_key)
2789 .ok_or_else(InternalError::store_corruption)?;
2790 let _verified = decode_verified_accepted_schema_revision_bundle(
2791 candidate.root(),
2792 persisted_bundle.as_bytes(),
2793 )?;
2794 self.apply_identity_state_transition(
2795 identity_transition,
2796 IdentityStateWriteTarget::Durable,
2797 )?;
2798
2799 let first = self.accepted_root_slot_bytes(0)?;
2802 let second = self.accepted_root_slot_bytes(1)?;
2803 let publication = prepare_accepted_schema_root_publication(
2804 [first.as_deref(), second.as_deref()],
2805 expected_revision,
2806 candidate,
2807 )
2808 .map_err(map_schema_publication_error)?;
2809 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
2810 self.insert_durable_raw_value(root_key, publication.encoded_root().to_vec());
2811
2812 let selected = self
2813 .current_accepted_schema_root()?
2814 .ok_or_else(InternalError::store_corruption)?;
2815 if selected.root() != candidate.root() {
2816 return Err(InternalError::store_corruption());
2817 }
2818 self.retain_durable_candidate_entries(candidate, selected.slot())?;
2819 Ok(())
2820 }
2821
2822 pub(in crate::db) fn restore_live_accepted_schema_checkpoint(
2825 &mut self,
2826 incarnation: DatabaseIncarnationId,
2827 candidate: &CandidateSchemaRevision,
2828 checkpoint_identity_states: &IdentityStateInventory,
2829 ) -> Result<(), InternalError> {
2830 if !matches!(self.backend, SchemaStoreBackend::Heap(_)) {
2831 return Err(InternalError::store_invariant());
2832 }
2833 let checkpoint_validation = prepare_identity_state_transition(
2834 incarnation,
2835 Some(candidate.bundle()),
2836 candidate.bundle(),
2837 checkpoint_identity_states.clone(),
2838 )?;
2839 if !checkpoint_validation.is_empty() {
2840 return Err(InternalError::identity_state_corruption());
2841 }
2842 if self.current_root_matches_candidate(candidate)? {
2843 for state in checkpoint_identity_states.values() {
2844 let key = RawSchemaKey::from_identity_state(
2845 state.owner().entity_tag(),
2846 state.owner().field_id(),
2847 );
2848 self.insert_durable_raw_value(key, encode_identity_state(state)?);
2849 }
2850 if self.identity_state_inventory(IdentityStateStorageView::Effective)?
2851 != *checkpoint_identity_states
2852 {
2853 return Err(InternalError::identity_state_corruption());
2854 }
2855 let selection = self
2856 .current_accepted_schema_root()?
2857 .ok_or_else(InternalError::store_corruption)?;
2858 self.retain_durable_candidate_entries(candidate, selection.slot())?;
2859 return Ok(());
2860 }
2861 if self.current_accepted_schema_root()?.is_some()
2862 || !self
2863 .identity_state_inventory(IdentityStateStorageView::Effective)?
2864 .is_empty()
2865 {
2866 return Err(InternalError::store_corruption());
2867 }
2868
2869 self.insert_durable_candidate_snapshots(candidate)?;
2870 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2871 self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2872 for state in checkpoint_identity_states.values() {
2873 let key = RawSchemaKey::from_identity_state(
2874 state.owner().entity_tag(),
2875 state.owner().field_id(),
2876 );
2877 self.insert_durable_raw_value(key, encode_identity_state(state)?);
2878 }
2879 let root_key = RawSchemaKey::from_accepted_root_slot(0)?;
2880 self.insert_durable_raw_value(root_key, candidate.encoded_root().to_vec());
2881
2882 let selected = self
2883 .current_accepted_schema_root()?
2884 .ok_or_else(InternalError::store_corruption)?;
2885 if selected.root() != candidate.root() {
2886 return Err(InternalError::store_corruption());
2887 }
2888 self.retain_durable_candidate_entries(candidate, selected.slot())?;
2889 Ok(())
2890 }
2891
2892 pub(in crate::db) fn preflight_accepted_schema_candidate(
2899 &self,
2900 incarnation: DatabaseIncarnationId,
2901 expected_revision: AcceptedSchemaRevision,
2902 candidate: &CandidateSchemaRevision,
2903 ) -> Result<bool, InternalError> {
2904 let identity_transition = self.prepare_identity_state_transition(
2905 incarnation,
2906 candidate,
2907 IdentityStateStorageView::Effective,
2908 )?;
2909 if self.current_root_matches_candidate(candidate)? {
2910 if !identity_transition.is_empty() {
2911 return Err(InternalError::identity_state_corruption());
2912 }
2913 return Ok(true);
2914 }
2915 let first = self.accepted_root_slot_bytes(0)?;
2916 let second = self.accepted_root_slot_bytes(1)?;
2917 prepare_accepted_schema_root_publication(
2918 [first.as_deref(), second.as_deref()],
2919 expected_revision,
2920 candidate,
2921 )
2922 .map_err(map_schema_publication_error)?;
2923
2924 Ok(false)
2925 }
2926
2927 pub(in crate::db) fn preflight_fold_journaled_accepted_schema_candidate(
2929 &self,
2930 incarnation: DatabaseIncarnationId,
2931 expected_revision: AcceptedSchemaRevision,
2932 candidate: &CandidateSchemaRevision,
2933 ) -> Result<(), InternalError> {
2934 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2935 return Err(InternalError::store_invariant());
2936 }
2937 let identity_transition = self.prepare_identity_state_transition(
2938 incarnation,
2939 candidate,
2940 IdentityStateStorageView::Canonical,
2941 )?;
2942 let candidate_is_current = self.canonical_root_matches_candidate(candidate)?;
2943 if candidate_is_current && !identity_transition.is_empty() {
2944 return Err(InternalError::identity_state_corruption());
2945 }
2946 for state in identity_transition.into_updates() {
2947 let _encoded = encode_identity_state(&state)?;
2948 }
2949 for snapshot in candidate.bundle().entity_snapshots().values() {
2950 let _encoded = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
2951 }
2952
2953 let first = self.canonical_root_slot_bytes(0)?;
2954 let second = self.canonical_root_slot_bytes(1)?;
2955 let root_slot = if candidate_is_current {
2956 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
2957 .ok_or_else(InternalError::store_corruption)?
2958 .slot()
2959 } else {
2960 prepare_accepted_schema_root_publication(
2961 [first.as_deref(), second.as_deref()],
2962 expected_revision,
2963 candidate,
2964 )
2965 .map_err(map_schema_publication_error)?
2966 .target_slot()
2967 };
2968 let _retained = Self::candidate_entry_keys(candidate, root_slot)?;
2969 Ok(())
2970 }
2971
2972 pub(in crate::db) fn projected_identity_state_count(
2974 &self,
2975 incarnation: DatabaseIncarnationId,
2976 candidate: &CandidateSchemaRevision,
2977 ) -> Result<usize, InternalError> {
2978 Ok(self
2979 .prepare_identity_state_transition(
2980 incarnation,
2981 candidate,
2982 IdentityStateStorageView::Effective,
2983 )?
2984 .projected_inventory_len())
2985 }
2986
2987 pub(in crate::db) fn apply_journaled_accepted_schema_candidate(
2989 &mut self,
2990 incarnation: DatabaseIncarnationId,
2991 expected_revision: AcceptedSchemaRevision,
2992 candidate: &CandidateSchemaRevision,
2993 ) -> Result<(), InternalError> {
2994 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2995 return Err(InternalError::store_invariant());
2996 }
2997 let identity_transition = self.prepare_identity_state_transition(
2998 incarnation,
2999 candidate,
3000 IdentityStateStorageView::Effective,
3001 )?;
3002 if self.current_root_matches_candidate(candidate)? {
3003 if !identity_transition.is_empty() {
3004 return Err(InternalError::identity_state_corruption());
3005 }
3006 let selection = self
3007 .current_accepted_schema_root()?
3008 .ok_or_else(InternalError::store_corruption)?;
3009 self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3010 return Ok(());
3011 }
3012
3013 let first = self.accepted_root_slot_bytes(0)?;
3014 let second = self.accepted_root_slot_bytes(1)?;
3015 prepare_accepted_schema_root_publication(
3016 [first.as_deref(), second.as_deref()],
3017 expected_revision,
3018 candidate,
3019 )
3020 .map_err(map_schema_publication_error)?;
3021
3022 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3023 self.insert_persisted_snapshot(*entity_tag, snapshot)?;
3024 }
3025 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3026 self.insert_raw_snapshot(
3027 bundle_key,
3028 RawSchemaSnapshot::from_encoded_control_record(candidate.encoded_bundle().to_vec()),
3029 );
3030 let persisted_bundle = self
3031 .get_raw_snapshot(&bundle_key)
3032 .ok_or_else(InternalError::store_corruption)?;
3033 let _verified = decode_verified_accepted_schema_revision_bundle(
3034 candidate.root(),
3035 persisted_bundle.as_bytes(),
3036 )?;
3037 self.apply_identity_state_transition(
3038 identity_transition,
3039 IdentityStateWriteTarget::Materialized,
3040 )?;
3041
3042 let first = self.accepted_root_slot_bytes(0)?;
3043 let second = self.accepted_root_slot_bytes(1)?;
3044 let publication = prepare_accepted_schema_root_publication(
3045 [first.as_deref(), second.as_deref()],
3046 expected_revision,
3047 candidate,
3048 )
3049 .map_err(map_schema_publication_error)?;
3050 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3051 self.insert_raw_snapshot(
3052 root_key,
3053 RawSchemaSnapshot::from_encoded_control_record(publication.encoded_root().to_vec()),
3054 );
3055
3056 if !self.current_root_matches_candidate(candidate)? {
3057 return Err(InternalError::store_corruption());
3058 }
3059 let selection = self
3060 .current_accepted_schema_root()?
3061 .ok_or_else(InternalError::store_corruption)?;
3062 self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3063 Ok(())
3064 }
3065
3066 pub(in crate::db) fn fold_journaled_accepted_schema_candidate(
3068 &mut self,
3069 incarnation: DatabaseIncarnationId,
3070 expected_revision: AcceptedSchemaRevision,
3071 candidate: &CandidateSchemaRevision,
3072 ) -> Result<(), InternalError> {
3073 let identity_transition = self.prepare_identity_state_transition(
3074 incarnation,
3075 candidate,
3076 IdentityStateStorageView::Canonical,
3077 )?;
3078 if self.canonical_root_matches_candidate(candidate)? {
3079 if !identity_transition.is_empty() {
3080 return Err(InternalError::identity_state_corruption());
3081 }
3082 let first = self.canonical_root_slot_bytes(0)?;
3083 let second = self.canonical_root_slot_bytes(1)?;
3084 let selection =
3085 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3086 .ok_or_else(InternalError::store_corruption)?;
3087 self.retain_canonical_candidate_entries(candidate, selection.slot())?;
3088 return Ok(());
3089 }
3090
3091 let first = self.canonical_root_slot_bytes(0)?;
3092 let second = self.canonical_root_slot_bytes(1)?;
3093 prepare_accepted_schema_root_publication(
3094 [first.as_deref(), second.as_deref()],
3095 expected_revision,
3096 candidate,
3097 )
3098 .map_err(map_schema_publication_error)?;
3099
3100 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3101 self.fold_persisted_snapshot(*entity_tag, snapshot)?;
3102 }
3103 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3104 self.insert_canonical_raw_value(bundle_key, candidate.encoded_bundle().to_vec())?;
3105 let persisted_bundle = self
3106 .get_canonical_raw_value(&bundle_key)?
3107 .ok_or_else(InternalError::store_corruption)?;
3108 let _verified = decode_verified_accepted_schema_revision_bundle(
3109 candidate.root(),
3110 persisted_bundle.as_bytes(),
3111 )?;
3112 self.apply_identity_state_transition(
3113 identity_transition,
3114 IdentityStateWriteTarget::Canonical,
3115 )?;
3116
3117 let first = self.canonical_root_slot_bytes(0)?;
3118 let second = self.canonical_root_slot_bytes(1)?;
3119 let publication = prepare_accepted_schema_root_publication(
3120 [first.as_deref(), second.as_deref()],
3121 expected_revision,
3122 candidate,
3123 )
3124 .map_err(map_schema_publication_error)?;
3125 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3126 self.insert_canonical_raw_value(root_key, publication.encoded_root().to_vec())?;
3127
3128 if !self.canonical_root_matches_candidate(candidate)? {
3129 return Err(InternalError::store_corruption());
3130 }
3131 let first = self.canonical_root_slot_bytes(0)?;
3132 let second = self.canonical_root_slot_bytes(1)?;
3133 let selection = select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3134 .ok_or_else(InternalError::store_corruption)?;
3135 self.retain_canonical_candidate_entries(candidate, selection.slot())?;
3136 Ok(())
3137 }
3138
3139 pub(in crate::db) fn get_persisted_snapshot(
3141 &self,
3142 entity: EntityTag,
3143 version: SchemaVersion,
3144 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3145 let key = RawSchemaKey::from_entity_version(entity, version);
3146 self.get_raw_snapshot(&key)
3147 .map(|snapshot| snapshot.decode_persisted_snapshot())
3148 .transpose()
3149 }
3150
3151 #[cfg(test)]
3152 fn latest_staged_persisted_snapshot(
3153 &self,
3154 entity: EntityTag,
3155 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3156 self.latest_raw_snapshots_by_entity()
3157 .remove(&entity)
3158 .map(|(_, snapshot)| snapshot.decode_persisted_snapshot())
3159 .transpose()
3160 }
3161
3162 pub(in crate::db) fn current_accepted_persisted_snapshot(
3165 &self,
3166 entity: EntityTag,
3167 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3168 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3169 return Ok(None);
3170 };
3171
3172 Ok(bundle.entity_snapshots().get(&entity).cloned())
3173 }
3174
3175 pub(in crate::db) fn current_accepted_catalog_selection(
3177 &self,
3178 entity: EntityTag,
3179 entity_path: &str,
3180 store_path: &'static str,
3181 ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3182 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3183 return Ok(None);
3184 };
3185 if bundle.store_path() != store_path {
3186 return Err(InternalError::store_corruption());
3187 }
3188 let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3189 return Ok(None);
3190 };
3191 if snapshot.entity_path() != entity_path {
3192 return Err(InternalError::store_corruption());
3193 }
3194
3195 let cache = self
3196 .accepted_bundle_cache
3197 .try_borrow()
3198 .map_err(|_| InternalError::store_invariant())?;
3199 let cached = cache.as_ref().ok_or_else(InternalError::store_invariant)?;
3200 if let Some(selection) = cached
3201 .entity_selections
3202 .try_borrow()
3203 .map_err(|_| InternalError::store_invariant())?
3204 .get(&entity)
3205 .cloned()
3206 {
3207 return Ok(Some(selection));
3208 }
3209
3210 let raw_snapshot = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3211 let fingerprint = raw_snapshot.accepted_schema_fingerprint()?;
3212 let identity = AcceptedCatalogIdentity::new(
3213 entity,
3214 entity_path,
3215 store_path,
3216 bundle.revision(),
3217 snapshot.version(),
3218 fingerprint,
3219 );
3220
3221 let selected = AcceptedCatalogSnapshotSelection::new(
3222 identity,
3223 cached.value_catalog.clone(),
3224 Rc::from(raw_snapshot.into_bytes()),
3225 );
3226 cached
3227 .entity_selections
3228 .try_borrow_mut()
3229 .map_err(|_| InternalError::store_invariant())?
3230 .insert(entity, selected.clone());
3231
3232 Ok(Some(selected))
3233 }
3234
3235 pub(in crate::db) fn current_canonical_accepted_catalog_selection(
3239 &self,
3240 entity: EntityTag,
3241 entity_path: &str,
3242 store_path: &'static str,
3243 ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3244 let first = self.canonical_root_slot_bytes(0)?;
3245 let second = self.canonical_root_slot_bytes(1)?;
3246 let Some(selection) =
3247 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3248 else {
3249 return Ok(None);
3250 };
3251 let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3252 let raw_bundle = self
3253 .get_canonical_raw_value(&bundle_key)?
3254 .ok_or_else(InternalError::store_corruption)?;
3255 let bundle = decode_verified_accepted_schema_revision_bundle(
3256 selection.root(),
3257 raw_bundle.as_bytes(),
3258 )?;
3259 if bundle.store_path() != store_path {
3260 return Err(InternalError::store_corruption());
3261 }
3262 let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3263 return Ok(None);
3264 };
3265 if snapshot.entity_path() != entity_path {
3266 return Err(InternalError::store_corruption());
3267 }
3268
3269 let raw_snapshot = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3270 let fingerprint = raw_snapshot.accepted_schema_fingerprint()?;
3271 let identity = AcceptedCatalogIdentity::new(
3272 entity,
3273 entity_path,
3274 store_path,
3275 bundle.revision(),
3276 snapshot.version(),
3277 fingerprint,
3278 );
3279
3280 Ok(Some(AcceptedCatalogSnapshotSelection::new(
3281 identity,
3282 AcceptedValueCatalogHandle::new(
3283 bundle.enum_catalog().clone(),
3284 bundle.composite_catalog().clone(),
3285 self.accepted_catalog_scope
3286 .get_or_init(AcceptedStoreCatalogScope::new)
3287 .clone(),
3288 bundle.revision(),
3289 selection.root().fingerprint(),
3290 ),
3291 Rc::from(raw_snapshot.into_bytes()),
3292 )))
3293 }
3294
3295 #[cfg(test)]
3301 pub(in crate::db) fn catalog_metadata(
3302 &self,
3303 ) -> Result<Option<SchemaStoreCatalogMetadata>, InternalError> {
3304 Ok(self
3305 .allocation_metadata()?
3306 .map(SchemaStoreAllocationMetadata::schema))
3307 }
3308
3309 pub(in crate::db) fn allocation_metadata(
3316 &self,
3317 ) -> Result<Option<SchemaStoreAllocationMetadata>, InternalError> {
3318 let latest_by_entity = self.latest_raw_snapshots_by_entity();
3319 if latest_by_entity.is_empty() {
3320 return Ok(None);
3321 }
3322
3323 Ok(Some(SchemaStoreAllocationMetadata::new(
3324 derive_data_allocation_metadata(&latest_by_entity)?,
3325 derive_index_allocation_metadata(&latest_by_entity)?,
3326 derive_schema_catalog_metadata(&latest_by_entity)?,
3327 )))
3328 }
3329
3330 fn insert_raw_snapshot(
3332 &mut self,
3333 key: RawSchemaKey,
3334 snapshot: RawSchemaSnapshot,
3335 ) -> Option<RawSchemaSnapshot> {
3336 self.invalidate_accepted_bundle_cache_for_key(key);
3337 let previous_journaled = if matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3338 self.get_raw_snapshot_for_backend(&key)
3339 } else {
3340 None
3341 };
3342 match &mut self.backend {
3343 SchemaStoreBackend::Heap(map) => map.insert(key, snapshot),
3344 SchemaStoreBackend::Journaled {
3345 live, tombstones, ..
3346 } => {
3347 tombstones.remove(&key);
3348 live.insert(key, snapshot);
3349 previous_journaled
3350 }
3351 }
3352 }
3353
3354 #[must_use]
3356 fn get_raw_snapshot(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
3357 match &self.backend {
3358 SchemaStoreBackend::Heap(map) => map.get(key).cloned(),
3359 SchemaStoreBackend::Journaled { .. } => self.get_raw_snapshot_for_backend(key),
3360 }
3361 }
3362
3363 fn accepted_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3364 let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3365 Ok(self
3366 .get_raw_snapshot(&key)
3367 .map(RawSchemaSnapshot::into_bytes))
3368 }
3369
3370 fn canonical_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3371 let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3372 Ok(self
3373 .get_canonical_raw_value(&key)?
3374 .map(RawSchemaSnapshot::into_bytes))
3375 }
3376
3377 fn current_root_matches_candidate(
3378 &self,
3379 candidate: &CandidateSchemaRevision,
3380 ) -> Result<bool, InternalError> {
3381 let Some(selection) = self.current_accepted_schema_root()? else {
3382 return Ok(false);
3383 };
3384 if selection.root() != candidate.root() {
3385 return Ok(false);
3386 }
3387 let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3388 let bundle = self
3389 .get_raw_snapshot(&key)
3390 .ok_or_else(InternalError::store_corruption)?;
3391 let _verified =
3392 decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3393 Ok(true)
3394 }
3395
3396 fn canonical_root_matches_candidate(
3397 &self,
3398 candidate: &CandidateSchemaRevision,
3399 ) -> Result<bool, InternalError> {
3400 let first = self.canonical_root_slot_bytes(0)?;
3401 let second = self.canonical_root_slot_bytes(1)?;
3402 let Some(selection) =
3403 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3404 else {
3405 return Ok(false);
3406 };
3407 if selection.root() != candidate.root() {
3408 return Ok(false);
3409 }
3410 let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3411 let bundle = self
3412 .get_canonical_raw_value(&key)?
3413 .ok_or_else(InternalError::store_corruption)?;
3414 let _verified =
3415 decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3416 Ok(true)
3417 }
3418
3419 fn get_canonical_raw_value(
3420 &self,
3421 key: &RawSchemaKey,
3422 ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
3423 match &self.backend {
3424 SchemaStoreBackend::Journaled { canonical, .. } => Ok(canonical.get(key)),
3425 SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
3426 }
3427 }
3428
3429 fn insert_canonical_raw_value(
3430 &mut self,
3431 key: RawSchemaKey,
3432 bytes: Vec<u8>,
3433 ) -> Result<(), InternalError> {
3434 self.invalidate_accepted_bundle_cache_for_key(key);
3435 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3436 return Err(InternalError::store_invariant());
3437 };
3438 canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
3439 Ok(())
3440 }
3441
3442 fn insert_durable_raw_value(&mut self, key: RawSchemaKey, bytes: Vec<u8>) {
3446 self.invalidate_accepted_bundle_cache_for_key(key);
3447 let value = RawSchemaSnapshot::from_encoded_control_record(bytes);
3448 match &mut self.backend {
3449 SchemaStoreBackend::Heap(map) => {
3450 map.insert(key, value);
3451 }
3452 SchemaStoreBackend::Journaled {
3453 canonical,
3454 live,
3455 tombstones,
3456 ..
3457 } => {
3458 live.remove(&key);
3459 tombstones.remove(&key);
3460 canonical.insert(key, value);
3461 }
3462 }
3463 }
3464
3465 fn invalidate_accepted_bundle_cache_for_key(&mut self, key: RawSchemaKey) {
3466 if key.is_accepted_root() {
3467 self.accepted_bundle_cache.get_mut().take();
3468 }
3469 }
3470
3471 fn insert_durable_candidate_snapshots(
3472 &mut self,
3473 candidate: &CandidateSchemaRevision,
3474 ) -> Result<(), InternalError> {
3475 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3476 let key = RawSchemaKey::from_entity_version(*entity_tag, snapshot.version());
3477 let value = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3478 match &mut self.backend {
3479 SchemaStoreBackend::Heap(map) => {
3480 map.insert(key, value);
3481 }
3482 SchemaStoreBackend::Journaled {
3483 canonical,
3484 live,
3485 tombstones,
3486 ..
3487 } => {
3488 live.remove(&key);
3489 tombstones.remove(&key);
3490 canonical.insert(key, value);
3491 }
3492 }
3493 }
3494 Ok(())
3495 }
3496
3497 fn candidate_entry_keys(
3498 candidate: &CandidateSchemaRevision,
3499 root_slot: usize,
3500 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3501 let mut keys = candidate
3502 .bundle()
3503 .entity_snapshots()
3504 .iter()
3505 .map(|(entity_tag, snapshot)| {
3506 RawSchemaKey::from_entity_version(*entity_tag, snapshot.version())
3507 })
3508 .collect::<BTreeSet<_>>();
3509 keys.insert(RawSchemaKey::from_accepted_bundle(
3510 candidate.root().bundle_key(),
3511 ));
3512 keys.insert(RawSchemaKey::from_accepted_root_slot(root_slot)?);
3513 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3514 for activation in snapshot
3515 .constraint_activations()
3516 .iter()
3517 .filter(|activation| activation.state() == ConstraintActivationState::Validating)
3518 {
3519 keys.insert(RawSchemaKey::from_constraint_validation_job(
3520 *entity_tag,
3521 activation.id(),
3522 ));
3523 }
3524 }
3525 Ok(keys)
3526 }
3527
3528 fn positioned_candidate_effect_keys(
3529 &self,
3530 incarnation: DatabaseIncarnationId,
3531 expected_revision: AcceptedSchemaRevision,
3532 candidate: &CandidateSchemaRevision,
3533 view: IdentityStateStorageView,
3534 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3535 let identity_transition =
3536 self.prepare_identity_state_transition(incarnation, candidate, view)?;
3537 let (first, second, candidate_is_current) = match view {
3538 IdentityStateStorageView::Effective => (
3539 self.accepted_root_slot_bytes(0)?,
3540 self.accepted_root_slot_bytes(1)?,
3541 self.current_root_matches_candidate(candidate)?,
3542 ),
3543 IdentityStateStorageView::Canonical => (
3544 self.canonical_root_slot_bytes(0)?,
3545 self.canonical_root_slot_bytes(1)?,
3546 self.canonical_root_matches_candidate(candidate)?,
3547 ),
3548 };
3549 let root_slot = if candidate_is_current {
3550 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3551 .ok_or_else(InternalError::store_corruption)?
3552 .slot()
3553 } else {
3554 prepare_accepted_schema_root_publication(
3555 [first.as_deref(), second.as_deref()],
3556 expected_revision,
3557 candidate,
3558 )
3559 .map_err(map_schema_publication_error)?
3560 .target_slot()
3561 };
3562 let mut keys = Self::candidate_entry_keys(candidate, root_slot)?;
3563 for state in identity_transition.into_updates() {
3564 keys.insert(RawSchemaKey::from_identity_state(
3565 state.owner().entity_tag(),
3566 state.owner().field_id(),
3567 ));
3568 }
3569
3570 let SchemaStoreBackend::Journaled {
3571 canonical,
3572 live,
3573 tombstones,
3574 positions,
3575 } = &self.backend
3576 else {
3577 return Err(InternalError::store_invariant());
3578 };
3579 for entry in canonical.iter() {
3580 let has_relevant_overlay = matches!(view, IdentityStateStorageView::Effective)
3581 || positions.is_positioned(entry.key())
3582 || live.contains_key(entry.key())
3583 || tombstones.contains(entry.key());
3584 if has_relevant_overlay
3585 && !keys.contains(entry.key())
3586 && !entry.key().is_identity_state()
3587 {
3588 keys.insert(*entry.key());
3589 }
3590 }
3591 if matches!(view, IdentityStateStorageView::Effective) {
3592 for key in live.keys() {
3593 if !keys.contains(key) && !key.is_identity_state() {
3594 keys.insert(*key);
3595 }
3596 }
3597 }
3598 Ok(keys)
3599 }
3600
3601 fn positioned_journal_batch_keys(
3602 &self,
3603 incarnation: DatabaseIncarnationId,
3604 batch: &JournalBatch,
3605 view: IdentityStateStorageView,
3606 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3607 let mut keys = BTreeSet::new();
3608 for record in batch.records() {
3609 match record {
3610 JournalRecord::SchemaPut {
3611 schema_snapshot_bytes,
3612 ..
3613 } => {
3614 let snapshot = decode_persisted_schema_snapshot(schema_snapshot_bytes)?;
3615 let entity_tag = match view {
3616 IdentityStateStorageView::Effective => self
3617 .current_accepted_schema_bundle_ref()?
3618 .ok_or_else(InternalError::store_corruption)?
3619 .entity_snapshots()
3620 .iter()
3621 .find_map(|(entity_tag, accepted)| {
3622 (accepted.entity_path() == snapshot.entity_path())
3623 .then_some(*entity_tag)
3624 }),
3625 IdentityStateStorageView::Canonical => self
3626 .current_canonical_accepted_schema_bundle()?
3627 .ok_or_else(InternalError::store_corruption)?
3628 .entity_snapshots()
3629 .iter()
3630 .find_map(|(entity_tag, accepted)| {
3631 (accepted.entity_path() == snapshot.entity_path())
3632 .then_some(*entity_tag)
3633 }),
3634 }
3635 .ok_or_else(InternalError::store_corruption)?;
3636 keys.insert(RawSchemaKey::from_entity_version(
3637 entity_tag,
3638 snapshot.version(),
3639 ));
3640 }
3641 JournalRecord::AcceptedSchemaPublish {
3642 expected_revision,
3643 schema_bundle_bytes,
3644 schema_root_bytes,
3645 ..
3646 } => {
3647 let candidate = CandidateSchemaRevision::from_encoded(
3648 schema_bundle_bytes.clone(),
3649 schema_root_bytes.clone(),
3650 )?;
3651 keys.extend(self.positioned_candidate_effect_keys(
3652 incarnation,
3653 *expected_revision,
3654 &candidate,
3655 view,
3656 )?);
3657 }
3658 JournalRecord::ConstraintValidationJobPut {
3659 entity_tag,
3660 constraint_id,
3661 ..
3662 }
3663 | JournalRecord::ConstraintValidationJobDelete {
3664 entity_tag,
3665 constraint_id,
3666 ..
3667 } => {
3668 keys.insert(RawSchemaKey::from_constraint_validation_job(
3669 *entity_tag,
3670 *constraint_id,
3671 ));
3672 }
3673 JournalRecord::IdentityRangeAdvance { range } => {
3674 keys.insert(RawSchemaKey::from_identity_state(
3675 range.owner().entity_tag(),
3676 range.owner().field_id(),
3677 ));
3678 }
3679 JournalRecord::RowPut { .. }
3680 | JournalRecord::RowDelete { .. }
3681 | JournalRecord::AcceptedSchemaIndexDelete { .. }
3682 | JournalRecord::AcceptedSchemaIndexPut { .. }
3683 | JournalRecord::ConstraintValidationIndexPut { .. } => {}
3684 #[cfg(any(test, feature = "migration"))]
3685 JournalRecord::SchemaMigrationRowPut { .. }
3686 | JournalRecord::SchemaMigrationIndexPut { .. } => {}
3687 }
3688 }
3689 Ok(keys)
3690 }
3691
3692 fn retain_durable_candidate_entries(
3696 &mut self,
3697 candidate: &CandidateSchemaRevision,
3698 root_slot: usize,
3699 ) -> Result<(), InternalError> {
3700 let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3701 self.accepted_bundle_cache.get_mut().take();
3702 match &mut self.backend {
3703 SchemaStoreBackend::Heap(map) => {
3704 map.retain(|key, _| keep.contains(key) || key.is_identity_state());
3705 }
3706 SchemaStoreBackend::Journaled {
3707 canonical,
3708 live,
3709 tombstones,
3710 ..
3711 } => {
3712 let stale = canonical
3713 .iter()
3714 .filter_map(|entry| {
3715 (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3716 .then_some(*entry.key())
3717 })
3718 .collect::<Vec<_>>();
3719 for key in stale {
3720 canonical.remove(&key);
3721 }
3722 live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3723 tombstones.clear();
3724 }
3725 }
3726 Ok(())
3727 }
3728
3729 fn retain_materialized_candidate_entries(
3730 &mut self,
3731 candidate: &CandidateSchemaRevision,
3732 root_slot: usize,
3733 ) -> Result<(), InternalError> {
3734 let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3735 self.accepted_bundle_cache.get_mut().take();
3736 let SchemaStoreBackend::Journaled {
3737 canonical,
3738 live,
3739 tombstones,
3740 ..
3741 } = &mut self.backend
3742 else {
3743 return Err(InternalError::store_invariant());
3744 };
3745 live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3746 let canonical_keys = canonical
3747 .iter()
3748 .map(|entry| *entry.key())
3749 .collect::<Vec<_>>();
3750 for key in canonical_keys {
3751 if keep.contains(&key) || key.is_identity_state() {
3752 tombstones.remove(&key);
3753 } else {
3754 tombstones.insert(key);
3755 }
3756 }
3757 Ok(())
3758 }
3759
3760 fn retain_canonical_candidate_entries(
3761 &mut self,
3762 candidate: &CandidateSchemaRevision,
3763 root_slot: usize,
3764 ) -> Result<(), InternalError> {
3765 let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3766 self.accepted_bundle_cache.get_mut().take();
3767 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3768 return Err(InternalError::store_invariant());
3769 };
3770 let stale = canonical
3771 .iter()
3772 .filter_map(|entry| {
3773 (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3774 .then_some(*entry.key())
3775 })
3776 .collect::<Vec<_>>();
3777 for key in stale {
3778 canonical.remove(&key);
3779 }
3780 Ok(())
3781 }
3782
3783 #[must_use]
3785 #[cfg(test)]
3786 fn contains_raw_snapshot(&self, key: &RawSchemaKey) -> bool {
3787 match &self.backend {
3788 SchemaStoreBackend::Heap(map) => map.contains_key(key),
3789 SchemaStoreBackend::Journaled { .. } => {
3790 self.get_raw_snapshot_for_backend(key).is_some()
3791 }
3792 }
3793 }
3794
3795 #[must_use]
3797 #[cfg(test)]
3798 pub(in crate::db) fn len(&self) -> u64 {
3799 match &self.backend {
3800 SchemaStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
3801 SchemaStoreBackend::Journaled { .. } => {
3802 let mut count = 0_u64;
3803 let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3804 count = count.saturating_add(1);
3805 Ok(SchemaStoreVisit::Continue)
3806 });
3807 count
3808 }
3809 }
3810 }
3811
3812 #[must_use]
3814 #[cfg(test)]
3815 pub(in crate::db) fn is_empty(&self) -> bool {
3816 match &self.backend {
3817 SchemaStoreBackend::Heap(map) => map.is_empty(),
3818 SchemaStoreBackend::Journaled { .. } => {
3819 let mut empty = true;
3820 let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3821 empty = false;
3822 Ok(SchemaStoreVisit::Stop)
3823 });
3824 empty
3825 }
3826 }
3827 }
3828
3829 #[cfg(test)]
3831 pub(in crate::db) fn clear(&mut self) {
3832 self.accepted_bundle_cache.get_mut().take();
3833 match &mut self.backend {
3834 SchemaStoreBackend::Heap(map) => map.clear(),
3835 SchemaStoreBackend::Journaled {
3836 canonical,
3837 live,
3838 tombstones,
3839 ..
3840 } => {
3841 live.clear();
3842 tombstones.clear();
3843 let keys = canonical
3844 .iter()
3845 .map(|entry| *entry.key())
3846 .collect::<Vec<_>>();
3847 for key in keys {
3848 if key.is_entity_snapshot() {
3849 tombstones.insert(key);
3850 } else {
3851 canonical.remove(&key);
3852 }
3853 }
3854 }
3855 }
3856 }
3857
3858 fn current_accepted_schema_bundle_ref(
3859 &self,
3860 ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
3861 self.current_accepted_schema_authority_ref()
3862 .map(|authority| authority.map(|(_selection, bundle)| bundle))
3863 }
3864
3865 pub(in crate::db) fn current_accepted_schema_authority_ref(
3867 &self,
3868 ) -> Result<
3869 Option<(
3870 AcceptedSchemaRootSelection,
3871 Ref<'_, AcceptedSchemaRevisionBundle>,
3872 )>,
3873 InternalError,
3874 > {
3875 let selection = self.current_accepted_schema_root()?;
3876 self.accepted_schema_authority_ref_for_selection(selection)
3877 }
3878
3879 fn accepted_schema_authority_ref_for_selection(
3880 &self,
3881 selection: Option<AcceptedSchemaRootSelection>,
3882 ) -> Result<
3883 Option<(
3884 AcceptedSchemaRootSelection,
3885 Ref<'_, AcceptedSchemaRevisionBundle>,
3886 )>,
3887 InternalError,
3888 > {
3889 let Some(selection) = selection else {
3890 self.accepted_bundle_cache
3891 .try_borrow_mut()
3892 .map_err(|_| InternalError::store_invariant())?
3893 .take();
3894 return Ok(None);
3895 };
3896
3897 let cache_matches = self
3898 .accepted_bundle_cache
3899 .try_borrow()
3900 .map_err(|_| InternalError::store_invariant())?
3901 .as_ref()
3902 .is_some_and(|cached| cached.selection == selection);
3903 if !cache_matches {
3904 let key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3905 let raw = self
3906 .get_raw_snapshot(&key)
3907 .ok_or_else(InternalError::store_corruption)?;
3908 let bundle =
3909 decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
3910 self.validate_constraint_validation_job_closure(&bundle)?;
3911 #[cfg(test)]
3912 ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES
3913 .with(|misses| misses.set(misses.get().saturating_add(1)));
3914 let value_catalog = AcceptedValueCatalogHandle::new(
3915 bundle.enum_catalog().clone(),
3916 bundle.composite_catalog().clone(),
3917 self.accepted_catalog_scope
3918 .get_or_init(AcceptedStoreCatalogScope::new)
3919 .clone(),
3920 bundle.revision(),
3921 selection.root().fingerprint(),
3922 );
3923 let cardinality_domain = Rc::new(CardinalityAcceptedDomain::derive(&bundle)?);
3924 *self
3925 .accepted_bundle_cache
3926 .try_borrow_mut()
3927 .map_err(|_| InternalError::store_invariant())? = Some(AcceptedSchemaBundleCache {
3928 selection,
3929 bundle,
3930 cardinality_domain,
3931 value_catalog,
3932 entity_selections: RefCell::new(StdBTreeMap::new()),
3933 });
3934 }
3935
3936 let cache = self
3937 .accepted_bundle_cache
3938 .try_borrow()
3939 .map_err(|_| InternalError::store_invariant())?;
3940 let bundle = Ref::filter_map(cache, |cache| {
3941 cache
3942 .as_ref()
3943 .filter(|cached| cached.selection == selection)
3944 .map(|cached| &cached.bundle)
3945 })
3946 .map_err(|_| InternalError::store_invariant())?;
3947 self.validate_identity_state_closure(&bundle)?;
3948 Ok(Some((selection, bundle)))
3949 }
3950
3951 pub(in crate::db) fn accepted_cardinality_domain_for_selection(
3953 &self,
3954 selection: Option<AcceptedSchemaRootSelection>,
3955 ) -> Result<Option<(AcceptedSchemaRootSelection, Rc<CardinalityAcceptedDomain>)>, InternalError>
3956 {
3957 let Some(selection) = selection else {
3958 self.accepted_bundle_cache
3959 .try_borrow_mut()
3960 .map_err(|_| InternalError::store_invariant())?
3961 .take();
3962 return Ok(None);
3963 };
3964 let cache_matches = self
3965 .accepted_bundle_cache
3966 .try_borrow()
3967 .map_err(|_| InternalError::store_invariant())?
3968 .as_ref()
3969 .is_some_and(|cached| cached.selection == selection);
3970 if !cache_matches {
3971 let authority = self
3972 .accepted_schema_authority_ref_for_selection(Some(selection))?
3973 .ok_or_else(InternalError::store_invariant)?;
3974 drop(authority);
3975 }
3976 let cache = self
3977 .accepted_bundle_cache
3978 .try_borrow()
3979 .map_err(|_| InternalError::store_invariant())?;
3980 let domain = cache
3981 .as_ref()
3982 .filter(|cached| cached.selection == selection)
3983 .map(|cached| Rc::clone(&cached.cardinality_domain))
3984 .ok_or_else(InternalError::store_invariant)?;
3985 Ok(Some((selection, domain)))
3986 }
3987
3988 pub(in crate::db) fn cached_cardinality_domain_for_root(
3990 &self,
3991 root: CardinalityAcceptedRootIdentity,
3992 ) -> Result<Option<Rc<CardinalityAcceptedDomain>>, InternalError> {
3993 let cache = self
3994 .accepted_bundle_cache
3995 .try_borrow()
3996 .map_err(|_| InternalError::store_invariant())?;
3997 Ok(cache
3998 .as_ref()
3999 .filter(|cached| root.matches(cached.selection.root()))
4000 .map(|cached| Rc::clone(&cached.cardinality_domain)))
4001 }
4002
4003 fn latest_raw_snapshots_by_entity(
4004 &self,
4005 ) -> StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)> {
4006 let mut latest_by_entity =
4007 StdBTreeMap::<EntityTag, (SchemaVersion, RawSchemaSnapshot)>::new();
4008
4009 let _: Result<(), std::convert::Infallible> = self.visit_raw_snapshots(|key, snapshot| {
4010 let version = SchemaVersion::new(key.version());
4011 match latest_by_entity.get_mut(&key.entity_tag()) {
4012 Some((latest_version, latest_snapshot)) if version > *latest_version => {
4013 *latest_version = version;
4014 *latest_snapshot = snapshot.clone();
4015 }
4016 None => {
4017 latest_by_entity.insert(key.entity_tag(), (version, snapshot.clone()));
4018 }
4019 Some(_) => {}
4020 }
4021 Ok(SchemaStoreVisit::Continue)
4022 });
4023
4024 latest_by_entity
4025 }
4026
4027 fn visit_raw_snapshots<E>(
4030 &self,
4031 visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4032 ) -> Result<(), E> {
4033 let bounds = RawSchemaKey::all_entity_range_bounds();
4034 match &self.backend {
4035 SchemaStoreBackend::Heap(map) => {
4036 let mut visitor = visitor;
4037 for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4038 if visitor(key, snapshot)?.should_stop() {
4039 break;
4040 }
4041 }
4042 }
4043 SchemaStoreBackend::Journaled {
4044 canonical,
4045 live,
4046 tombstones,
4047 ..
4048 } => Self::visit_journaled_raw_snapshot_range(
4049 canonical,
4050 live,
4051 tombstones,
4052 bounds,
4053 Direction::Asc,
4054 visitor,
4055 )?,
4056 }
4057
4058 Ok(())
4059 }
4060
4061 fn visit_constraint_validation_jobs_in_view<E>(
4062 &self,
4063 view: IdentityStateStorageView,
4064 visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4065 ) -> Result<(), E> {
4066 let bounds = RawSchemaKey::all_constraint_validation_job_range_bounds();
4067 match (&self.backend, view) {
4068 (SchemaStoreBackend::Heap(map), _) => {
4069 let mut visitor = visitor;
4070 for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4071 if visitor(key, snapshot)?.should_stop() {
4072 break;
4073 }
4074 }
4075 }
4076 (
4077 SchemaStoreBackend::Journaled {
4078 canonical,
4079 live,
4080 tombstones,
4081 ..
4082 },
4083 IdentityStateStorageView::Effective,
4084 ) => Self::visit_journaled_raw_snapshot_range(
4085 canonical,
4086 live,
4087 tombstones,
4088 bounds,
4089 Direction::Asc,
4090 visitor,
4091 )?,
4092 (
4093 SchemaStoreBackend::Journaled { canonical, .. },
4094 IdentityStateStorageView::Canonical,
4095 ) => {
4096 let mut visitor = visitor;
4097 for entry in canonical.range((bounds.0, bounds.1)) {
4098 if visitor(entry.key(), &entry.value())?.should_stop() {
4099 break;
4100 }
4101 }
4102 }
4103 }
4104 Ok(())
4105 }
4106
4107 #[cfg(test)]
4108 #[must_use]
4109 pub(in crate::db) fn canonical_len_for_tests(&self) -> u64 {
4110 match &self.backend {
4111 SchemaStoreBackend::Journaled { canonical: map, .. } => map.len(),
4112 SchemaStoreBackend::Heap(_) => 0,
4113 }
4114 }
4115
4116 fn get_raw_snapshot_for_backend(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
4117 let SchemaStoreBackend::Journaled {
4118 canonical,
4119 live,
4120 tombstones,
4121 ..
4122 } = &self.backend
4123 else {
4124 return None;
4125 };
4126
4127 if tombstones.contains(key) {
4128 return None;
4129 }
4130 live.get(key).cloned().or_else(|| canonical.get(key))
4131 }
4132
4133 fn visit_journaled_raw_snapshot_range<E>(
4134 canonical: &StableBTreeMap<
4135 RawSchemaKey,
4136 RawSchemaSnapshot,
4137 VirtualMemory<DefaultMemoryImpl>,
4138 >,
4139 live: &StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
4140 tombstones: &BTreeSet<RawSchemaKey>,
4141 bounds: (RangeBound<RawSchemaKey>, RangeBound<RawSchemaKey>),
4142 direction: Direction,
4143 mut visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4144 ) -> Result<(), E> {
4145 match direction {
4146 Direction::Asc => {
4147 for entry in ordered_overlay_entries(
4148 canonical.range((bounds.0, bounds.1)),
4149 live.range((bounds.0, bounds.1)),
4150 Direction::Asc,
4151 |entry| entry.key(),
4152 |entry| entry.0,
4153 tombstones,
4154 ) {
4155 let visit = match entry {
4156 OrderedOverlayEntry::Canonical(canonical_entry) => {
4157 visitor(canonical_entry.key(), &canonical_entry.value())?
4158 }
4159 OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4160 };
4161 if visit.should_stop() {
4162 return Ok(());
4163 }
4164 }
4165 }
4166 Direction::Desc => {
4167 for entry in ordered_overlay_entries(
4168 canonical.range((bounds.0, bounds.1)).rev(),
4169 live.range((bounds.0, bounds.1)).rev(),
4170 Direction::Desc,
4171 |entry| entry.key(),
4172 |entry| entry.0,
4173 tombstones,
4174 ) {
4175 let visit = match entry {
4176 OrderedOverlayEntry::Canonical(canonical_entry) => {
4177 visitor(canonical_entry.key(), &canonical_entry.value())?
4178 }
4179 OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4180 };
4181 if visit.should_stop() {
4182 return Ok(());
4183 }
4184 }
4185 }
4186 }
4187
4188 Ok(())
4189 }
4190}
4191
4192fn map_schema_publication_error(error: AcceptedSchemaPublicationError) -> InternalError {
4193 match error {
4194 AcceptedSchemaPublicationError::StaleSchemaRevision { .. }
4195 | AcceptedSchemaPublicationError::RevisionExhausted => InternalError::store_unsupported(),
4196 AcceptedSchemaPublicationError::InvalidCandidate => InternalError::store_invariant(),
4197 AcceptedSchemaPublicationError::CorruptRootSlots => InternalError::store_corruption(),
4198 }
4199}
4200
4201fn derive_data_allocation_metadata(
4202 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4203) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4204 let mut max_version = SchemaVersion::initial();
4205 let mut hasher = new_hash_sha256();
4206 write_hash_tag_u8(&mut hasher, SCHEMA_STORE_DATA_ALLOCATION_FINGERPRINT_DOMAIN);
4207
4208 for (entity, (_, snapshot)) in latest_by_entity {
4209 let persisted = snapshot.decode_persisted_snapshot()?;
4210 if persisted.version() > max_version {
4211 max_version = persisted.version();
4212 }
4213
4214 let data_projection = PersistedSchemaSnapshot::new_with_primary_key_fields_and_indexes(
4215 persisted.version(),
4216 persisted.entity_path().to_string(),
4217 persisted.entity_name().to_string(),
4218 persisted.primary_key_field_ids().to_vec(),
4219 persisted.row_layout().clone(),
4220 persisted.fields().to_vec(),
4221 Vec::new(),
4222 );
4223 let constraint_catalog = crate::db::schema::AcceptedConstraintCatalog::initial(
4224 data_projection.fields(),
4225 data_projection.indexes(),
4226 data_projection.relations(),
4227 )
4228 .map_err(|_| InternalError::store_invariant())?;
4229 let data_projection = data_projection.with_constraint_catalog(constraint_catalog);
4230 let encoded = encode_persisted_schema_snapshot(&data_projection)?;
4231
4232 write_hash_u64(&mut hasher, entity.value());
4233 write_hash_u32(&mut hasher, persisted.version().get());
4234 write_hash_len_u32(&mut hasher, encoded.len());
4235 hasher.update(encoded);
4236 }
4237
4238 Ok(finalize_schema_metadata(
4239 max_version,
4240 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4241 hasher,
4242 latest_by_entity.len(),
4243 ))
4244}
4245
4246fn derive_index_allocation_metadata(
4247 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4248) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4249 let mut max_version = SchemaVersion::initial();
4250 let mut hasher = new_hash_sha256();
4251 write_hash_tag_u8(
4252 &mut hasher,
4253 SCHEMA_STORE_INDEX_ALLOCATION_FINGERPRINT_DOMAIN,
4254 );
4255
4256 for (entity, (_, snapshot)) in latest_by_entity {
4257 let persisted = snapshot.decode_persisted_snapshot()?;
4258 if persisted.version() > max_version {
4259 max_version = persisted.version();
4260 }
4261
4262 write_hash_u64(&mut hasher, entity.value());
4263 write_hash_u32(&mut hasher, persisted.version().get());
4264 write_hash_len_u32(&mut hasher, persisted.indexes().len());
4265 for index in persisted.indexes() {
4266 write_hash_u32(&mut hasher, u32::from(index.ordinal()));
4267 write_hash_str_u32(&mut hasher, index.name());
4268 write_hash_str_u32(&mut hasher, index.store());
4269 write_hash_tag_u8(&mut hasher, u8::from(index.unique()));
4270 write_hash_str_u32(&mut hasher, persisted_index_origin_name(index.origin()));
4271 match index.predicate_sql() {
4272 Some(predicate_sql) => {
4273 write_hash_tag_u8(&mut hasher, 1);
4274 write_hash_str_u32(&mut hasher, predicate_sql);
4275 }
4276 None => write_hash_tag_u8(&mut hasher, 0),
4277 }
4278 hash_persisted_index_key(&mut hasher, index.key());
4279 }
4280 }
4281
4282 Ok(finalize_schema_metadata(
4283 max_version,
4284 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4285 hasher,
4286 latest_by_entity.len(),
4287 ))
4288}
4289
4290fn derive_schema_catalog_metadata(
4291 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4292) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4293 let mut max_version = SchemaVersion::initial();
4294 let mut hasher = new_hash_sha256();
4295 write_hash_tag_u8(&mut hasher, SCHEMA_STORE_CATALOG_FINGERPRINT_DOMAIN);
4296
4297 for (entity, (version, snapshot)) in latest_by_entity {
4298 let persisted = snapshot.decode_persisted_snapshot()?;
4299 if persisted.version() > max_version {
4300 max_version = persisted.version();
4301 }
4302
4303 write_hash_u64(&mut hasher, entity.value());
4304 write_hash_u32(&mut hasher, version.get());
4305 write_hash_len_u32(&mut hasher, snapshot.as_bytes().len());
4306 hasher.update(snapshot.as_bytes());
4307 }
4308
4309 Ok(finalize_schema_metadata(
4310 max_version,
4311 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4312 hasher,
4313 latest_by_entity.len(),
4314 ))
4315}
4316
4317fn finalize_schema_metadata(
4318 schema_version: SchemaVersion,
4319 schema_fingerprint_method_version: u8,
4320 hasher: sha2::Sha256,
4321 entity_count: usize,
4322) -> SchemaStoreCatalogMetadata {
4323 let digest = finalize_hash_sha256(hasher);
4324 let mut schema_fingerprint = [0u8; 16];
4325 schema_fingerprint.copy_from_slice(&digest[..16]);
4326
4327 SchemaStoreCatalogMetadata::new(
4328 schema_version,
4329 schema_fingerprint_method_version,
4330 schema_fingerprint,
4331 u64::try_from(entity_count).unwrap_or(u64::MAX),
4332 )
4333}
4334
4335fn hash_persisted_index_key(hasher: &mut sha2::Sha256, key: &PersistedIndexKeySnapshot) {
4336 match key {
4337 PersistedIndexKeySnapshot::FieldPath(paths) => {
4338 write_hash_tag_u8(hasher, 1);
4339 write_hash_len_u32(hasher, paths.len());
4340 for path in paths {
4341 hash_persisted_index_field_path(hasher, path);
4342 }
4343 }
4344 PersistedIndexKeySnapshot::Items(items) => {
4345 write_hash_tag_u8(hasher, 2);
4346 write_hash_len_u32(hasher, items.len());
4347 for item in items {
4348 match item {
4349 PersistedIndexKeyItemSnapshot::FieldPath(path) => {
4350 write_hash_tag_u8(hasher, 1);
4351 hash_persisted_index_field_path(hasher, path);
4352 }
4353 PersistedIndexKeyItemSnapshot::Expression(expression) => {
4354 write_hash_tag_u8(hasher, 2);
4355 write_hash_str_u32(hasher, persisted_expression_op_name(expression.op()));
4356 hash_persisted_index_field_path(hasher, expression.source());
4357 hash_accepted_field_kind(hasher, expression.input_kind());
4358 hash_accepted_field_kind(hasher, expression.output_kind());
4359 write_hash_str_u32(hasher, expression.canonical_text());
4360 }
4361 }
4362 }
4363 }
4364 }
4365}
4366
4367fn hash_persisted_index_field_path(
4368 hasher: &mut sha2::Sha256,
4369 path: &crate::db::schema::PersistedIndexFieldPathSnapshot,
4370) {
4371 write_hash_u32(hasher, path.field_id().get());
4372 write_hash_u32(hasher, u32::from(path.slot().get()));
4373 write_hash_len_u32(hasher, path.path().len());
4374 for segment in path.path() {
4375 write_hash_str_u32(hasher, segment);
4376 }
4377 hash_accepted_field_kind(hasher, path.kind());
4378 write_hash_tag_u8(hasher, u8::from(path.nullable()));
4379}
4380
4381fn hash_accepted_field_kind(hasher: &mut sha2::Sha256, kind: &AcceptedFieldKind) {
4382 match kind {
4383 AcceptedFieldKind::Account => write_hash_tag_u8(hasher, 1),
4384 AcceptedFieldKind::Blob { max_len } => {
4385 write_hash_tag_u8(hasher, 2);
4386 hash_optional_u32(hasher, *max_len);
4387 }
4388 AcceptedFieldKind::Bool => {
4389 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_BOOL);
4390 }
4391 AcceptedFieldKind::Date => write_hash_tag_u8(hasher, 4),
4392 AcceptedFieldKind::Decimal { scale } => {
4393 write_hash_tag_u8(hasher, 5);
4394 write_hash_u32(hasher, *scale);
4395 }
4396 AcceptedFieldKind::Duration => write_hash_tag_u8(hasher, 6),
4397 AcceptedFieldKind::Enum { type_id } => {
4398 write_hash_tag_u8(hasher, 7);
4399 write_hash_u32(hasher, type_id.get());
4400 }
4401 AcceptedFieldKind::Float32 => write_hash_tag_u8(hasher, 8),
4402 AcceptedFieldKind::Float64 => write_hash_tag_u8(hasher, 9),
4403 AcceptedFieldKind::Int8 => write_hash_tag_u8(hasher, 10),
4404 AcceptedFieldKind::Int16 => write_hash_tag_u8(hasher, 11),
4405 AcceptedFieldKind::Int32 => write_hash_tag_u8(hasher, 12),
4406 AcceptedFieldKind::Int64 => write_hash_tag_u8(hasher, 13),
4407 AcceptedFieldKind::Int128 => write_hash_tag_u8(hasher, 14),
4408 AcceptedFieldKind::IntBig { max_bytes } => {
4409 write_hash_tag_u8(hasher, 15);
4410 write_hash_u32(hasher, *max_bytes);
4411 }
4412 AcceptedFieldKind::Principal => write_hash_tag_u8(hasher, 16),
4413 AcceptedFieldKind::Subaccount => write_hash_tag_u8(hasher, 17),
4414 AcceptedFieldKind::Text { max_len } => {
4415 write_hash_tag_u8(hasher, 18);
4416 hash_optional_u32(hasher, *max_len);
4417 }
4418 AcceptedFieldKind::Timestamp => write_hash_tag_u8(hasher, 19),
4419 AcceptedFieldKind::Nat8 => write_hash_tag_u8(hasher, 20),
4420 AcceptedFieldKind::Nat16 => write_hash_tag_u8(hasher, 21),
4421 AcceptedFieldKind::Nat32 => write_hash_tag_u8(hasher, 22),
4422 AcceptedFieldKind::Nat64 => write_hash_tag_u8(hasher, 23),
4423 AcceptedFieldKind::Nat128 => write_hash_tag_u8(hasher, 24),
4424 AcceptedFieldKind::NatBig { max_bytes } => {
4425 write_hash_tag_u8(hasher, 25);
4426 write_hash_u32(hasher, *max_bytes);
4427 }
4428 AcceptedFieldKind::Ulid => write_hash_tag_u8(hasher, 26),
4429 AcceptedFieldKind::Unit => write_hash_tag_u8(hasher, 27),
4430 AcceptedFieldKind::Relation {
4431 target_path,
4432 target_entity_name,
4433 target_entity_tag,
4434 target_store_path,
4435 key_kind,
4436 } => {
4437 write_hash_tag_u8(hasher, 28);
4438 write_hash_str_u32(hasher, target_path);
4439 write_hash_str_u32(hasher, target_entity_name);
4440 write_hash_u64(hasher, target_entity_tag.value());
4441 write_hash_str_u32(hasher, target_store_path);
4442 hash_accepted_field_kind(hasher, key_kind);
4443 }
4444 AcceptedFieldKind::List(inner) => {
4445 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_LIST);
4446 hash_accepted_field_kind(hasher, inner);
4447 }
4448 AcceptedFieldKind::Set(inner) => {
4449 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_SET);
4450 hash_accepted_field_kind(hasher, inner);
4451 }
4452 AcceptedFieldKind::Map { key, value } => {
4453 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_MAP);
4454 hash_accepted_field_kind(hasher, key);
4455 hash_accepted_field_kind(hasher, value);
4456 }
4457 AcceptedFieldKind::Composite { type_id } => {
4458 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_COMPOSITE);
4459 write_hash_u32(hasher, type_id.get());
4460 }
4461 AcceptedFieldKind::U256 => write_hash_tag_u8(hasher, 33),
4462 }
4463}
4464
4465fn hash_optional_u32(hasher: &mut sha2::Sha256, value: Option<u32>) {
4466 match value {
4467 Some(value) => {
4468 write_hash_tag_u8(hasher, 1);
4469 write_hash_u32(hasher, value);
4470 }
4471 None => write_hash_tag_u8(hasher, 0),
4472 }
4473}
4474
4475const fn persisted_index_origin_name(
4476 origin: crate::db::schema::PersistedIndexOrigin,
4477) -> &'static str {
4478 match origin {
4479 crate::db::schema::PersistedIndexOrigin::Generated => "generated",
4480 crate::db::schema::PersistedIndexOrigin::SqlDdl => "sql_ddl",
4481 }
4482}
4483
4484const fn persisted_expression_op_name(
4485 op: crate::db::schema::PersistedIndexExpressionOp,
4486) -> &'static str {
4487 match op {
4488 crate::db::schema::PersistedIndexExpressionOp::Lower => "lower",
4489 crate::db::schema::PersistedIndexExpressionOp::Upper => "upper",
4490 crate::db::schema::PersistedIndexExpressionOp::Trim => "trim",
4491 crate::db::schema::PersistedIndexExpressionOp::LowerTrim => "lower_trim",
4492 crate::db::schema::PersistedIndexExpressionOp::Date => "date",
4493 crate::db::schema::PersistedIndexExpressionOp::Year => "year",
4494 crate::db::schema::PersistedIndexExpressionOp::Month => "month",
4495 crate::db::schema::PersistedIndexExpressionOp::Day => "day",
4496 }
4497}
4498
4499#[cfg(test)]
4504mod tests;