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