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