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
863#[derive(Clone)]
865pub(in crate::db) struct PreparedSchemaPositionPublication {
866 keys: Vec<RawSchemaKey>,
867 position: JournalOverlayPosition,
868}
869
870pub(in crate::db) struct PreparedSchemaPositionRetirement {
872 entries: Vec<(RawSchemaKey, PositionedOverlayRetirement)>,
873}
874
875pub(in crate::db) struct PreparedCardinalityCountWrites {
877 slot: CardinalityCountSlot,
878 generation: CardinalityGenerationId,
879 entries: Vec<(RawSchemaKey, RawSchemaSnapshot)>,
880 new_count_keys: u64,
881}
882
883impl PreparedCardinalityCountWrites {
884 #[must_use]
885 pub(in crate::db) const fn new_count_keys(&self) -> u64 {
886 self.new_count_keys
887 }
888}
889
890pub(in crate::db) struct PreparedCardinalityBuildPage {
892 count_entries: Vec<(RawSchemaKey, RawSchemaSnapshot)>,
893 cursor: (RawSchemaKey, RawSchemaSnapshot),
894}
895
896pub(in crate::db) struct PreparedCardinalityMaintenance {
898 count_entries: Vec<(RawSchemaKey, Option<RawSchemaSnapshot>)>,
899 header: (RawSchemaKey, RawSchemaSnapshot),
900}
901
902#[derive(Clone, Copy)]
903enum IdentityStateWriteTarget {
904 Durable,
905 Materialized,
906 Canonical,
907}
908
909impl SchemaStore {
910 #[must_use]
912 pub const fn init_heap() -> Self {
913 Self {
914 backend: SchemaStoreBackend::Heap(StdBTreeMap::new()),
915 accepted_bundle_cache: RefCell::new(None),
916 cardinality_header_cache: RefCell::new(None),
917 accepted_catalog_scope: OnceCell::new(),
918 }
919 }
920
921 #[must_use]
926 pub fn init_journaled(memory: RuntimeMemory<DefaultMemoryImpl>) -> Self {
927 Self {
928 backend: SchemaStoreBackend::Journaled {
929 canonical: StableBTreeMap::init(memory),
930 live: StdBTreeMap::new(),
931 tombstones: BTreeSet::new(),
932 positions: PositionedOverlayMetadata::new(),
933 },
934 accepted_bundle_cache: RefCell::new(None),
935 cardinality_header_cache: RefCell::new(None),
936 accepted_catalog_scope: OnceCell::new(),
937 }
938 }
939
940 pub(in crate::db) fn cardinality_generation_header(
942 &self,
943 ) -> Result<Option<CardinalityGenerationHeader>, InternalError> {
944 let key = RawSchemaKey::from_cardinality_generation_header();
945 let raw = self.get_canonical_raw_value(&key)?;
946 if let Some(raw) = raw.as_ref() {
947 self.decode_cardinality_header_cached(raw).map(Some)
948 } else {
949 self.cardinality_header_cache
950 .try_borrow_mut()
951 .map_err(|_| InternalError::store_invariant())?
952 .take();
953 Ok(None)
954 }
955 }
956
957 pub(in crate::db) fn cardinality_build_cursor(
959 &self,
960 ) -> Result<Option<CardinalityBuildCursor>, InternalError> {
961 let key = RawSchemaKey::from_cardinality_build_cursor();
962 self.get_canonical_raw_value(&key)?
963 .map(|raw| CardinalityBuildCursor::decode(raw.as_bytes()))
964 .transpose()
965 }
966
967 pub(in crate::db) fn cardinality_generation_control(
969 &self,
970 ) -> Result<
971 (
972 Option<CardinalityGenerationHeader>,
973 Option<CardinalityBuildCursor>,
974 ),
975 InternalError,
976 > {
977 let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
978 return Err(InternalError::store_invariant());
979 };
980 let header_key = RawSchemaKey::from_cardinality_generation_header();
981 let cursor_key = RawSchemaKey::from_cardinality_build_cursor();
982 let mut header = None;
983 let mut cursor = None;
984 for entry in canonical.range(header_key..=cursor_key) {
985 if *entry.key() == header_key {
986 header = Some(self.decode_cardinality_header_cached(&entry.value())?);
987 } else if *entry.key() == cursor_key {
988 cursor = Some(CardinalityBuildCursor::decode(entry.value().as_bytes())?);
989 } else {
990 return Err(InternalError::store_corruption());
991 }
992 }
993 Ok((header, cursor))
994 }
995
996 pub(in crate::db) fn cardinality_generation_lifecycle_control(
1001 &self,
1002 ) -> Result<(Option<CardinalityGenerationHeader>, bool), InternalError> {
1003 let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
1004 return Err(InternalError::store_invariant());
1005 };
1006 let header_key = RawSchemaKey::from_cardinality_generation_header();
1007 let cursor_key = RawSchemaKey::from_cardinality_build_cursor();
1008 let mut header = None;
1009 let mut cursor_present = false;
1010 for entry in canonical.range(header_key..=cursor_key) {
1011 if *entry.key() == header_key {
1012 header = Some(self.decode_cardinality_header_cached(&entry.value())?);
1013 } else if *entry.key() == cursor_key {
1014 cursor_present = true;
1015 } else {
1016 return Err(InternalError::store_corruption());
1017 }
1018 }
1019
1020 Ok((header, cursor_present))
1021 }
1022
1023 fn decode_cardinality_header_cached(
1024 &self,
1025 raw: &RawSchemaSnapshot,
1026 ) -> Result<CardinalityGenerationHeader, InternalError> {
1027 if let Some((_, header)) = self
1028 .cardinality_header_cache
1029 .try_borrow()
1030 .map_err(|_| InternalError::store_invariant())?
1031 .as_ref()
1032 .filter(|(bytes, _)| bytes.as_slice() == raw.as_bytes())
1033 {
1034 return Ok(*header);
1035 }
1036 let header = CardinalityGenerationHeader::decode(raw.as_bytes())?;
1037 *self
1038 .cardinality_header_cache
1039 .try_borrow_mut()
1040 .map_err(|_| InternalError::store_invariant())? =
1041 Some((raw.as_bytes().to_vec(), header));
1042 Ok(header)
1043 }
1044
1045 pub(in crate::db) fn cardinality_storage_is_pristine(&self) -> Result<bool, InternalError> {
1047 if self.cardinality_generation_header()?.is_some()
1048 || self.cardinality_build_cursor()?.is_some()
1049 {
1050 return Ok(false);
1051 }
1052 Ok(
1053 self.cardinality_count_slot_is_empty(CardinalityCountSlot::A)?
1054 && self.cardinality_count_slot_is_empty(CardinalityCountSlot::B)?,
1055 )
1056 }
1057
1058 pub(in crate::db) fn write_cardinality_generation_header(
1060 &mut self,
1061 header: CardinalityGenerationHeader,
1062 ) -> Result<(), InternalError> {
1063 self.insert_canonical_raw_value(
1064 RawSchemaKey::from_cardinality_generation_header(),
1065 header.encode(),
1066 )
1067 }
1068
1069 pub(in crate::db) fn restart_cardinality_generation(
1075 &mut self,
1076 current: CardinalityGenerationHeader,
1077 source: CardinalitySourceIdentity,
1078 ) -> Result<CardinalityGenerationHeader, InternalError> {
1079 if self.cardinality_generation_header()? != Some(current) {
1080 return Err(InternalError::store_corruption());
1081 }
1082 if current.validate_source(source).is_ok() {
1083 return Err(InternalError::store_invariant());
1084 }
1085 if let Some(cursor) = self.cardinality_build_cursor()? {
1086 cursor.validate_header(current)?;
1087 }
1088 let next = CardinalityGenerationHeader::new(
1089 current.generation().checked_next()?,
1090 CardinalityGenerationState::Building,
1091 current.slot().alternate(),
1092 source,
1093 );
1094 let encoded = RawSchemaSnapshot::from_encoded_control_record(next.encode());
1095 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1096 return Err(InternalError::store_invariant());
1097 };
1098 canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1099 canonical.remove(&RawSchemaKey::from_cardinality_build_cursor());
1100 Ok(next)
1101 }
1102
1103 pub(in crate::db) fn publish_empty_cardinality_generation(
1105 &mut self,
1106 candidate: &EmptyCardinalityReadyCandidate,
1107 ) -> Result<CardinalityGenerationHeader, InternalError> {
1108 if !self.cardinality_storage_is_pristine()? {
1109 return Err(InternalError::store_corruption());
1110 }
1111 let ready = CardinalityGenerationHeader::new(
1112 CardinalityGenerationId::INITIAL,
1113 CardinalityGenerationState::Ready,
1114 CardinalityCountSlot::A,
1115 candidate.source(),
1116 );
1117 let encoded = RawSchemaSnapshot::from_encoded_control_record(ready.encode());
1118 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1119 return Err(InternalError::store_invariant());
1120 };
1121 canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1122 Ok(ready)
1123 }
1124
1125 pub(in crate::db) fn publish_ready_cardinality_generation(
1127 &mut self,
1128 candidate: &CardinalityReadyCandidate,
1129 source: CardinalitySourceIdentity,
1130 ) -> Result<CardinalityGenerationHeader, InternalError> {
1131 let building = candidate.header();
1132 candidate.cursor().validate_header(building)?;
1133 if building.state() != CardinalityGenerationState::Building
1134 || building.validate_source(source).is_err()
1135 || self.cardinality_generation_header()? != Some(building)
1136 || self.cardinality_build_cursor()?.as_ref() != Some(candidate.cursor())
1137 {
1138 return Err(InternalError::store_corruption());
1139 }
1140 let ready = CardinalityGenerationHeader::new(
1141 building.generation(),
1142 CardinalityGenerationState::Ready,
1143 building.slot(),
1144 building.source(),
1145 );
1146 let encoded = RawSchemaSnapshot::from_encoded_control_record(ready.encode());
1147 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1148 return Err(InternalError::store_invariant());
1149 };
1150 canonical.insert(RawSchemaKey::from_cardinality_generation_header(), encoded);
1151 canonical.remove(&RawSchemaKey::from_cardinality_build_cursor());
1152 Ok(ready)
1153 }
1154
1155 pub(in crate::db) fn clear_cardinality_count_slot_page(
1161 &mut self,
1162 header: CardinalityGenerationHeader,
1163 initial_cursor: &CardinalityBuildCursor,
1164 limit: usize,
1165 ) -> Result<bool, InternalError> {
1166 if limit == 0 {
1167 return Err(InternalError::store_invariant());
1168 }
1169 initial_cursor.validate_header(header)?;
1170 let encoded_cursor = initial_cursor.encode()?;
1171 if self.cardinality_generation_header()? != Some(header)
1172 || self.cardinality_build_cursor()?.is_some()
1173 {
1174 return Err(InternalError::store_corruption());
1175 }
1176 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1177 return Err(InternalError::store_invariant());
1178 };
1179 let bounds = RawSchemaKey::cardinality_count_range_bounds(header.slot());
1180 let collect_limit = limit
1181 .checked_add(1)
1182 .ok_or_else(InternalError::store_unsupported)?;
1183 let mut keys = Vec::new();
1184 keys.try_reserve_exact(collect_limit)
1185 .map_err(|_| InternalError::store_unsupported())?;
1186 for entry in canonical.range(bounds).take(collect_limit) {
1187 keys.push(*entry.key());
1188 }
1189 let has_more = keys.len() > limit;
1190 for key in keys.into_iter().take(limit) {
1191 canonical.remove(&key);
1192 }
1193 if !has_more {
1194 canonical.insert(
1195 RawSchemaKey::from_cardinality_build_cursor(),
1196 RawSchemaSnapshot::from_encoded_control_record(encoded_cursor),
1197 );
1198 }
1199 Ok(has_more)
1200 }
1201
1202 pub(in crate::db) fn prepare_cardinality_count_increments(
1204 &self,
1205 slot: CardinalityCountSlot,
1206 generation: CardinalityGenerationId,
1207 increments: &[(CardinalityCountDigest, u64)],
1208 ) -> Result<PreparedCardinalityCountWrites, InternalError> {
1209 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1210 return Err(InternalError::store_invariant());
1211 }
1212 let mut entries = Vec::new();
1213 entries
1214 .try_reserve_exact(increments.len())
1215 .map_err(|_| InternalError::store_unsupported())?;
1216 let mut physical_keys = StdBTreeMap::new();
1217 let mut new_count_keys = 0_u64;
1218 for (digest, increment) in increments {
1219 if *increment == 0 {
1220 return Err(InternalError::store_invariant());
1221 }
1222 let key = RawSchemaKey::from_cardinality_count(slot, *digest);
1223 if let Some(previous_digest) = physical_keys.insert(key, *digest) {
1224 return Err(if previous_digest == *digest {
1225 InternalError::store_invariant()
1226 } else {
1227 InternalError::store_corruption()
1228 });
1229 }
1230 let current = self.get_canonical_raw_value(&key)?;
1231 let count = if let Some(raw) = current {
1232 let record = CardinalityCountRecord::decode(raw.as_bytes())?;
1233 record
1234 .validate_identity(generation, *digest)
1235 .map_err(|_| InternalError::store_corruption())?
1236 .checked_add(*increment)
1237 .ok_or_else(InternalError::store_unsupported)?
1238 } else {
1239 new_count_keys = new_count_keys
1240 .checked_add(1)
1241 .ok_or_else(InternalError::store_unsupported)?;
1242 *increment
1243 };
1244 let record = CardinalityCountRecord::new(generation, *digest, count)?;
1245 entries.push((
1246 key,
1247 RawSchemaSnapshot::from_encoded_control_record(record.encode().to_vec()),
1248 ));
1249 }
1250 Ok(PreparedCardinalityCountWrites {
1251 slot,
1252 generation,
1253 entries,
1254 new_count_keys,
1255 })
1256 }
1257
1258 pub(in crate::db) fn prepare_cardinality_build_page(
1260 &self,
1261 header: CardinalityGenerationHeader,
1262 current_cursor: &CardinalityBuildCursor,
1263 counts: PreparedCardinalityCountWrites,
1264 next_cursor: &CardinalityBuildCursor,
1265 ) -> Result<PreparedCardinalityBuildPage, InternalError> {
1266 current_cursor.validate_header(header)?;
1267 next_cursor.validate_header(header)?;
1268 if counts.slot != header.slot() || counts.generation != header.generation() {
1269 return Err(InternalError::store_invariant());
1270 }
1271 if self.cardinality_generation_header()? != Some(header)
1272 || self.cardinality_build_cursor()?.as_ref() != Some(current_cursor)
1273 {
1274 return Err(InternalError::store_corruption());
1275 }
1276 let encoded_cursor = next_cursor.encode()?;
1277 Ok(PreparedCardinalityBuildPage {
1278 count_entries: counts.entries,
1279 cursor: (
1280 RawSchemaKey::from_cardinality_build_cursor(),
1281 RawSchemaSnapshot::from_encoded_control_record(encoded_cursor),
1282 ),
1283 })
1284 }
1285
1286 pub(in crate::db) fn apply_prepared_cardinality_build_page(
1288 &mut self,
1289 prepared: PreparedCardinalityBuildPage,
1290 ) -> Result<(), InternalError> {
1291 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1292 return Err(InternalError::store_invariant());
1293 };
1294 for (key, value) in prepared.count_entries {
1295 canonical.insert(key, value);
1296 }
1297 canonical.insert(prepared.cursor.0, prepared.cursor.1);
1298 Ok(())
1299 }
1300
1301 pub(in crate::db) fn cardinality_count(
1303 &self,
1304 slot: CardinalityCountSlot,
1305 generation: CardinalityGenerationId,
1306 digest: CardinalityCountDigest,
1307 ) -> Result<Option<u64>, InternalError> {
1308 let key = RawSchemaKey::from_cardinality_count(slot, digest);
1309 self.get_canonical_raw_value(&key)?
1310 .map(|raw| {
1311 CardinalityCountRecord::decode(raw.as_bytes())?
1312 .validate_identity(generation, digest)
1313 .map_err(|_| InternalError::store_corruption())
1314 })
1315 .transpose()
1316 }
1317
1318 pub(in crate::db) fn prepare_cardinality_maintenance(
1320 &self,
1321 current: CardinalityGenerationHeader,
1322 current_source: CardinalitySourceIdentity,
1323 next_source: CardinalitySourceIdentity,
1324 changes: &[(CardinalityCountDigest, i64)],
1325 ) -> Result<PreparedCardinalityMaintenance, InternalError> {
1326 if current.state() != CardinalityGenerationState::Ready
1327 || current.validate_source(current_source).is_err()
1328 || self.cardinality_generation_header()? != Some(current)
1329 || self.cardinality_build_cursor()?.is_some()
1330 {
1331 return Err(InternalError::store_corruption());
1332 }
1333 let next = CardinalityGenerationHeader::new(
1334 current.generation(),
1335 CardinalityGenerationState::Ready,
1336 current.slot(),
1337 next_source,
1338 );
1339 let mut count_entries = Vec::new();
1340 count_entries
1341 .try_reserve_exact(changes.len())
1342 .map_err(|_| InternalError::store_unsupported())?;
1343 let mut physical_keys = StdBTreeMap::new();
1344 for (digest, delta) in changes {
1345 if *delta == 0 {
1346 return Err(InternalError::store_invariant());
1347 }
1348 let key = RawSchemaKey::from_cardinality_count(current.slot(), *digest);
1349 if let Some(previous_digest) = physical_keys.insert(key, *digest) {
1350 return Err(if previous_digest == *digest {
1351 InternalError::store_invariant()
1352 } else {
1353 InternalError::store_corruption()
1354 });
1355 }
1356 let base = self
1357 .cardinality_count(current.slot(), current.generation(), *digest)?
1358 .unwrap_or(0);
1359 let count = if *delta > 0 {
1360 base.checked_add(
1361 u64::try_from(*delta).map_err(|_| InternalError::store_invariant())?,
1362 )
1363 } else {
1364 base.checked_sub(delta.unsigned_abs())
1365 }
1366 .ok_or_else(InternalError::store_corruption)?;
1367 let value = if count == 0 {
1368 None
1369 } else {
1370 Some(RawSchemaSnapshot::from_encoded_control_record(
1371 CardinalityCountRecord::new(current.generation(), *digest, count)?
1372 .encode()
1373 .to_vec(),
1374 ))
1375 };
1376 count_entries.push((key, value));
1377 }
1378 Ok(PreparedCardinalityMaintenance {
1379 count_entries,
1380 header: (
1381 RawSchemaKey::from_cardinality_generation_header(),
1382 RawSchemaSnapshot::from_encoded_control_record(next.encode()),
1383 ),
1384 })
1385 }
1386
1387 pub(in crate::db) fn apply_prepared_cardinality_maintenance(
1389 &mut self,
1390 prepared: PreparedCardinalityMaintenance,
1391 ) -> Result<(), InternalError> {
1392 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1393 return Err(InternalError::store_invariant());
1394 };
1395 for (key, value) in prepared.count_entries {
1396 if let Some(value) = value {
1397 canonical.insert(key, value);
1398 } else {
1399 canonical.remove(&key);
1400 }
1401 }
1402 canonical.insert(prepared.header.0, prepared.header.1);
1403 Ok(())
1404 }
1405
1406 pub(in crate::db) fn cardinality_count_slot_is_empty(
1407 &self,
1408 slot: CardinalityCountSlot,
1409 ) -> Result<bool, InternalError> {
1410 let SchemaStoreBackend::Journaled { canonical, .. } = &self.backend else {
1411 return Err(InternalError::store_invariant());
1412 };
1413 Ok(canonical
1414 .range(RawSchemaKey::cardinality_count_range_bounds(slot))
1415 .next()
1416 .is_none())
1417 }
1418
1419 pub(in crate::db) fn preflight_fold_recovered_journal(&self) -> Result<(), InternalError> {
1421 match self.backend {
1422 SchemaStoreBackend::Journaled { .. } => Ok(()),
1423 SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
1424 }
1425 }
1426
1427 fn prepare_identity_state_transition(
1428 &self,
1429 incarnation: DatabaseIncarnationId,
1430 candidate: &CandidateSchemaRevision,
1431 view: IdentityStateStorageView,
1432 ) -> Result<IdentityStateTransition, InternalError> {
1433 let current = match view {
1434 IdentityStateStorageView::Effective => self
1435 .current_accepted_schema_bundle_ref()?
1436 .as_ref()
1437 .map(|bundle| (*bundle).clone()),
1438 IdentityStateStorageView::Canonical => {
1439 self.current_canonical_accepted_schema_bundle()?
1440 }
1441 };
1442 let inventory = self.identity_state_inventory(view)?;
1443 prepare_identity_state_transition(
1444 incarnation,
1445 current.as_ref(),
1446 candidate.bundle(),
1447 inventory,
1448 )
1449 }
1450
1451 fn validate_identity_state_closure(
1452 &self,
1453 bundle: &AcceptedSchemaRevisionBundle,
1454 ) -> Result<(), InternalError> {
1455 let inventory = self.identity_state_inventory(IdentityStateStorageView::Effective)?;
1456 validate_identity_state_closure(bundle, &inventory)
1457 }
1458
1459 pub(in crate::db) fn identity_statement_cursor(
1461 &self,
1462 database_incarnation_id: DatabaseIncarnationId,
1463 entity_tag: EntityTag,
1464 field_id: FieldId,
1465 accepted_kind: &AcceptedFieldKind,
1466 ) -> Result<IdentityStatementCursor, InternalError> {
1467 let key = RawSchemaKey::from_identity_state(entity_tag, field_id);
1468 let raw = self
1469 .get_raw_snapshot(&key)
1470 .ok_or_else(InternalError::identity_state_corruption)?;
1471 let state = decode_identity_state(raw.as_bytes())?;
1472 let owner = state.owner();
1473 if owner.database_incarnation_id() != database_incarnation_id
1474 || owner.entity_tag() != entity_tag
1475 || owner.field_id() != field_id
1476 || state.accepted_kind() != accepted_kind
1477 || state.lifecycle() != IdentityStateLifecycle::Active
1478 {
1479 return Err(InternalError::identity_state_corruption());
1480 }
1481 IdentityStatementCursor::from_active_state(&state)
1482 }
1483
1484 pub(in crate::db) fn identity_high_water_for_integrity(
1486 &self,
1487 database_incarnation_id: DatabaseIncarnationId,
1488 entity_tag: EntityTag,
1489 field_id: FieldId,
1490 accepted_kind: &AcceptedFieldKind,
1491 ) -> Result<u128, InternalError> {
1492 let key = RawSchemaKey::from_identity_state(entity_tag, field_id);
1493 let raw = self
1494 .get_raw_snapshot(&key)
1495 .ok_or_else(InternalError::identity_state_corruption)?;
1496 let state = decode_identity_state(raw.as_bytes())?;
1497 let owner = state.owner();
1498 if owner.database_incarnation_id() != database_incarnation_id
1499 || owner.entity_tag() != entity_tag
1500 || owner.field_id() != field_id
1501 || state.accepted_kind() != accepted_kind
1502 || state.lifecycle() != IdentityStateLifecycle::Active
1503 {
1504 return Err(InternalError::identity_state_corruption());
1505 }
1506 Ok(state.materialized_high_water())
1507 }
1508
1509 pub(in crate::db) fn preflight_identity_range_advance(
1511 &self,
1512 range: IdentityRangeAdvance,
1513 ) -> Result<(), InternalError> {
1514 let state =
1515 self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Effective)?;
1516 state.preflight_range_advance(range)
1517 }
1518
1519 pub(in crate::db) fn apply_identity_range_advance(
1521 &mut self,
1522 range: IdentityRangeAdvance,
1523 advance_id: IdentityAdvanceId,
1524 ) -> Result<(), InternalError> {
1525 self.apply_identity_range_advance_to(
1526 range,
1527 advance_id,
1528 IdentityStateStorageView::Effective,
1529 IdentityStateWriteTarget::Materialized,
1530 )
1531 }
1532
1533 pub(in crate::db) fn fold_identity_range_advance(
1535 &mut self,
1536 range: IdentityRangeAdvance,
1537 advance_id: IdentityAdvanceId,
1538 ) -> Result<(), InternalError> {
1539 self.apply_identity_range_advance_to(
1540 range,
1541 advance_id,
1542 IdentityStateStorageView::Canonical,
1543 IdentityStateWriteTarget::Canonical,
1544 )
1545 }
1546
1547 pub(in crate::db) fn preflight_fold_identity_range_advance(
1549 &self,
1550 range: IdentityRangeAdvance,
1551 advance_id: IdentityAdvanceId,
1552 ) -> Result<(), InternalError> {
1553 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1554 return Err(InternalError::store_invariant());
1555 }
1556 let state =
1557 self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Canonical)?;
1558 let advanced = state.apply_range_advance(range, advance_id)?;
1559 let _encoded = encode_identity_state(&advanced)?;
1560 Ok(())
1561 }
1562
1563 pub(in crate::db) fn verify_identity_range_advance(
1565 &self,
1566 range: IdentityRangeAdvance,
1567 advance_id: IdentityAdvanceId,
1568 ) -> Result<(), InternalError> {
1569 let state =
1570 self.identity_state_for_owner(range.owner(), IdentityStateStorageView::Effective)?;
1571 if state.materialized_high_water() != range.new_high_water()
1572 || state.last_applied_advance() != Some(advance_id)
1573 {
1574 return Err(InternalError::recovery_effect_verification_failed());
1575 }
1576 Ok(())
1577 }
1578
1579 pub(in crate::db) fn identity_range_commit_state(
1581 &self,
1582 range: IdentityRangeAdvance,
1583 advance_id: IdentityAdvanceId,
1584 canonical: bool,
1585 ) -> Result<IdentityRangeCommitState, InternalError> {
1586 let view = if canonical {
1587 IdentityStateStorageView::Canonical
1588 } else {
1589 IdentityStateStorageView::Effective
1590 };
1591 self.identity_state_for_owner(range.owner(), view)?
1592 .range_commit_state(range, advance_id)
1593 }
1594
1595 pub(in crate::db) fn identity_state_inventory_for_integrity(
1598 &self,
1599 incarnation: DatabaseIncarnationId,
1600 ) -> Result<Vec<IdentityState>, InternalError> {
1601 let has_accepted_bundle = self.current_accepted_schema_bundle_ref()?.is_some();
1602 let inventory = self.identity_state_inventory(IdentityStateStorageView::Effective)?;
1603 if !has_accepted_bundle && !inventory.is_empty() {
1604 return Err(InternalError::identity_state_corruption());
1605 }
1606 if inventory
1607 .values()
1608 .any(|state| state.owner().database_incarnation_id() != incarnation)
1609 {
1610 return Err(InternalError::identity_state_corruption());
1611 }
1612 Ok(inventory.into_values().collect())
1613 }
1614
1615 fn identity_state_for_owner(
1616 &self,
1617 owner: crate::db::schema::identity_state::IdentityStateOwner,
1618 view: IdentityStateStorageView,
1619 ) -> Result<IdentityState, InternalError> {
1620 let key = RawSchemaKey::from_identity_state(owner.entity_tag(), owner.field_id());
1621 let raw = match view {
1622 IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
1623 IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
1624 }
1625 .ok_or_else(InternalError::identity_state_corruption)?;
1626 let state = decode_identity_state(raw.as_bytes())?;
1627 if state.owner() != owner {
1628 return Err(InternalError::identity_state_corruption());
1629 }
1630 Ok(state)
1631 }
1632
1633 fn apply_identity_range_advance_to(
1634 &mut self,
1635 range: IdentityRangeAdvance,
1636 advance_id: IdentityAdvanceId,
1637 view: IdentityStateStorageView,
1638 target: IdentityStateWriteTarget,
1639 ) -> Result<(), InternalError> {
1640 let state = self.identity_state_for_owner(range.owner(), view)?;
1641 let advanced = state.apply_range_advance(range, advance_id)?;
1642 let key = RawSchemaKey::from_identity_state(
1643 advanced.owner().entity_tag(),
1644 advanced.owner().field_id(),
1645 );
1646 let bytes = encode_identity_state(&advanced)?;
1647 match target {
1648 IdentityStateWriteTarget::Materialized => {
1649 self.insert_raw_snapshot(
1650 key,
1651 RawSchemaSnapshot::from_encoded_control_record(bytes),
1652 );
1653 }
1654 IdentityStateWriteTarget::Canonical => {
1655 self.insert_canonical_raw_value(key, bytes)?;
1656 }
1657 IdentityStateWriteTarget::Durable => {
1658 return Err(InternalError::store_invariant());
1659 }
1660 }
1661 Ok(())
1662 }
1663
1664 fn identity_state_inventory(
1665 &self,
1666 view: IdentityStateStorageView,
1667 ) -> Result<IdentityStateInventory, InternalError> {
1668 let bounds = RawSchemaKey::all_identity_state_range_bounds();
1669 let mut inventory = StdBTreeMap::new();
1670 let mut collect = |key: &RawSchemaKey,
1671 raw: &RawSchemaSnapshot|
1672 -> Result<SchemaStoreVisit, InternalError> {
1673 if inventory.len() >= MAX_IDENTITY_STATE_RECORDS_PER_DATABASE {
1674 return Err(InternalError::identity_state_corruption());
1675 }
1676 let state = decode_identity_state(raw.as_bytes())?;
1677 let state_key = (key.entity_tag(), FieldId::new(key.version()));
1678 if !key.is_identity_state()
1679 || state.owner().entity_tag() != state_key.0
1680 || state.owner().field_id() != state_key.1
1681 || inventory.insert(state_key, state).is_some()
1682 {
1683 return Err(InternalError::identity_state_corruption());
1684 }
1685 Ok(SchemaStoreVisit::Continue)
1686 };
1687
1688 match (&self.backend, view) {
1689 (SchemaStoreBackend::Heap(map), IdentityStateStorageView::Effective) => {
1690 for (key, raw) in map.range((bounds.0, bounds.1)) {
1691 collect(key, raw)?;
1692 }
1693 }
1694 (
1695 SchemaStoreBackend::Journaled {
1696 canonical,
1697 live,
1698 tombstones,
1699 ..
1700 },
1701 IdentityStateStorageView::Effective,
1702 ) => Self::visit_journaled_raw_snapshot_range(
1703 canonical,
1704 live,
1705 tombstones,
1706 bounds,
1707 Direction::Asc,
1708 &mut collect,
1709 )?,
1710 (
1711 SchemaStoreBackend::Journaled { canonical, .. },
1712 IdentityStateStorageView::Canonical,
1713 ) => {
1714 for entry in canonical.range((bounds.0, bounds.1)) {
1715 collect(entry.key(), &entry.value())?;
1716 }
1717 }
1718 (SchemaStoreBackend::Heap(_), IdentityStateStorageView::Canonical) => {
1719 return Err(InternalError::store_invariant());
1720 }
1721 }
1722
1723 Ok(inventory)
1724 }
1725
1726 fn apply_identity_state_transition(
1727 &mut self,
1728 transition: IdentityStateTransition,
1729 target: IdentityStateWriteTarget,
1730 ) -> Result<(), InternalError> {
1731 for state in transition.into_updates() {
1732 let key = RawSchemaKey::from_identity_state(
1733 state.owner().entity_tag(),
1734 state.owner().field_id(),
1735 );
1736 let bytes = encode_identity_state(&state)?;
1737 match target {
1738 IdentityStateWriteTarget::Durable => {
1739 self.insert_durable_raw_value(key, bytes);
1740 }
1741 IdentityStateWriteTarget::Materialized => {
1742 self.insert_raw_snapshot(
1743 key,
1744 RawSchemaSnapshot::from_encoded_control_record(bytes),
1745 );
1746 }
1747 IdentityStateWriteTarget::Canonical => {
1748 self.insert_canonical_raw_value(key, bytes)?;
1749 }
1750 }
1751 }
1752 Ok(())
1753 }
1754
1755 pub(in crate::db) fn current_canonical_accepted_schema_bundle(
1756 &self,
1757 ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
1758 self.current_canonical_accepted_schema_authority()
1759 .map(|authority| authority.map(|(_, bundle)| bundle))
1760 }
1761
1762 pub(in crate::db) fn current_canonical_accepted_schema_authority(
1764 &self,
1765 ) -> Result<Option<(AcceptedSchemaRootSelection, AcceptedSchemaRevisionBundle)>, InternalError>
1766 {
1767 let Some(selection) = self.current_canonical_accepted_schema_root()? else {
1768 return Ok(None);
1769 };
1770 let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
1771 let raw = self
1772 .get_canonical_raw_value(&bundle_key)?
1773 .ok_or_else(InternalError::store_corruption)?;
1774 let bundle =
1775 decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
1776 Ok(Some((selection, bundle)))
1777 }
1778
1779 pub(in crate::db) fn current_canonical_accepted_schema_root(
1781 &self,
1782 ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
1783 let first = self.canonical_root_slot_bytes(0)?;
1784 let second = self.canonical_root_slot_bytes(1)?;
1785 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
1786 }
1787
1788 pub(in crate::db) fn current_effective_and_canonical_accepted_schema_roots(
1790 &self,
1791 ) -> Result<
1792 (
1793 Option<AcceptedSchemaRootSelection>,
1794 Option<AcceptedSchemaRootSelection>,
1795 ),
1796 InternalError,
1797 > {
1798 let SchemaStoreBackend::Journaled {
1799 canonical,
1800 live,
1801 tombstones,
1802 ..
1803 } = &self.backend
1804 else {
1805 return Err(InternalError::store_invariant());
1806 };
1807 let first_key = RawSchemaKey::from_accepted_root_slot(0)?;
1808 let second_key = RawSchemaKey::from_accepted_root_slot(1)?;
1809 let mut canonical_first = None;
1810 let mut canonical_second = None;
1811 for entry in canonical.range(first_key..=second_key) {
1812 if *entry.key() == first_key {
1813 canonical_first = Some(entry.value().clone());
1814 } else if *entry.key() == second_key {
1815 canonical_second = Some(entry.value().clone());
1816 } else {
1817 return Err(InternalError::store_corruption());
1818 }
1819 }
1820 let effective_first = if tombstones.contains(&first_key) {
1821 None
1822 } else {
1823 live.get(&first_key)
1824 .cloned()
1825 .or_else(|| canonical_first.clone())
1826 };
1827 let effective_second = if tombstones.contains(&second_key) {
1828 None
1829 } else {
1830 live.get(&second_key)
1831 .cloned()
1832 .or_else(|| canonical_second.clone())
1833 };
1834 let effective_first = effective_first.map(RawSchemaSnapshot::into_bytes);
1835 let effective_second = effective_second.map(RawSchemaSnapshot::into_bytes);
1836 let canonical_first = canonical_first.map(RawSchemaSnapshot::into_bytes);
1837 let canonical_second = canonical_second.map(RawSchemaSnapshot::into_bytes);
1838 Ok((
1839 select_current_accepted_schema_root([
1840 effective_first.as_deref(),
1841 effective_second.as_deref(),
1842 ])?,
1843 select_current_accepted_schema_root([
1844 canonical_first.as_deref(),
1845 canonical_second.as_deref(),
1846 ])?,
1847 ))
1848 }
1849
1850 pub(in crate::db) fn insert_persisted_snapshot(
1852 &mut self,
1853 entity: EntityTag,
1854 snapshot: &PersistedSchemaSnapshot,
1855 ) -> Result<(), InternalError> {
1856 let prepared = Self::prepare_persisted_snapshot(entity, snapshot)?;
1857 self.apply_prepared_persisted_snapshot(prepared);
1858
1859 Ok(())
1860 }
1861
1862 pub(in crate::db) fn prepare_persisted_snapshot(
1864 entity: EntityTag,
1865 snapshot: &PersistedSchemaSnapshot,
1866 ) -> Result<PreparedSchemaSnapshot, InternalError> {
1867 Ok(PreparedSchemaSnapshot {
1868 key: RawSchemaKey::from_entity_version(entity, snapshot.version()),
1869 snapshot: RawSchemaSnapshot::from_persisted_snapshot(snapshot)?,
1870 })
1871 }
1872
1873 pub(in crate::db) fn apply_prepared_persisted_snapshot(
1875 &mut self,
1876 prepared: PreparedSchemaSnapshot,
1877 ) {
1878 let _ = self.insert_raw_snapshot(prepared.key, prepared.snapshot);
1879 }
1880
1881 pub(in crate::db) fn constraint_validation_job(
1883 &self,
1884 entity: EntityTag,
1885 constraint_id: ConstraintId,
1886 ) -> Result<Option<ConstraintValidationJob>, InternalError> {
1887 let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1888 self.get_raw_snapshot(&key)
1889 .map(|raw| decode_constraint_validation_job(raw.as_bytes()))
1890 .transpose()
1891 }
1892
1893 pub(in crate::db) fn apply_constraint_validation_job(
1895 &mut self,
1896 job: &ConstraintValidationJob,
1897 ) -> Result<(), InternalError> {
1898 let key =
1899 RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1900 let bytes = encode_constraint_validation_job(job)?;
1901 let _ =
1902 self.insert_raw_snapshot(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1903 Ok(())
1904 }
1905
1906 #[expect(
1908 clippy::unnecessary_wraps,
1909 reason = "marker apply operations share one fallible callback contract"
1910 )]
1911 pub(in crate::db) fn apply_constraint_validation_job_removal(
1912 &mut self,
1913 entity: EntityTag,
1914 constraint_id: ConstraintId,
1915 ) -> Result<(), InternalError> {
1916 let key = RawSchemaKey::from_constraint_validation_job(entity, constraint_id);
1917 match &mut self.backend {
1918 SchemaStoreBackend::Heap(map) => {
1919 map.remove(&key);
1920 }
1921 SchemaStoreBackend::Journaled {
1922 live, tombstones, ..
1923 } => {
1924 live.remove(&key);
1925 tombstones.insert(key);
1926 }
1927 }
1928 Ok(())
1929 }
1930
1931 pub(in crate::db) fn fold_constraint_validation_job(
1933 &mut self,
1934 job: &ConstraintValidationJob,
1935 ) -> Result<(), InternalError> {
1936 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1937 return Err(InternalError::store_invariant());
1938 };
1939 let key =
1940 RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id());
1941 let bytes = encode_constraint_validation_job(job)?;
1942 canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
1943 Ok(())
1944 }
1945
1946 pub(in crate::db) fn preflight_fold_constraint_validation_job(
1948 &self,
1949 job: &ConstraintValidationJob,
1950 ) -> Result<(), InternalError> {
1951 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
1952 return Err(InternalError::store_invariant());
1953 }
1954 let _encoded = encode_constraint_validation_job(job)?;
1955 Ok(())
1956 }
1957
1958 pub(in crate::db) fn fold_constraint_validation_job_removal(
1960 &mut self,
1961 entity: EntityTag,
1962 constraint_id: ConstraintId,
1963 ) -> Result<(), InternalError> {
1964 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
1965 return Err(InternalError::store_invariant());
1966 };
1967 canonical.remove(&RawSchemaKey::from_constraint_validation_job(
1968 entity,
1969 constraint_id,
1970 ));
1971 Ok(())
1972 }
1973
1974 pub(in crate::db) fn preflight_fold_constraint_validation_job_removal(
1976 &self,
1977 ) -> Result<(), InternalError> {
1978 match self.backend {
1979 SchemaStoreBackend::Journaled { .. } => Ok(()),
1980 SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
1981 }
1982 }
1983
1984 pub(in crate::db) fn reset_journaled_live_projection(&mut self) -> Result<(), InternalError> {
1987 let SchemaStoreBackend::Journaled {
1988 live,
1989 tombstones,
1990 positions,
1991 ..
1992 } = &mut self.backend
1993 else {
1994 return Err(InternalError::store_invariant());
1995 };
1996
1997 live.clear();
1998 tombstones.clear();
1999 positions.clear();
2000 self.accepted_bundle_cache.get_mut().take();
2001
2002 Ok(())
2003 }
2004
2005 pub(in crate::db) fn prepare_positioned_journal_batch_publication(
2007 &self,
2008 incarnation: DatabaseIncarnationId,
2009 batch: &JournalBatch,
2010 position: JournalOverlayPosition,
2011 ) -> Result<PreparedSchemaPositionPublication, InternalError> {
2012 let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2013 return Err(InternalError::store_invariant());
2014 };
2015 let keys = self.positioned_journal_batch_keys(
2016 incarnation,
2017 batch,
2018 IdentityStateStorageView::Effective,
2019 )?;
2020 for key in &keys {
2021 positions.preflight_publish(key, position)?;
2022 }
2023 Ok(PreparedSchemaPositionPublication {
2024 keys: keys.into_iter().collect(),
2025 position,
2026 })
2027 }
2028
2029 pub(in crate::db) fn prepare_positioned_journal_batch_retirement(
2031 &self,
2032 incarnation: DatabaseIncarnationId,
2033 batch: &JournalBatch,
2034 position: JournalOverlayPosition,
2035 ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2036 let keys = self.positioned_journal_batch_keys(
2037 incarnation,
2038 batch,
2039 IdentityStateStorageView::Canonical,
2040 )?;
2041 self.prepare_positioned_key_retirements(keys, position)
2042 }
2043
2044 fn prepare_positioned_key_retirements(
2045 &self,
2046 keys: impl IntoIterator<Item = RawSchemaKey>,
2047 position: JournalOverlayPosition,
2048 ) -> Result<PreparedSchemaPositionRetirement, InternalError> {
2049 let SchemaStoreBackend::Journaled {
2050 live,
2051 tombstones,
2052 positions,
2053 ..
2054 } = &self.backend
2055 else {
2056 return Err(InternalError::store_invariant());
2057 };
2058 let mut entries = Vec::new();
2059 for key in keys {
2060 if !positions.is_positioned(&key) {
2061 if live.contains_key(&key) || tombstones.contains(&key) {
2065 return Err(InternalError::store_invariant());
2066 }
2067 continue;
2068 }
2069 let retirement = positions.preflight_retirement(&key, position)?;
2070 entries.push((key, retirement));
2071 }
2072 Ok(PreparedSchemaPositionRetirement { entries })
2073 }
2074
2075 pub(in crate::db) fn publish_prepared_journal_batch_positions(
2077 &mut self,
2078 prepared: PreparedSchemaPositionPublication,
2079 ) {
2080 let SchemaStoreBackend::Journaled { positions, .. } = &mut self.backend else {
2081 debug_assert!(
2082 false,
2083 "preflighted schema positions require a journaled store"
2084 );
2085 return;
2086 };
2087 for key in prepared.keys {
2088 positions.publish_preflighted(key, prepared.position);
2089 }
2090 }
2091
2092 pub(in crate::db) fn apply_prepared_journal_batch_retirement(
2094 &mut self,
2095 prepared: PreparedSchemaPositionRetirement,
2096 ) {
2097 for (key, retirement) in prepared.entries {
2098 if retirement != PositionedOverlayRetirement::Exact {
2099 continue;
2100 }
2101 self.invalidate_accepted_bundle_cache_for_key(key);
2102 let SchemaStoreBackend::Journaled {
2103 live,
2104 tombstones,
2105 positions,
2106 ..
2107 } = &mut self.backend
2108 else {
2109 debug_assert!(
2110 false,
2111 "preflighted schema retirement requires a journaled store"
2112 );
2113 return;
2114 };
2115 live.remove(&key);
2116 tombstones.remove(&key);
2117 positions.retire_preflighted(&key, retirement);
2118 }
2119 }
2120
2121 #[cfg(test)]
2122 fn publish_positioned_journal_entry(
2123 &mut self,
2124 key: RawSchemaKey,
2125 snapshot: Option<RawSchemaSnapshot>,
2126 position: JournalOverlayPosition,
2127 ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
2128 let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2129 return Err(InternalError::store_invariant());
2130 };
2131 positions.preflight_publish(&key, position)?;
2132 self.invalidate_accepted_bundle_cache_for_key(key);
2133 let SchemaStoreBackend::Journaled {
2134 canonical,
2135 live,
2136 tombstones,
2137 positions,
2138 } = &mut self.backend
2139 else {
2140 return Err(InternalError::store_invariant());
2141 };
2142 let previous = if tombstones.contains(&key) {
2143 None
2144 } else {
2145 live.get(&key).cloned().or_else(|| canonical.get(&key))
2146 };
2147 if let Some(snapshot) = snapshot {
2148 tombstones.remove(&key);
2149 live.insert(key, snapshot);
2150 } else {
2151 live.remove(&key);
2152 tombstones.insert(key);
2153 }
2154 positions.publish_preflighted(key, position);
2155 Ok(previous)
2156 }
2157
2158 #[cfg(test)]
2159 fn retire_positioned_journal_effect(
2160 &mut self,
2161 key: RawSchemaKey,
2162 position: JournalOverlayPosition,
2163 ) -> Result<PositionedOverlayRetirement, InternalError> {
2164 let SchemaStoreBackend::Journaled { positions, .. } = &self.backend else {
2165 return Err(InternalError::store_invariant());
2166 };
2167 let retirement = positions.preflight_retirement(&key, position)?;
2168 let prepared = PreparedSchemaPositionRetirement {
2169 entries: vec![(key, retirement)],
2170 };
2171 self.apply_prepared_journal_batch_retirement(prepared);
2172 Ok(retirement)
2173 }
2174
2175 #[cfg(test)]
2177 pub(in crate::db) fn fold_persisted_snapshot(
2178 &mut self,
2179 entity: EntityTag,
2180 snapshot: &PersistedSchemaSnapshot,
2181 ) -> Result<(), InternalError> {
2182 let prepared = self.prepare_fold_persisted_snapshot(entity, snapshot)?;
2183 self.apply_prepared_fold_persisted_snapshot(prepared)
2184 }
2185
2186 pub(in crate::db) fn apply_prepared_fold_persisted_snapshot(
2188 &mut self,
2189 prepared: PreparedSchemaSnapshot,
2190 ) -> Result<(), InternalError> {
2191 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
2192 return Err(InternalError::store_invariant());
2193 };
2194 canonical.insert(prepared.key, prepared.snapshot);
2195
2196 Ok(())
2197 }
2198
2199 pub(in crate::db) fn prepare_fold_persisted_snapshot(
2201 &self,
2202 entity: EntityTag,
2203 snapshot: &PersistedSchemaSnapshot,
2204 ) -> Result<PreparedSchemaSnapshot, InternalError> {
2205 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2206 return Err(InternalError::store_invariant());
2207 }
2208 Self::prepare_persisted_snapshot(entity, snapshot)
2209 }
2210
2211 pub(in crate::db) fn current_accepted_schema_root(
2213 &self,
2214 ) -> Result<Option<AcceptedSchemaRootSelection>, InternalError> {
2215 let first = self.accepted_root_slot_bytes(0)?;
2216 let second = self.accepted_root_slot_bytes(1)?;
2217 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])
2218 }
2219
2220 pub(in crate::db) fn current_accepted_schema_bundle(
2222 &self,
2223 ) -> Result<Option<AcceptedSchemaRevisionBundle>, InternalError> {
2224 self.borrow_current_accepted_schema_bundle()
2225 .map(|bundle| bundle.map(|bundle| bundle.clone()))
2226 }
2227
2228 pub(in crate::db) fn borrow_current_accepted_schema_bundle(
2231 &self,
2232 ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
2233 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2234 return Ok(None);
2235 };
2236 self.validate_constraint_validation_job_closure(&bundle)?;
2237
2238 Ok(Some(bundle))
2239 }
2240
2241 pub(in crate::db) fn current_accepted_runtime_entities(
2243 &self,
2244 registered_store_path: &'static str,
2245 ) -> Result<Vec<AcceptedRuntimeEntity>, InternalError> {
2246 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2247 return Ok(Vec::new());
2248 };
2249 if bundle.store_path() != registered_store_path {
2250 return Err(InternalError::store_corruption());
2251 }
2252
2253 bundle
2254 .entity_snapshots()
2255 .iter()
2256 .map(|(entity_tag, snapshot)| {
2257 AcceptedRuntimeEntity::from_accepted_snapshot(
2258 &bundle,
2259 *entity_tag,
2260 snapshot,
2261 registered_store_path,
2262 )
2263 })
2264 .collect()
2265 }
2266
2267 pub(in crate::db) fn current_accepted_runtime_entity_for_tag(
2269 &self,
2270 registered_store_path: &'static str,
2271 entity_tag: EntityTag,
2272 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2273 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2274 return Ok(None);
2275 };
2276 if bundle.store_path() != registered_store_path {
2277 return Err(InternalError::store_corruption());
2278 }
2279 let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2280 return Ok(None);
2281 };
2282
2283 AcceptedRuntimeEntity::from_accepted_snapshot(
2284 &bundle,
2285 entity_tag,
2286 snapshot,
2287 registered_store_path,
2288 )
2289 .map(Some)
2290 }
2291
2292 pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_tag(
2294 &self,
2295 registered_store_path: &'static str,
2296 entity_tag: EntityTag,
2297 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2298 let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2299 return Ok(None);
2300 };
2301 if bundle.store_path() != registered_store_path {
2302 return Err(InternalError::store_corruption());
2303 }
2304 let Some(snapshot) = bundle.entity_snapshots().get(&entity_tag) else {
2305 return Ok(None);
2306 };
2307
2308 AcceptedRuntimeEntity::from_accepted_snapshot(
2309 &bundle,
2310 entity_tag,
2311 snapshot,
2312 registered_store_path,
2313 )
2314 .map(Some)
2315 }
2316
2317 pub(in crate::db) fn current_accepted_runtime_entity_for_path(
2319 &self,
2320 registered_store_path: &'static str,
2321 entity_path: &str,
2322 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2323 self.current_accepted_runtime_entity_matching(registered_store_path, |snapshot_path, _| {
2324 snapshot_path == entity_path
2325 })
2326 }
2327
2328 pub(in crate::db) fn current_canonical_accepted_runtime_entity_for_path(
2330 &self,
2331 registered_store_path: &'static str,
2332 entity_path: &str,
2333 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2334 let Some(bundle) = self.current_canonical_accepted_schema_bundle()? else {
2335 return Ok(None);
2336 };
2337 if bundle.store_path() != registered_store_path {
2338 return Err(InternalError::store_corruption());
2339 }
2340
2341 let mut matched = None;
2342 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2343 if snapshot.entity_path() != entity_path {
2344 continue;
2345 }
2346 let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2347 &bundle,
2348 *entity_tag,
2349 snapshot,
2350 registered_store_path,
2351 )?;
2352 if matched.replace(entity).is_some() {
2353 return Err(InternalError::store_corruption());
2354 }
2355 }
2356
2357 Ok(matched)
2358 }
2359
2360 #[cfg(test)]
2362 pub(in crate::db) fn current_accepted_runtime_entity_for_name(
2363 &self,
2364 registered_store_path: &'static str,
2365 entity_name: &str,
2366 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2367 self.current_accepted_runtime_entity_matching(registered_store_path, |_, snapshot_name| {
2368 snapshot_name == entity_name
2369 })
2370 }
2371
2372 fn current_accepted_runtime_entity_matching(
2373 &self,
2374 registered_store_path: &'static str,
2375 mut predicate: impl FnMut(&str, &str) -> bool,
2376 ) -> Result<Option<AcceptedRuntimeEntity>, InternalError> {
2377 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2378 return Ok(None);
2379 };
2380 if bundle.store_path() != registered_store_path {
2381 return Err(InternalError::store_corruption());
2382 }
2383
2384 let mut matched = None;
2385 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2386 if !predicate(snapshot.entity_path(), snapshot.entity_name()) {
2387 continue;
2388 }
2389 let entity = AcceptedRuntimeEntity::from_accepted_snapshot(
2390 &bundle,
2391 *entity_tag,
2392 snapshot,
2393 registered_store_path,
2394 )?;
2395 if matched.replace(entity).is_some() {
2396 return Err(InternalError::store_corruption());
2397 }
2398 }
2399
2400 Ok(matched)
2401 }
2402
2403 pub(in crate::db) fn current_accepted_schema_revision(
2405 &self,
2406 ) -> Result<Option<AcceptedSchemaRevision>, InternalError> {
2407 Ok(self
2408 .current_accepted_schema_root()?
2409 .map(|selection| selection.root().revision()))
2410 }
2411
2412 pub(in crate::db) fn pending_relation_activation_for_target(
2418 &self,
2419 target_path: &str,
2420 ) -> Result<Option<PendingRelationActivationDeleteBarrier>, InternalError> {
2421 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2422 return Ok(None);
2423 };
2424 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2425 let Some(candidate) = snapshot
2426 .candidate_relations()
2427 .iter()
2428 .find(|candidate| candidate.target_path() == target_path)
2429 else {
2430 continue;
2431 };
2432 let activation = snapshot
2433 .constraint_activations()
2434 .iter()
2435 .find(|activation| {
2436 matches!(
2437 activation.kind(),
2438 ConstraintActivationKind::Relation { relation_id }
2439 if *relation_id == candidate.id()
2440 )
2441 })
2442 .ok_or_else(InternalError::store_corruption)?;
2443 return Ok(Some(PendingRelationActivationDeleteBarrier {
2444 accepted_schema_fingerprint:
2445 accepted_schema_cache_fingerprint_for_persisted_snapshot(snapshot)?,
2446 source_entity_tag: *entity_tag,
2447 constraint_id: activation.id(),
2448 }));
2449 }
2450
2451 Ok(None)
2452 }
2453
2454 pub(in crate::db) fn entity_has_relation_to_target(
2456 &self,
2457 source_entity: EntityTag,
2458 target_path: &str,
2459 ) -> Result<bool, InternalError> {
2460 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
2461 return Ok(false);
2462 };
2463 let Some(snapshot) = bundle.entity_snapshots().get(&source_entity) else {
2464 return Ok(false);
2465 };
2466
2467 Ok(snapshot
2468 .relations()
2469 .iter()
2470 .any(|relation| relation.target_path() == target_path))
2471 }
2472
2473 pub(in crate::db) fn validate_live_activation_transition(
2475 &self,
2476 candidate: &AcceptedSchemaRevisionBundle,
2477 ) -> Result<(), InternalError> {
2478 let Some(current) = self.current_accepted_schema_bundle()? else {
2479 return Ok(());
2480 };
2481 Self::validate_activation_transition_from(¤t, candidate)
2482 }
2483
2484 pub(in crate::db) fn validate_canonical_activation_transition(
2486 &self,
2487 candidate: &AcceptedSchemaRevisionBundle,
2488 ) -> Result<(), InternalError> {
2489 let Some(current) = self.current_canonical_accepted_schema_bundle()? else {
2490 return Ok(());
2491 };
2492 Self::validate_activation_transition_from(¤t, candidate)
2493 }
2494
2495 fn validate_activation_transition_from(
2496 current: &AcceptedSchemaRevisionBundle,
2497 candidate: &AcceptedSchemaRevisionBundle,
2498 ) -> Result<(), InternalError> {
2499 for (entity_tag, before) in current.entity_snapshots() {
2500 if before.constraint_activations().is_empty() {
2501 continue;
2502 }
2503 let after = candidate
2504 .entity_snapshots()
2505 .get(entity_tag)
2506 .ok_or_else(InternalError::store_invariant)?;
2507 if before == after {
2508 continue;
2509 }
2510 let expected_shape = before
2511 .clone()
2512 .with_constraint_catalog(after.constraint_catalog().clone());
2513 let catalog_only_transition = expected_shape == *after
2514 && before
2515 .constraint_catalog()
2516 .permits_live_activation_transition_to(after.constraint_catalog());
2517 let sql_row_local_abort_with_version =
2518 before.constraint_activations().iter().any(|activation| {
2519 activation.origin() == ConstraintOrigin::SqlDdl
2520 && matches!(
2521 activation.kind(),
2522 ConstraintActivationKind::Check { .. }
2523 | ConstraintActivationKind::NotNull { .. }
2524 )
2525 && before.version().get().checked_add(1) == Some(after.version().get())
2526 && before
2527 .constraint_catalog()
2528 .clone()
2529 .with_aborted_activation(activation.id())
2530 .is_ok_and(|catalog| catalog == *after.constraint_catalog())
2531 && before
2532 .clone()
2533 .with_constraint_catalog(after.constraint_catalog().clone())
2534 .with_schema_version(after.version())
2535 == *after
2536 });
2537 let sql_unique_abort_with_version =
2538 before.constraint_activations().iter().any(|activation| {
2539 activation.origin() == ConstraintOrigin::SqlDdl
2540 && matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2541 && before.version().get().checked_add(1) == Some(after.version().get())
2542 && before
2543 .with_aborted_unique_activation(activation.id(), after.version())
2544 .is_ok_and(|expected| expected == *after)
2545 });
2546 let not_null_promotion = before.constraint_activations().iter().any(|activation| {
2547 matches!(activation.kind(), ConstraintActivationKind::NotNull { .. })
2548 && before
2549 .with_promoted_not_null_activation(activation.id(), after.version())
2550 .is_ok_and(|expected| expected == *after)
2551 });
2552 let unique_promotion = before.constraint_activations().iter().any(|activation| {
2553 matches!(activation.kind(), ConstraintActivationKind::Unique { .. })
2554 && before
2555 .with_promoted_unique_activation(activation.id(), after.version())
2556 .is_ok_and(|expected| expected == *after)
2557 });
2558 let relation_promotion = before.constraint_activations().iter().any(|activation| {
2559 matches!(activation.kind(), ConstraintActivationKind::Relation { .. })
2560 && before
2561 .with_promoted_relation_activation(activation.id(), after.version())
2562 .is_ok_and(|expected| expected == *after)
2563 });
2564 if !catalog_only_transition
2565 && !sql_row_local_abort_with_version
2566 && !sql_unique_abort_with_version
2567 && !not_null_promotion
2568 && !unique_promotion
2569 && !relation_promotion
2570 {
2571 return Err(InternalError::store_invariant());
2572 }
2573 }
2574 Ok(())
2575 }
2576
2577 pub(in crate::db) fn validate_constraint_validation_job_closure(
2579 &self,
2580 bundle: &AcceptedSchemaRevisionBundle,
2581 ) -> Result<(), InternalError> {
2582 self.validate_constraint_validation_job_closure_with_change(bundle, None, None)
2583 }
2584
2585 pub(in crate::db) fn validate_constraint_validation_job_closure_with_change(
2588 &self,
2589 bundle: &AcceptedSchemaRevisionBundle,
2590 replacement: Option<&ConstraintValidationJob>,
2591 removal: Option<(EntityTag, ConstraintId)>,
2592 ) -> Result<(), InternalError> {
2593 self.validate_constraint_validation_job_closure_with_change_in_view(
2594 bundle,
2595 replacement,
2596 removal,
2597 IdentityStateStorageView::Effective,
2598 )
2599 }
2600
2601 pub(in crate::db) fn validate_canonical_constraint_validation_job_closure_with_change(
2603 &self,
2604 bundle: &AcceptedSchemaRevisionBundle,
2605 replacement: Option<&ConstraintValidationJob>,
2606 removal: Option<(EntityTag, ConstraintId)>,
2607 ) -> Result<(), InternalError> {
2608 self.validate_constraint_validation_job_closure_with_change_in_view(
2609 bundle,
2610 replacement,
2611 removal,
2612 IdentityStateStorageView::Canonical,
2613 )
2614 }
2615
2616 fn validate_constraint_validation_job_closure_with_change_in_view(
2617 &self,
2618 bundle: &AcceptedSchemaRevisionBundle,
2619 replacement: Option<&ConstraintValidationJob>,
2620 removal: Option<(EntityTag, ConstraintId)>,
2621 view: IdentityStateStorageView,
2622 ) -> Result<(), InternalError> {
2623 if replacement.is_some() && removal.is_some() {
2624 return Err(InternalError::store_invariant());
2625 }
2626 let replacement_key = replacement.map(|job| {
2627 RawSchemaKey::from_constraint_validation_job(job.entity_tag(), job.constraint_id())
2628 });
2629 let removal_key = removal.map(|(entity_tag, constraint_id)| {
2630 RawSchemaKey::from_constraint_validation_job(entity_tag, constraint_id)
2631 });
2632 let mut expected = BTreeSet::new();
2633 for (entity_tag, snapshot) in bundle.entity_snapshots() {
2634 for activation in snapshot.constraint_activations() {
2635 let key =
2636 RawSchemaKey::from_constraint_validation_job(*entity_tag, activation.id());
2637 match activation.state() {
2638 ConstraintActivationState::EnforcingNewWrites => {
2639 if self
2640 .constraint_validation_job_after_change(
2641 key,
2642 replacement,
2643 replacement_key,
2644 removal_key,
2645 view,
2646 )?
2647 .is_some()
2648 {
2649 return Err(InternalError::store_corruption());
2650 }
2651 }
2652 ConstraintActivationState::Validating => {
2653 let job = self
2654 .constraint_validation_job_after_change(
2655 key,
2656 replacement,
2657 replacement_key,
2658 removal_key,
2659 view,
2660 )?
2661 .ok_or_else(InternalError::store_corruption)?;
2662 if job.entity_tag() != *entity_tag
2663 || job.entity_path() != snapshot.entity_path()
2664 {
2665 return Err(InternalError::store_corruption());
2666 }
2667 job.validate(Some(activation))?;
2668 expected.insert(key);
2669 }
2670 }
2671 }
2672 }
2673
2674 self.visit_constraint_validation_jobs_in_view(view, |key, raw| {
2675 if removal_key == Some(*key) || replacement_key == Some(*key) {
2676 return Ok(SchemaStoreVisit::Continue);
2677 }
2678 if !expected.contains(key) {
2679 return Err(InternalError::store_corruption());
2680 }
2681 let job = decode_constraint_validation_job(raw.as_bytes())?;
2682 if job.entity_tag() != key.entity_tag()
2683 || key.constraint_id() != Some(job.constraint_id())
2684 {
2685 return Err(InternalError::store_corruption());
2686 }
2687 Ok(SchemaStoreVisit::Continue)
2688 })?;
2689
2690 if let Some(key) = replacement_key
2691 && !expected.contains(&key)
2692 {
2693 return Err(InternalError::store_corruption());
2694 }
2695 if let Some(key) = removal_key
2696 && expected.contains(&key)
2697 {
2698 return Err(InternalError::store_corruption());
2699 }
2700
2701 Ok(())
2702 }
2703
2704 fn constraint_validation_job_after_change(
2705 &self,
2706 key: RawSchemaKey,
2707 replacement: Option<&ConstraintValidationJob>,
2708 replacement_key: Option<RawSchemaKey>,
2709 removal_key: Option<RawSchemaKey>,
2710 view: IdentityStateStorageView,
2711 ) -> Result<Option<ConstraintValidationJob>, InternalError> {
2712 if removal_key == Some(key) {
2713 return Ok(None);
2714 }
2715 if replacement_key == Some(key) {
2716 return Ok(replacement.cloned());
2717 }
2718 let raw = match view {
2719 IdentityStateStorageView::Effective => self.get_raw_snapshot(&key),
2720 IdentityStateStorageView::Canonical => self.get_canonical_raw_value(&key)?,
2721 };
2722 raw.map(|raw| decode_constraint_validation_job(raw.as_bytes()))
2723 .transpose()
2724 }
2725
2726 pub(in crate::db) fn current_accepted_schema_authority_matches(
2729 &self,
2730 expected: &AcceptedSchemaAuthority,
2731 ) -> Result<bool, InternalError> {
2732 let Some(store_scope) = self.accepted_catalog_scope.get() else {
2733 return Ok(false);
2734 };
2735
2736 if let Some(cached) = self
2739 .accepted_bundle_cache
2740 .try_borrow()
2741 .map_err(|_| InternalError::store_invariant())?
2742 .as_ref()
2743 {
2744 let root = cached.selection.root();
2745 return Ok(expected.matches_store_root(
2746 store_scope,
2747 root.revision(),
2748 root.fingerprint(),
2749 ));
2750 }
2751
2752 let Some(selection) = self.current_accepted_schema_root()? else {
2753 return Ok(false);
2754 };
2755 let root = selection.root();
2756
2757 Ok(expected.matches_store_root(store_scope, root.revision(), root.fingerprint()))
2758 }
2759
2760 pub(in crate::db) fn publish_accepted_schema_candidate(
2766 &mut self,
2767 incarnation: DatabaseIncarnationId,
2768 expected_revision: AcceptedSchemaRevision,
2769 candidate: &CandidateSchemaRevision,
2770 ) -> Result<(), InternalError> {
2771 let identity_transition = self.prepare_identity_state_transition(
2772 incarnation,
2773 candidate,
2774 IdentityStateStorageView::Effective,
2775 )?;
2776 if self.current_root_matches_candidate(candidate)? {
2777 if !identity_transition.is_empty() {
2778 return Err(InternalError::identity_state_corruption());
2779 }
2780 let selection = self
2781 .current_accepted_schema_root()?
2782 .ok_or_else(InternalError::store_corruption)?;
2783 self.retain_durable_candidate_entries(candidate, selection.slot())?;
2784 return Ok(());
2785 }
2786 let first = self.accepted_root_slot_bytes(0)?;
2787 let second = self.accepted_root_slot_bytes(1)?;
2788 prepare_accepted_schema_root_publication(
2789 [first.as_deref(), second.as_deref()],
2790 expected_revision,
2791 candidate,
2792 )
2793 .map_err(map_schema_publication_error)?;
2794
2795 self.insert_durable_candidate_snapshots(candidate)?;
2796 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2797 self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2798 let persisted_bundle = self
2799 .get_raw_snapshot(&bundle_key)
2800 .ok_or_else(InternalError::store_corruption)?;
2801 let _verified = decode_verified_accepted_schema_revision_bundle(
2802 candidate.root(),
2803 persisted_bundle.as_bytes(),
2804 )?;
2805 self.apply_identity_state_transition(
2806 identity_transition,
2807 IdentityStateWriteTarget::Durable,
2808 )?;
2809
2810 let first = self.accepted_root_slot_bytes(0)?;
2813 let second = self.accepted_root_slot_bytes(1)?;
2814 let publication = prepare_accepted_schema_root_publication(
2815 [first.as_deref(), second.as_deref()],
2816 expected_revision,
2817 candidate,
2818 )
2819 .map_err(map_schema_publication_error)?;
2820 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
2821 self.insert_durable_raw_value(root_key, publication.encoded_root().to_vec());
2822
2823 let selected = self
2824 .current_accepted_schema_root()?
2825 .ok_or_else(InternalError::store_corruption)?;
2826 if selected.root() != candidate.root() {
2827 return Err(InternalError::store_corruption());
2828 }
2829 self.retain_durable_candidate_entries(candidate, selected.slot())?;
2830 Ok(())
2831 }
2832
2833 pub(in crate::db) fn restore_live_accepted_schema_checkpoint(
2836 &mut self,
2837 incarnation: DatabaseIncarnationId,
2838 candidate: &CandidateSchemaRevision,
2839 checkpoint_identity_states: &IdentityStateInventory,
2840 ) -> Result<(), InternalError> {
2841 if !matches!(self.backend, SchemaStoreBackend::Heap(_)) {
2842 return Err(InternalError::store_invariant());
2843 }
2844 let checkpoint_validation = prepare_identity_state_transition(
2845 incarnation,
2846 Some(candidate.bundle()),
2847 candidate.bundle(),
2848 checkpoint_identity_states.clone(),
2849 )?;
2850 if !checkpoint_validation.is_empty() {
2851 return Err(InternalError::identity_state_corruption());
2852 }
2853 if self.current_root_matches_candidate(candidate)? {
2854 for state in checkpoint_identity_states.values() {
2855 let key = RawSchemaKey::from_identity_state(
2856 state.owner().entity_tag(),
2857 state.owner().field_id(),
2858 );
2859 self.insert_durable_raw_value(key, encode_identity_state(state)?);
2860 }
2861 if self.identity_state_inventory(IdentityStateStorageView::Effective)?
2862 != *checkpoint_identity_states
2863 {
2864 return Err(InternalError::identity_state_corruption());
2865 }
2866 let selection = self
2867 .current_accepted_schema_root()?
2868 .ok_or_else(InternalError::store_corruption)?;
2869 self.retain_durable_candidate_entries(candidate, selection.slot())?;
2870 return Ok(());
2871 }
2872 if self.current_accepted_schema_root()?.is_some()
2873 || !self
2874 .identity_state_inventory(IdentityStateStorageView::Effective)?
2875 .is_empty()
2876 {
2877 return Err(InternalError::store_corruption());
2878 }
2879
2880 self.insert_durable_candidate_snapshots(candidate)?;
2881 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
2882 self.insert_durable_raw_value(bundle_key, candidate.encoded_bundle().to_vec());
2883 for state in checkpoint_identity_states.values() {
2884 let key = RawSchemaKey::from_identity_state(
2885 state.owner().entity_tag(),
2886 state.owner().field_id(),
2887 );
2888 self.insert_durable_raw_value(key, encode_identity_state(state)?);
2889 }
2890 let root_key = RawSchemaKey::from_accepted_root_slot(0)?;
2891 self.insert_durable_raw_value(root_key, candidate.encoded_root().to_vec());
2892
2893 let selected = self
2894 .current_accepted_schema_root()?
2895 .ok_or_else(InternalError::store_corruption)?;
2896 if selected.root() != candidate.root() {
2897 return Err(InternalError::store_corruption());
2898 }
2899 self.retain_durable_candidate_entries(candidate, selected.slot())?;
2900 Ok(())
2901 }
2902
2903 pub(in crate::db) fn preflight_accepted_schema_candidate(
2910 &self,
2911 incarnation: DatabaseIncarnationId,
2912 expected_revision: AcceptedSchemaRevision,
2913 candidate: &CandidateSchemaRevision,
2914 ) -> Result<bool, InternalError> {
2915 let identity_transition = self.prepare_identity_state_transition(
2916 incarnation,
2917 candidate,
2918 IdentityStateStorageView::Effective,
2919 )?;
2920 if self.current_root_matches_candidate(candidate)? {
2921 if !identity_transition.is_empty() {
2922 return Err(InternalError::identity_state_corruption());
2923 }
2924 return Ok(true);
2925 }
2926 let first = self.accepted_root_slot_bytes(0)?;
2927 let second = self.accepted_root_slot_bytes(1)?;
2928 prepare_accepted_schema_root_publication(
2929 [first.as_deref(), second.as_deref()],
2930 expected_revision,
2931 candidate,
2932 )
2933 .map_err(map_schema_publication_error)?;
2934
2935 Ok(false)
2936 }
2937
2938 pub(in crate::db) fn prepare_fold_journaled_accepted_schema_candidate(
2940 &self,
2941 incarnation: DatabaseIncarnationId,
2942 expected_revision: AcceptedSchemaRevision,
2943 candidate: CandidateSchemaRevision,
2944 ) -> Result<PreparedAcceptedSchemaFold, InternalError> {
2945 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
2946 return Err(InternalError::store_invariant());
2947 }
2948 let identity_transition = self.prepare_identity_state_transition(
2949 incarnation,
2950 &candidate,
2951 IdentityStateStorageView::Canonical,
2952 )?;
2953 let candidate_is_current = self.canonical_root_matches_candidate(&candidate)?;
2954 if candidate_is_current && !identity_transition.is_empty() {
2955 return Err(InternalError::identity_state_corruption());
2956 }
2957 let identity_updates = identity_transition
2958 .into_updates()
2959 .into_iter()
2960 .map(|state| {
2961 Ok((
2962 RawSchemaKey::from_identity_state(
2963 state.owner().entity_tag(),
2964 state.owner().field_id(),
2965 ),
2966 encode_identity_state(&state)?,
2967 ))
2968 })
2969 .collect::<Result<Vec<_>, InternalError>>()?;
2970 let snapshots = candidate
2971 .bundle()
2972 .entity_snapshots()
2973 .iter()
2974 .map(|(entity, snapshot)| Self::prepare_persisted_snapshot(*entity, snapshot))
2975 .collect::<Result<Vec<_>, _>>()?;
2976
2977 let first = self.canonical_root_slot_bytes(0)?;
2978 let second = self.canonical_root_slot_bytes(1)?;
2979 let root_slot = if candidate_is_current {
2980 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
2981 .ok_or_else(InternalError::store_corruption)?
2982 .slot()
2983 } else {
2984 prepare_accepted_schema_root_publication(
2985 [first.as_deref(), second.as_deref()],
2986 expected_revision,
2987 &candidate,
2988 )
2989 .map_err(map_schema_publication_error)?
2990 .target_slot()
2991 };
2992 let retained = Self::candidate_entry_keys(&candidate, root_slot)?;
2993 Ok(PreparedAcceptedSchemaFold {
2994 candidate,
2995 expected_revision,
2996 snapshots,
2997 identity_updates,
2998 retained,
2999 root_slot,
3000 })
3001 }
3002
3003 pub(in crate::db) fn projected_identity_state_count(
3005 &self,
3006 incarnation: DatabaseIncarnationId,
3007 candidate: &CandidateSchemaRevision,
3008 ) -> Result<usize, InternalError> {
3009 Ok(self
3010 .prepare_identity_state_transition(
3011 incarnation,
3012 candidate,
3013 IdentityStateStorageView::Effective,
3014 )?
3015 .projected_inventory_len())
3016 }
3017
3018 pub(in crate::db) fn apply_journaled_accepted_schema_candidate(
3020 &mut self,
3021 incarnation: DatabaseIncarnationId,
3022 expected_revision: AcceptedSchemaRevision,
3023 candidate: &CandidateSchemaRevision,
3024 ) -> Result<(), InternalError> {
3025 if !matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3026 return Err(InternalError::store_invariant());
3027 }
3028 let identity_transition = self.prepare_identity_state_transition(
3029 incarnation,
3030 candidate,
3031 IdentityStateStorageView::Effective,
3032 )?;
3033 if self.current_root_matches_candidate(candidate)? {
3034 if !identity_transition.is_empty() {
3035 return Err(InternalError::identity_state_corruption());
3036 }
3037 let selection = self
3038 .current_accepted_schema_root()?
3039 .ok_or_else(InternalError::store_corruption)?;
3040 self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3041 return Ok(());
3042 }
3043
3044 let first = self.accepted_root_slot_bytes(0)?;
3045 let second = self.accepted_root_slot_bytes(1)?;
3046 prepare_accepted_schema_root_publication(
3047 [first.as_deref(), second.as_deref()],
3048 expected_revision,
3049 candidate,
3050 )
3051 .map_err(map_schema_publication_error)?;
3052
3053 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3054 self.insert_persisted_snapshot(*entity_tag, snapshot)?;
3055 }
3056 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3057 self.insert_raw_snapshot(
3058 bundle_key,
3059 RawSchemaSnapshot::from_encoded_control_record(candidate.encoded_bundle().to_vec()),
3060 );
3061 let persisted_bundle = self
3062 .get_raw_snapshot(&bundle_key)
3063 .ok_or_else(InternalError::store_corruption)?;
3064 let _verified = decode_verified_accepted_schema_revision_bundle(
3065 candidate.root(),
3066 persisted_bundle.as_bytes(),
3067 )?;
3068 self.apply_identity_state_transition(
3069 identity_transition,
3070 IdentityStateWriteTarget::Materialized,
3071 )?;
3072
3073 let first = self.accepted_root_slot_bytes(0)?;
3074 let second = self.accepted_root_slot_bytes(1)?;
3075 let publication = prepare_accepted_schema_root_publication(
3076 [first.as_deref(), second.as_deref()],
3077 expected_revision,
3078 candidate,
3079 )
3080 .map_err(map_schema_publication_error)?;
3081 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3082 self.insert_raw_snapshot(
3083 root_key,
3084 RawSchemaSnapshot::from_encoded_control_record(publication.encoded_root().to_vec()),
3085 );
3086
3087 if !self.current_root_matches_candidate(candidate)? {
3088 return Err(InternalError::store_corruption());
3089 }
3090 let selection = self
3091 .current_accepted_schema_root()?
3092 .ok_or_else(InternalError::store_corruption)?;
3093 self.retain_materialized_candidate_entries(candidate, selection.slot())?;
3094 Ok(())
3095 }
3096
3097 pub(in crate::db) fn apply_prepared_accepted_schema_fold(
3099 &mut self,
3100 prepared: PreparedAcceptedSchemaFold,
3101 ) -> Result<(), InternalError> {
3102 let PreparedAcceptedSchemaFold {
3103 candidate,
3104 expected_revision,
3105 snapshots,
3106 identity_updates,
3107 retained,
3108 root_slot,
3109 } = prepared;
3110 if self.canonical_root_matches_candidate(&candidate)? {
3111 if !identity_updates.is_empty() {
3112 return Err(InternalError::identity_state_corruption());
3113 }
3114 let first = self.canonical_root_slot_bytes(0)?;
3115 let second = self.canonical_root_slot_bytes(1)?;
3116 let selection =
3117 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3118 .ok_or_else(InternalError::store_corruption)?;
3119 if selection.slot() != root_slot {
3120 return Err(InternalError::store_invariant());
3121 }
3122 self.retain_canonical_candidate_entries(&retained)?;
3123 return Ok(());
3124 }
3125
3126 let first = self.canonical_root_slot_bytes(0)?;
3127 let second = self.canonical_root_slot_bytes(1)?;
3128 let publication = prepare_accepted_schema_root_publication(
3129 [first.as_deref(), second.as_deref()],
3130 expected_revision,
3131 &candidate,
3132 )
3133 .map_err(map_schema_publication_error)?;
3134 if publication.target_slot() != root_slot {
3135 return Err(InternalError::store_invariant());
3136 }
3137
3138 for snapshot in snapshots {
3139 self.apply_prepared_fold_persisted_snapshot(snapshot)?;
3140 }
3141 let bundle_key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3142 self.insert_canonical_raw_value(bundle_key, candidate.encoded_bundle().to_vec())?;
3143 let persisted_bundle = self
3144 .get_canonical_raw_value(&bundle_key)?
3145 .ok_or_else(InternalError::store_corruption)?;
3146 let _verified = decode_verified_accepted_schema_revision_bundle(
3147 candidate.root(),
3148 persisted_bundle.as_bytes(),
3149 )?;
3150 for (key, bytes) in identity_updates {
3151 self.insert_canonical_raw_value(key, bytes)?;
3152 }
3153
3154 let first = self.canonical_root_slot_bytes(0)?;
3155 let second = self.canonical_root_slot_bytes(1)?;
3156 let publication = prepare_accepted_schema_root_publication(
3157 [first.as_deref(), second.as_deref()],
3158 expected_revision,
3159 &candidate,
3160 )
3161 .map_err(map_schema_publication_error)?;
3162 let root_key = RawSchemaKey::from_accepted_root_slot(publication.target_slot())?;
3163 self.insert_canonical_raw_value(root_key, publication.encoded_root().to_vec())?;
3164
3165 if !self.canonical_root_matches_candidate(&candidate)? {
3166 return Err(InternalError::store_corruption());
3167 }
3168 let first = self.canonical_root_slot_bytes(0)?;
3169 let second = self.canonical_root_slot_bytes(1)?;
3170 let selection = select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3171 .ok_or_else(InternalError::store_corruption)?;
3172 if selection.slot() != root_slot {
3173 return Err(InternalError::store_invariant());
3174 }
3175 self.retain_canonical_candidate_entries(&retained)?;
3176 Ok(())
3177 }
3178
3179 pub(in crate::db) fn get_persisted_snapshot(
3181 &self,
3182 entity: EntityTag,
3183 version: SchemaVersion,
3184 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3185 let key = RawSchemaKey::from_entity_version(entity, version);
3186 self.get_raw_snapshot(&key)
3187 .map(|snapshot| snapshot.decode_persisted_snapshot())
3188 .transpose()
3189 }
3190
3191 #[cfg(test)]
3192 fn latest_staged_persisted_snapshot(
3193 &self,
3194 entity: EntityTag,
3195 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3196 self.latest_raw_snapshots_by_entity()
3197 .remove(&entity)
3198 .map(|(_, snapshot)| snapshot.decode_persisted_snapshot())
3199 .transpose()
3200 }
3201
3202 pub(in crate::db) fn current_accepted_persisted_snapshot(
3205 &self,
3206 entity: EntityTag,
3207 ) -> Result<Option<PersistedSchemaSnapshot>, InternalError> {
3208 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3209 return Ok(None);
3210 };
3211
3212 Ok(bundle.entity_snapshots().get(&entity).cloned())
3213 }
3214
3215 pub(in crate::db) fn current_accepted_catalog_selection(
3217 &self,
3218 entity: EntityTag,
3219 entity_path: &str,
3220 store_path: &'static str,
3221 ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3222 let Some(bundle) = self.current_accepted_schema_bundle_ref()? else {
3223 return Ok(None);
3224 };
3225 if bundle.store_path() != store_path {
3226 return Err(InternalError::store_corruption());
3227 }
3228 let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3229 return Ok(None);
3230 };
3231 if snapshot.entity_path() != entity_path {
3232 return Err(InternalError::store_corruption());
3233 }
3234
3235 let cache = self
3236 .accepted_bundle_cache
3237 .try_borrow()
3238 .map_err(|_| InternalError::store_invariant())?;
3239 let cached = cache.as_ref().ok_or_else(InternalError::store_invariant)?;
3240 if let Some(selection) = cached
3241 .entity_selections
3242 .try_borrow()
3243 .map_err(|_| InternalError::store_invariant())?
3244 .get(&entity)
3245 .cloned()
3246 {
3247 return Ok(Some(selection));
3248 }
3249
3250 let selected = AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3251 entity,
3252 store_path,
3253 snapshot,
3254 cached.value_catalog.clone(),
3255 )?;
3256 cached
3257 .entity_selections
3258 .try_borrow_mut()
3259 .map_err(|_| InternalError::store_invariant())?
3260 .insert(entity, selected.clone());
3261
3262 Ok(Some(selected))
3263 }
3264
3265 pub(in crate::db) fn current_canonical_accepted_catalog_selection(
3269 &self,
3270 entity: EntityTag,
3271 entity_path: &str,
3272 store_path: &'static str,
3273 ) -> Result<Option<AcceptedCatalogSnapshotSelection>, InternalError> {
3274 let first = self.canonical_root_slot_bytes(0)?;
3275 let second = self.canonical_root_slot_bytes(1)?;
3276 let Some(selection) =
3277 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3278 else {
3279 return Ok(None);
3280 };
3281 let bundle_key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3282 let raw_bundle = self
3283 .get_canonical_raw_value(&bundle_key)?
3284 .ok_or_else(InternalError::store_corruption)?;
3285 let bundle = decode_verified_accepted_schema_revision_bundle(
3286 selection.root(),
3287 raw_bundle.as_bytes(),
3288 )?;
3289 if bundle.store_path() != store_path {
3290 return Err(InternalError::store_corruption());
3291 }
3292 let Some(snapshot) = bundle.entity_snapshots().get(&entity) else {
3293 return Ok(None);
3294 };
3295 if snapshot.entity_path() != entity_path {
3296 return Err(InternalError::store_corruption());
3297 }
3298
3299 AcceptedCatalogSnapshotSelection::from_verified_snapshot(
3300 entity,
3301 store_path,
3302 snapshot,
3303 AcceptedValueCatalogHandle::new(
3304 bundle.enum_catalog().clone(),
3305 bundle.composite_catalog().clone(),
3306 self.accepted_catalog_scope
3307 .get_or_init(AcceptedStoreCatalogScope::new)
3308 .clone(),
3309 bundle.revision(),
3310 selection.root().fingerprint(),
3311 ),
3312 )
3313 .map(Some)
3314 }
3315
3316 #[cfg(test)]
3322 pub(in crate::db) fn catalog_metadata(
3323 &self,
3324 ) -> Result<Option<SchemaStoreCatalogMetadata>, InternalError> {
3325 Ok(self
3326 .allocation_metadata()?
3327 .map(SchemaStoreAllocationMetadata::schema))
3328 }
3329
3330 pub(in crate::db) fn allocation_metadata(
3337 &self,
3338 ) -> Result<Option<SchemaStoreAllocationMetadata>, InternalError> {
3339 let latest_by_entity = self.latest_raw_snapshots_by_entity();
3340 if latest_by_entity.is_empty() {
3341 return Ok(None);
3342 }
3343
3344 Ok(Some(SchemaStoreAllocationMetadata::new(
3345 derive_data_allocation_metadata(&latest_by_entity)?,
3346 derive_index_allocation_metadata(&latest_by_entity)?,
3347 derive_schema_catalog_metadata(&latest_by_entity)?,
3348 )))
3349 }
3350
3351 fn insert_raw_snapshot(
3353 &mut self,
3354 key: RawSchemaKey,
3355 snapshot: RawSchemaSnapshot,
3356 ) -> Option<RawSchemaSnapshot> {
3357 self.invalidate_accepted_bundle_cache_for_key(key);
3358 let previous_journaled = if matches!(self.backend, SchemaStoreBackend::Journaled { .. }) {
3359 self.get_raw_snapshot_for_backend(&key)
3360 } else {
3361 None
3362 };
3363 match &mut self.backend {
3364 SchemaStoreBackend::Heap(map) => map.insert(key, snapshot),
3365 SchemaStoreBackend::Journaled {
3366 live, tombstones, ..
3367 } => {
3368 tombstones.remove(&key);
3369 live.insert(key, snapshot);
3370 previous_journaled
3371 }
3372 }
3373 }
3374
3375 #[must_use]
3377 fn get_raw_snapshot(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
3378 match &self.backend {
3379 SchemaStoreBackend::Heap(map) => map.get(key).cloned(),
3380 SchemaStoreBackend::Journaled { .. } => self.get_raw_snapshot_for_backend(key),
3381 }
3382 }
3383
3384 fn accepted_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3385 let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3386 Ok(self
3387 .get_raw_snapshot(&key)
3388 .map(RawSchemaSnapshot::into_bytes))
3389 }
3390
3391 fn canonical_root_slot_bytes(&self, slot: usize) -> Result<Option<Vec<u8>>, InternalError> {
3392 let key = RawSchemaKey::from_accepted_root_slot(slot)?;
3393 Ok(self
3394 .get_canonical_raw_value(&key)?
3395 .map(RawSchemaSnapshot::into_bytes))
3396 }
3397
3398 fn current_root_matches_candidate(
3399 &self,
3400 candidate: &CandidateSchemaRevision,
3401 ) -> Result<bool, InternalError> {
3402 let Some(selection) = self.current_accepted_schema_root()? else {
3403 return Ok(false);
3404 };
3405 if selection.root() != candidate.root() {
3406 return Ok(false);
3407 }
3408 let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3409 let bundle = self
3410 .get_raw_snapshot(&key)
3411 .ok_or_else(InternalError::store_corruption)?;
3412 let _verified =
3413 decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3414 Ok(true)
3415 }
3416
3417 fn canonical_root_matches_candidate(
3418 &self,
3419 candidate: &CandidateSchemaRevision,
3420 ) -> Result<bool, InternalError> {
3421 let first = self.canonical_root_slot_bytes(0)?;
3422 let second = self.canonical_root_slot_bytes(1)?;
3423 let Some(selection) =
3424 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3425 else {
3426 return Ok(false);
3427 };
3428 if selection.root() != candidate.root() {
3429 return Ok(false);
3430 }
3431 let key = RawSchemaKey::from_accepted_bundle(candidate.root().bundle_key());
3432 let bundle = self
3433 .get_canonical_raw_value(&key)?
3434 .ok_or_else(InternalError::store_corruption)?;
3435 let _verified =
3436 decode_verified_accepted_schema_revision_bundle(candidate.root(), bundle.as_bytes())?;
3437 Ok(true)
3438 }
3439
3440 fn get_canonical_raw_value(
3441 &self,
3442 key: &RawSchemaKey,
3443 ) -> Result<Option<RawSchemaSnapshot>, InternalError> {
3444 match &self.backend {
3445 SchemaStoreBackend::Journaled { canonical, .. } => Ok(canonical.get(key)),
3446 SchemaStoreBackend::Heap(_) => Err(InternalError::store_invariant()),
3447 }
3448 }
3449
3450 fn insert_canonical_raw_value(
3451 &mut self,
3452 key: RawSchemaKey,
3453 bytes: Vec<u8>,
3454 ) -> Result<(), InternalError> {
3455 self.invalidate_accepted_bundle_cache_for_key(key);
3456 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3457 return Err(InternalError::store_invariant());
3458 };
3459 canonical.insert(key, RawSchemaSnapshot::from_encoded_control_record(bytes));
3460 Ok(())
3461 }
3462
3463 fn insert_durable_raw_value(&mut self, key: RawSchemaKey, bytes: Vec<u8>) {
3467 self.invalidate_accepted_bundle_cache_for_key(key);
3468 let value = RawSchemaSnapshot::from_encoded_control_record(bytes);
3469 match &mut self.backend {
3470 SchemaStoreBackend::Heap(map) => {
3471 map.insert(key, value);
3472 }
3473 SchemaStoreBackend::Journaled {
3474 canonical,
3475 live,
3476 tombstones,
3477 ..
3478 } => {
3479 live.remove(&key);
3480 tombstones.remove(&key);
3481 canonical.insert(key, value);
3482 }
3483 }
3484 }
3485
3486 fn invalidate_accepted_bundle_cache_for_key(&mut self, key: RawSchemaKey) {
3487 if key.is_accepted_root() {
3488 self.accepted_bundle_cache.get_mut().take();
3489 }
3490 }
3491
3492 fn insert_durable_candidate_snapshots(
3493 &mut self,
3494 candidate: &CandidateSchemaRevision,
3495 ) -> Result<(), InternalError> {
3496 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3497 let key = RawSchemaKey::from_entity_version(*entity_tag, snapshot.version());
3498 let value = RawSchemaSnapshot::from_persisted_snapshot(snapshot)?;
3499 match &mut self.backend {
3500 SchemaStoreBackend::Heap(map) => {
3501 map.insert(key, value);
3502 }
3503 SchemaStoreBackend::Journaled {
3504 canonical,
3505 live,
3506 tombstones,
3507 ..
3508 } => {
3509 live.remove(&key);
3510 tombstones.remove(&key);
3511 canonical.insert(key, value);
3512 }
3513 }
3514 }
3515 Ok(())
3516 }
3517
3518 fn candidate_entry_keys(
3519 candidate: &CandidateSchemaRevision,
3520 root_slot: usize,
3521 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3522 let mut keys = candidate
3523 .bundle()
3524 .entity_snapshots()
3525 .iter()
3526 .map(|(entity_tag, snapshot)| {
3527 RawSchemaKey::from_entity_version(*entity_tag, snapshot.version())
3528 })
3529 .collect::<BTreeSet<_>>();
3530 keys.insert(RawSchemaKey::from_accepted_bundle(
3531 candidate.root().bundle_key(),
3532 ));
3533 keys.insert(RawSchemaKey::from_accepted_root_slot(root_slot)?);
3534 for (entity_tag, snapshot) in candidate.bundle().entity_snapshots() {
3535 for activation in snapshot
3536 .constraint_activations()
3537 .iter()
3538 .filter(|activation| activation.state() == ConstraintActivationState::Validating)
3539 {
3540 keys.insert(RawSchemaKey::from_constraint_validation_job(
3541 *entity_tag,
3542 activation.id(),
3543 ));
3544 }
3545 }
3546 Ok(keys)
3547 }
3548
3549 fn positioned_candidate_effect_keys(
3550 &self,
3551 incarnation: DatabaseIncarnationId,
3552 expected_revision: AcceptedSchemaRevision,
3553 candidate: &CandidateSchemaRevision,
3554 view: IdentityStateStorageView,
3555 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3556 let identity_transition =
3557 self.prepare_identity_state_transition(incarnation, candidate, view)?;
3558 let (first, second, candidate_is_current) = match view {
3559 IdentityStateStorageView::Effective => (
3560 self.accepted_root_slot_bytes(0)?,
3561 self.accepted_root_slot_bytes(1)?,
3562 self.current_root_matches_candidate(candidate)?,
3563 ),
3564 IdentityStateStorageView::Canonical => (
3565 self.canonical_root_slot_bytes(0)?,
3566 self.canonical_root_slot_bytes(1)?,
3567 self.canonical_root_matches_candidate(candidate)?,
3568 ),
3569 };
3570 let root_slot = if candidate_is_current {
3571 select_current_accepted_schema_root([first.as_deref(), second.as_deref()])?
3572 .ok_or_else(InternalError::store_corruption)?
3573 .slot()
3574 } else {
3575 prepare_accepted_schema_root_publication(
3576 [first.as_deref(), second.as_deref()],
3577 expected_revision,
3578 candidate,
3579 )
3580 .map_err(map_schema_publication_error)?
3581 .target_slot()
3582 };
3583 let mut keys = Self::candidate_entry_keys(candidate, root_slot)?;
3584 for state in identity_transition.into_updates() {
3585 keys.insert(RawSchemaKey::from_identity_state(
3586 state.owner().entity_tag(),
3587 state.owner().field_id(),
3588 ));
3589 }
3590
3591 let SchemaStoreBackend::Journaled {
3592 canonical,
3593 live,
3594 tombstones,
3595 positions,
3596 } = &self.backend
3597 else {
3598 return Err(InternalError::store_invariant());
3599 };
3600 for entry in canonical.iter() {
3601 let has_relevant_overlay = matches!(view, IdentityStateStorageView::Effective)
3602 || positions.is_positioned(entry.key())
3603 || live.contains_key(entry.key())
3604 || tombstones.contains(entry.key());
3605 if has_relevant_overlay
3606 && !keys.contains(entry.key())
3607 && !entry.key().is_identity_state()
3608 {
3609 keys.insert(*entry.key());
3610 }
3611 }
3612 if matches!(view, IdentityStateStorageView::Effective) {
3613 for key in live.keys() {
3614 if !keys.contains(key) && !key.is_identity_state() {
3615 keys.insert(*key);
3616 }
3617 }
3618 }
3619 Ok(keys)
3620 }
3621
3622 fn positioned_journal_batch_keys(
3623 &self,
3624 incarnation: DatabaseIncarnationId,
3625 batch: &JournalBatch,
3626 view: IdentityStateStorageView,
3627 ) -> Result<BTreeSet<RawSchemaKey>, InternalError> {
3628 let mut keys = BTreeSet::new();
3629 for record in batch.records() {
3630 match record {
3631 JournalRecord::SchemaPut {
3632 schema_snapshot_bytes,
3633 ..
3634 } => {
3635 let snapshot = decode_persisted_schema_snapshot(schema_snapshot_bytes)?;
3636 let entity_tag = match view {
3637 IdentityStateStorageView::Effective => self
3638 .current_accepted_schema_bundle_ref()?
3639 .ok_or_else(InternalError::store_corruption)?
3640 .entity_snapshots()
3641 .iter()
3642 .find_map(|(entity_tag, accepted)| {
3643 (accepted.entity_path() == snapshot.entity_path())
3644 .then_some(*entity_tag)
3645 }),
3646 IdentityStateStorageView::Canonical => self
3647 .current_canonical_accepted_schema_bundle()?
3648 .ok_or_else(InternalError::store_corruption)?
3649 .entity_snapshots()
3650 .iter()
3651 .find_map(|(entity_tag, accepted)| {
3652 (accepted.entity_path() == snapshot.entity_path())
3653 .then_some(*entity_tag)
3654 }),
3655 }
3656 .ok_or_else(InternalError::store_corruption)?;
3657 keys.insert(RawSchemaKey::from_entity_version(
3658 entity_tag,
3659 snapshot.version(),
3660 ));
3661 }
3662 JournalRecord::AcceptedSchemaPublish {
3663 expected_revision,
3664 schema_bundle_bytes,
3665 schema_root_bytes,
3666 ..
3667 } => {
3668 let candidate = CandidateSchemaRevision::from_encoded(
3669 schema_bundle_bytes.clone(),
3670 schema_root_bytes.clone(),
3671 )?;
3672 keys.extend(self.positioned_candidate_effect_keys(
3673 incarnation,
3674 *expected_revision,
3675 &candidate,
3676 view,
3677 )?);
3678 }
3679 JournalRecord::ConstraintValidationJobPut {
3680 entity_tag,
3681 constraint_id,
3682 ..
3683 }
3684 | JournalRecord::ConstraintValidationJobDelete {
3685 entity_tag,
3686 constraint_id,
3687 ..
3688 } => {
3689 keys.insert(RawSchemaKey::from_constraint_validation_job(
3690 *entity_tag,
3691 *constraint_id,
3692 ));
3693 }
3694 JournalRecord::IdentityRangeAdvance { range } => {
3695 keys.insert(RawSchemaKey::from_identity_state(
3696 range.owner().entity_tag(),
3697 range.owner().field_id(),
3698 ));
3699 }
3700 JournalRecord::RowPut { .. }
3701 | JournalRecord::RowDelete { .. }
3702 | JournalRecord::AcceptedSchemaIndexDelete { .. }
3703 | JournalRecord::AcceptedSchemaIndexPut { .. }
3704 | JournalRecord::ConstraintValidationIndexPut { .. } => {}
3705 #[cfg(any(test, feature = "migration"))]
3706 JournalRecord::SchemaMigrationRowPut { .. }
3707 | JournalRecord::SchemaMigrationIndexPut { .. } => {}
3708 }
3709 }
3710 Ok(keys)
3711 }
3712
3713 fn retain_durable_candidate_entries(
3717 &mut self,
3718 candidate: &CandidateSchemaRevision,
3719 root_slot: usize,
3720 ) -> Result<(), InternalError> {
3721 let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3722 self.accepted_bundle_cache.get_mut().take();
3723 match &mut self.backend {
3724 SchemaStoreBackend::Heap(map) => {
3725 map.retain(|key, _| keep.contains(key) || key.is_identity_state());
3726 }
3727 SchemaStoreBackend::Journaled {
3728 canonical,
3729 live,
3730 tombstones,
3731 ..
3732 } => {
3733 let stale = canonical
3734 .iter()
3735 .filter_map(|entry| {
3736 (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3737 .then_some(*entry.key())
3738 })
3739 .collect::<Vec<_>>();
3740 for key in stale {
3741 canonical.remove(&key);
3742 }
3743 live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3744 tombstones.clear();
3745 }
3746 }
3747 Ok(())
3748 }
3749
3750 fn retain_materialized_candidate_entries(
3751 &mut self,
3752 candidate: &CandidateSchemaRevision,
3753 root_slot: usize,
3754 ) -> Result<(), InternalError> {
3755 let keep = Self::candidate_entry_keys(candidate, root_slot)?;
3756 self.accepted_bundle_cache.get_mut().take();
3757 let SchemaStoreBackend::Journaled {
3758 canonical,
3759 live,
3760 tombstones,
3761 ..
3762 } = &mut self.backend
3763 else {
3764 return Err(InternalError::store_invariant());
3765 };
3766 live.retain(|key, _| keep.contains(key) || key.is_identity_state());
3767 let canonical_keys = canonical
3768 .iter()
3769 .map(|entry| *entry.key())
3770 .collect::<Vec<_>>();
3771 for key in canonical_keys {
3772 if keep.contains(&key) || key.is_identity_state() {
3773 tombstones.remove(&key);
3774 } else {
3775 tombstones.insert(key);
3776 }
3777 }
3778 Ok(())
3779 }
3780
3781 fn retain_canonical_candidate_entries(
3782 &mut self,
3783 keep: &BTreeSet<RawSchemaKey>,
3784 ) -> Result<(), InternalError> {
3785 self.accepted_bundle_cache.get_mut().take();
3786 let SchemaStoreBackend::Journaled { canonical, .. } = &mut self.backend else {
3787 return Err(InternalError::store_invariant());
3788 };
3789 let stale = canonical
3790 .iter()
3791 .filter_map(|entry| {
3792 (!keep.contains(entry.key()) && !entry.key().is_identity_state())
3793 .then_some(*entry.key())
3794 })
3795 .collect::<Vec<_>>();
3796 for key in stale {
3797 canonical.remove(&key);
3798 }
3799 Ok(())
3800 }
3801
3802 #[must_use]
3804 #[cfg(test)]
3805 fn contains_raw_snapshot(&self, key: &RawSchemaKey) -> bool {
3806 match &self.backend {
3807 SchemaStoreBackend::Heap(map) => map.contains_key(key),
3808 SchemaStoreBackend::Journaled { .. } => {
3809 self.get_raw_snapshot_for_backend(key).is_some()
3810 }
3811 }
3812 }
3813
3814 #[must_use]
3816 #[cfg(test)]
3817 pub(in crate::db) fn len(&self) -> u64 {
3818 match &self.backend {
3819 SchemaStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
3820 SchemaStoreBackend::Journaled { .. } => {
3821 let mut count = 0_u64;
3822 let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3823 count = count.saturating_add(1);
3824 Ok(SchemaStoreVisit::Continue)
3825 });
3826 count
3827 }
3828 }
3829 }
3830
3831 #[must_use]
3833 #[cfg(test)]
3834 pub(in crate::db) fn is_empty(&self) -> bool {
3835 match &self.backend {
3836 SchemaStoreBackend::Heap(map) => map.is_empty(),
3837 SchemaStoreBackend::Journaled { .. } => {
3838 let mut empty = true;
3839 let _: Result<(), Infallible> = self.visit_raw_snapshots(|_key, _snapshot| {
3840 empty = false;
3841 Ok(SchemaStoreVisit::Stop)
3842 });
3843 empty
3844 }
3845 }
3846 }
3847
3848 #[cfg(test)]
3850 pub(in crate::db) fn clear(&mut self) {
3851 self.accepted_bundle_cache.get_mut().take();
3852 match &mut self.backend {
3853 SchemaStoreBackend::Heap(map) => map.clear(),
3854 SchemaStoreBackend::Journaled {
3855 canonical,
3856 live,
3857 tombstones,
3858 ..
3859 } => {
3860 live.clear();
3861 tombstones.clear();
3862 let keys = canonical
3863 .iter()
3864 .map(|entry| *entry.key())
3865 .collect::<Vec<_>>();
3866 for key in keys {
3867 if key.is_entity_snapshot() {
3868 tombstones.insert(key);
3869 } else {
3870 canonical.remove(&key);
3871 }
3872 }
3873 }
3874 }
3875 }
3876
3877 fn current_accepted_schema_bundle_ref(
3878 &self,
3879 ) -> Result<Option<Ref<'_, AcceptedSchemaRevisionBundle>>, InternalError> {
3880 self.current_accepted_schema_authority_ref()
3881 .map(|authority| authority.map(|(_selection, bundle)| bundle))
3882 }
3883
3884 pub(in crate::db) fn current_accepted_schema_authority_ref(
3886 &self,
3887 ) -> Result<
3888 Option<(
3889 AcceptedSchemaRootSelection,
3890 Ref<'_, AcceptedSchemaRevisionBundle>,
3891 )>,
3892 InternalError,
3893 > {
3894 let selection = self.current_accepted_schema_root()?;
3895 self.accepted_schema_authority_ref_for_selection(selection)
3896 }
3897
3898 fn accepted_schema_authority_ref_for_selection(
3899 &self,
3900 selection: Option<AcceptedSchemaRootSelection>,
3901 ) -> Result<
3902 Option<(
3903 AcceptedSchemaRootSelection,
3904 Ref<'_, AcceptedSchemaRevisionBundle>,
3905 )>,
3906 InternalError,
3907 > {
3908 let Some(selection) = selection else {
3909 self.accepted_bundle_cache
3910 .try_borrow_mut()
3911 .map_err(|_| InternalError::store_invariant())?
3912 .take();
3913 return Ok(None);
3914 };
3915
3916 let cache_matches = self
3917 .accepted_bundle_cache
3918 .try_borrow()
3919 .map_err(|_| InternalError::store_invariant())?
3920 .as_ref()
3921 .is_some_and(|cached| cached.selection == selection);
3922 if !cache_matches {
3923 let key = RawSchemaKey::from_accepted_bundle(selection.root().bundle_key());
3924 let raw = self
3925 .get_raw_snapshot(&key)
3926 .ok_or_else(InternalError::store_corruption)?;
3927 let bundle =
3928 decode_verified_accepted_schema_revision_bundle(selection.root(), raw.as_bytes())?;
3929 self.validate_constraint_validation_job_closure(&bundle)?;
3930 #[cfg(test)]
3931 ACCEPTED_SCHEMA_BUNDLE_CACHE_MISSES
3932 .with(|misses| misses.set(misses.get().saturating_add(1)));
3933 let value_catalog = AcceptedValueCatalogHandle::new(
3934 bundle.enum_catalog().clone(),
3935 bundle.composite_catalog().clone(),
3936 self.accepted_catalog_scope
3937 .get_or_init(AcceptedStoreCatalogScope::new)
3938 .clone(),
3939 bundle.revision(),
3940 selection.root().fingerprint(),
3941 );
3942 let cardinality_domain = Rc::new(CardinalityAcceptedDomain::derive(&bundle)?);
3943 *self
3944 .accepted_bundle_cache
3945 .try_borrow_mut()
3946 .map_err(|_| InternalError::store_invariant())? = Some(AcceptedSchemaBundleCache {
3947 selection,
3948 bundle,
3949 cardinality_domain,
3950 value_catalog,
3951 entity_selections: RefCell::new(StdBTreeMap::new()),
3952 });
3953 }
3954
3955 let cache = self
3956 .accepted_bundle_cache
3957 .try_borrow()
3958 .map_err(|_| InternalError::store_invariant())?;
3959 let bundle = Ref::filter_map(cache, |cache| {
3960 cache
3961 .as_ref()
3962 .filter(|cached| cached.selection == selection)
3963 .map(|cached| &cached.bundle)
3964 })
3965 .map_err(|_| InternalError::store_invariant())?;
3966 self.validate_identity_state_closure(&bundle)?;
3967 Ok(Some((selection, bundle)))
3968 }
3969
3970 pub(in crate::db) fn accepted_cardinality_domain_for_selection(
3972 &self,
3973 selection: Option<AcceptedSchemaRootSelection>,
3974 ) -> Result<Option<(AcceptedSchemaRootSelection, Rc<CardinalityAcceptedDomain>)>, InternalError>
3975 {
3976 let Some(selection) = selection else {
3977 self.accepted_bundle_cache
3978 .try_borrow_mut()
3979 .map_err(|_| InternalError::store_invariant())?
3980 .take();
3981 return Ok(None);
3982 };
3983 let cache_matches = self
3984 .accepted_bundle_cache
3985 .try_borrow()
3986 .map_err(|_| InternalError::store_invariant())?
3987 .as_ref()
3988 .is_some_and(|cached| cached.selection == selection);
3989 if !cache_matches {
3990 let authority = self
3991 .accepted_schema_authority_ref_for_selection(Some(selection))?
3992 .ok_or_else(InternalError::store_invariant)?;
3993 drop(authority);
3994 }
3995 let cache = self
3996 .accepted_bundle_cache
3997 .try_borrow()
3998 .map_err(|_| InternalError::store_invariant())?;
3999 let domain = cache
4000 .as_ref()
4001 .filter(|cached| cached.selection == selection)
4002 .map(|cached| Rc::clone(&cached.cardinality_domain))
4003 .ok_or_else(InternalError::store_invariant)?;
4004 Ok(Some((selection, domain)))
4005 }
4006
4007 pub(in crate::db) fn cached_cardinality_domain_for_root(
4009 &self,
4010 root: CardinalityAcceptedRootIdentity,
4011 ) -> Result<Option<Rc<CardinalityAcceptedDomain>>, InternalError> {
4012 let cache = self
4013 .accepted_bundle_cache
4014 .try_borrow()
4015 .map_err(|_| InternalError::store_invariant())?;
4016 Ok(cache
4017 .as_ref()
4018 .filter(|cached| root.matches(cached.selection.root()))
4019 .map(|cached| Rc::clone(&cached.cardinality_domain)))
4020 }
4021
4022 fn latest_raw_snapshots_by_entity(
4023 &self,
4024 ) -> StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)> {
4025 let mut latest_by_entity =
4026 StdBTreeMap::<EntityTag, (SchemaVersion, RawSchemaSnapshot)>::new();
4027
4028 let _: Result<(), std::convert::Infallible> = self.visit_raw_snapshots(|key, snapshot| {
4029 let version = SchemaVersion::new(key.version());
4030 match latest_by_entity.get_mut(&key.entity_tag()) {
4031 Some((latest_version, latest_snapshot)) if version > *latest_version => {
4032 *latest_version = version;
4033 *latest_snapshot = snapshot.clone();
4034 }
4035 None => {
4036 latest_by_entity.insert(key.entity_tag(), (version, snapshot.clone()));
4037 }
4038 Some(_) => {}
4039 }
4040 Ok(SchemaStoreVisit::Continue)
4041 });
4042
4043 latest_by_entity
4044 }
4045
4046 fn visit_raw_snapshots<E>(
4049 &self,
4050 visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4051 ) -> Result<(), E> {
4052 let bounds = RawSchemaKey::all_entity_range_bounds();
4053 match &self.backend {
4054 SchemaStoreBackend::Heap(map) => {
4055 let mut visitor = visitor;
4056 for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4057 if visitor(key, snapshot)?.should_stop() {
4058 break;
4059 }
4060 }
4061 }
4062 SchemaStoreBackend::Journaled {
4063 canonical,
4064 live,
4065 tombstones,
4066 ..
4067 } => Self::visit_journaled_raw_snapshot_range(
4068 canonical,
4069 live,
4070 tombstones,
4071 bounds,
4072 Direction::Asc,
4073 visitor,
4074 )?,
4075 }
4076
4077 Ok(())
4078 }
4079
4080 fn visit_constraint_validation_jobs_in_view<E>(
4081 &self,
4082 view: IdentityStateStorageView,
4083 visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4084 ) -> Result<(), E> {
4085 let bounds = RawSchemaKey::all_constraint_validation_job_range_bounds();
4086 match (&self.backend, view) {
4087 (SchemaStoreBackend::Heap(map), _) => {
4088 let mut visitor = visitor;
4089 for (key, snapshot) in map.range((bounds.0, bounds.1)) {
4090 if visitor(key, snapshot)?.should_stop() {
4091 break;
4092 }
4093 }
4094 }
4095 (
4096 SchemaStoreBackend::Journaled {
4097 canonical,
4098 live,
4099 tombstones,
4100 ..
4101 },
4102 IdentityStateStorageView::Effective,
4103 ) => Self::visit_journaled_raw_snapshot_range(
4104 canonical,
4105 live,
4106 tombstones,
4107 bounds,
4108 Direction::Asc,
4109 visitor,
4110 )?,
4111 (
4112 SchemaStoreBackend::Journaled { canonical, .. },
4113 IdentityStateStorageView::Canonical,
4114 ) => {
4115 let mut visitor = visitor;
4116 for entry in canonical.range((bounds.0, bounds.1)) {
4117 if visitor(entry.key(), &entry.value())?.should_stop() {
4118 break;
4119 }
4120 }
4121 }
4122 }
4123 Ok(())
4124 }
4125
4126 #[cfg(test)]
4127 #[must_use]
4128 pub(in crate::db) fn canonical_len_for_tests(&self) -> u64 {
4129 match &self.backend {
4130 SchemaStoreBackend::Journaled { canonical: map, .. } => map.len(),
4131 SchemaStoreBackend::Heap(_) => 0,
4132 }
4133 }
4134
4135 fn get_raw_snapshot_for_backend(&self, key: &RawSchemaKey) -> Option<RawSchemaSnapshot> {
4136 let SchemaStoreBackend::Journaled {
4137 canonical,
4138 live,
4139 tombstones,
4140 ..
4141 } = &self.backend
4142 else {
4143 return None;
4144 };
4145
4146 if tombstones.contains(key) {
4147 return None;
4148 }
4149 live.get(key).cloned().or_else(|| canonical.get(key))
4150 }
4151
4152 fn visit_journaled_raw_snapshot_range<E>(
4153 canonical: &StableBTreeMap<
4154 RawSchemaKey,
4155 RawSchemaSnapshot,
4156 RuntimeMemory<DefaultMemoryImpl>,
4157 >,
4158 live: &StdBTreeMap<RawSchemaKey, RawSchemaSnapshot>,
4159 tombstones: &BTreeSet<RawSchemaKey>,
4160 bounds: (RangeBound<RawSchemaKey>, RangeBound<RawSchemaKey>),
4161 direction: Direction,
4162 mut visitor: impl FnMut(&RawSchemaKey, &RawSchemaSnapshot) -> Result<SchemaStoreVisit, E>,
4163 ) -> Result<(), E> {
4164 match direction {
4165 Direction::Asc => {
4166 for entry in ordered_overlay_entries(
4167 canonical.range((bounds.0, bounds.1)),
4168 live.range((bounds.0, bounds.1)),
4169 Direction::Asc,
4170 |entry| entry.key(),
4171 |entry| entry.0,
4172 tombstones,
4173 ) {
4174 let visit = match entry {
4175 OrderedOverlayEntry::Canonical(canonical_entry) => {
4176 visitor(canonical_entry.key(), &canonical_entry.value())?
4177 }
4178 OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4179 };
4180 if visit.should_stop() {
4181 return Ok(());
4182 }
4183 }
4184 }
4185 Direction::Desc => {
4186 for entry in ordered_overlay_entries(
4187 canonical.range((bounds.0, bounds.1)).rev(),
4188 live.range((bounds.0, bounds.1)).rev(),
4189 Direction::Desc,
4190 |entry| entry.key(),
4191 |entry| entry.0,
4192 tombstones,
4193 ) {
4194 let visit = match entry {
4195 OrderedOverlayEntry::Canonical(canonical_entry) => {
4196 visitor(canonical_entry.key(), &canonical_entry.value())?
4197 }
4198 OrderedOverlayEntry::Live((key, snapshot)) => visitor(key, snapshot)?,
4199 };
4200 if visit.should_stop() {
4201 return Ok(());
4202 }
4203 }
4204 }
4205 }
4206
4207 Ok(())
4208 }
4209}
4210
4211fn map_schema_publication_error(error: AcceptedSchemaPublicationError) -> InternalError {
4212 match error {
4213 AcceptedSchemaPublicationError::StaleSchemaRevision { .. }
4214 | AcceptedSchemaPublicationError::RevisionExhausted => InternalError::store_unsupported(),
4215 AcceptedSchemaPublicationError::InvalidCandidate => InternalError::store_invariant(),
4216 AcceptedSchemaPublicationError::CorruptRootSlots => InternalError::store_corruption(),
4217 }
4218}
4219
4220fn derive_data_allocation_metadata(
4221 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4222) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4223 let mut max_version = SchemaVersion::initial();
4224 let mut hasher = new_hash_sha256();
4225 write_hash_tag_u8(&mut hasher, SCHEMA_STORE_DATA_ALLOCATION_FINGERPRINT_DOMAIN);
4226
4227 for (entity, (_, snapshot)) in latest_by_entity {
4228 let persisted = snapshot.decode_persisted_snapshot()?;
4229 if persisted.version() > max_version {
4230 max_version = persisted.version();
4231 }
4232
4233 let data_projection = PersistedSchemaSnapshot::new_with_primary_key_fields_and_indexes(
4234 persisted.version(),
4235 persisted.entity_path().to_string(),
4236 persisted.entity_name().to_string(),
4237 persisted.primary_key_field_ids().to_vec(),
4238 persisted.row_layout().clone(),
4239 persisted.fields().to_vec(),
4240 Vec::new(),
4241 );
4242 let constraint_catalog = crate::db::schema::AcceptedConstraintCatalog::initial(
4243 data_projection.fields(),
4244 data_projection.indexes(),
4245 data_projection.relations(),
4246 )
4247 .map_err(|_| InternalError::store_invariant())?;
4248 let data_projection = data_projection.with_constraint_catalog(constraint_catalog);
4249 let encoded = encode_persisted_schema_snapshot(&data_projection)?;
4250
4251 write_hash_u64(&mut hasher, entity.value());
4252 write_hash_u32(&mut hasher, persisted.version().get());
4253 write_hash_len_u32(&mut hasher, encoded.len());
4254 hasher.update(encoded);
4255 }
4256
4257 Ok(finalize_schema_metadata(
4258 max_version,
4259 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4260 hasher,
4261 latest_by_entity.len(),
4262 ))
4263}
4264
4265fn derive_index_allocation_metadata(
4266 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4267) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4268 let mut max_version = SchemaVersion::initial();
4269 let mut hasher = new_hash_sha256();
4270 write_hash_tag_u8(
4271 &mut hasher,
4272 SCHEMA_STORE_INDEX_ALLOCATION_FINGERPRINT_DOMAIN,
4273 );
4274
4275 for (entity, (_, snapshot)) in latest_by_entity {
4276 let persisted = snapshot.decode_persisted_snapshot()?;
4277 if persisted.version() > max_version {
4278 max_version = persisted.version();
4279 }
4280
4281 write_hash_u64(&mut hasher, entity.value());
4282 write_hash_u32(&mut hasher, persisted.version().get());
4283 write_hash_len_u32(&mut hasher, persisted.indexes().len());
4284 for index in persisted.indexes() {
4285 write_hash_u32(&mut hasher, u32::from(index.ordinal()));
4286 write_hash_str_u32(&mut hasher, index.name());
4287 write_hash_str_u32(&mut hasher, index.store());
4288 write_hash_tag_u8(&mut hasher, u8::from(index.unique()));
4289 write_hash_str_u32(&mut hasher, persisted_index_origin_name(index.origin()));
4290 match index.predicate_sql() {
4291 Some(predicate_sql) => {
4292 write_hash_tag_u8(&mut hasher, 1);
4293 write_hash_str_u32(&mut hasher, predicate_sql);
4294 }
4295 None => write_hash_tag_u8(&mut hasher, 0),
4296 }
4297 hash_persisted_index_key(&mut hasher, index.key());
4298 }
4299 }
4300
4301 Ok(finalize_schema_metadata(
4302 max_version,
4303 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4304 hasher,
4305 latest_by_entity.len(),
4306 ))
4307}
4308
4309fn derive_schema_catalog_metadata(
4310 latest_by_entity: &StdBTreeMap<EntityTag, (SchemaVersion, RawSchemaSnapshot)>,
4311) -> Result<SchemaStoreCatalogMetadata, InternalError> {
4312 let mut max_version = SchemaVersion::initial();
4313 let mut hasher = new_hash_sha256();
4314 write_hash_tag_u8(&mut hasher, SCHEMA_STORE_CATALOG_FINGERPRINT_DOMAIN);
4315
4316 for (entity, (version, snapshot)) in latest_by_entity {
4317 let persisted = snapshot.decode_persisted_snapshot()?;
4318 if persisted.version() > max_version {
4319 max_version = persisted.version();
4320 }
4321
4322 write_hash_u64(&mut hasher, entity.value());
4323 write_hash_u32(&mut hasher, version.get());
4324 write_hash_len_u32(&mut hasher, snapshot.as_bytes().len());
4325 hasher.update(snapshot.as_bytes());
4326 }
4327
4328 Ok(finalize_schema_metadata(
4329 max_version,
4330 SCHEMA_STORE_FINGERPRINT_METHOD_VERSION,
4331 hasher,
4332 latest_by_entity.len(),
4333 ))
4334}
4335
4336fn finalize_schema_metadata(
4337 schema_version: SchemaVersion,
4338 schema_fingerprint_method_version: u8,
4339 hasher: sha2::Sha256,
4340 entity_count: usize,
4341) -> SchemaStoreCatalogMetadata {
4342 let digest = finalize_hash_sha256(hasher);
4343 let mut schema_fingerprint = [0u8; 16];
4344 schema_fingerprint.copy_from_slice(&digest[..16]);
4345
4346 SchemaStoreCatalogMetadata::new(
4347 schema_version,
4348 schema_fingerprint_method_version,
4349 schema_fingerprint,
4350 u64::try_from(entity_count).unwrap_or(u64::MAX),
4351 )
4352}
4353
4354fn hash_persisted_index_key(hasher: &mut sha2::Sha256, key: &PersistedIndexKeySnapshot) {
4355 match key {
4356 PersistedIndexKeySnapshot::FieldPath(paths) => {
4357 write_hash_tag_u8(hasher, 1);
4358 write_hash_len_u32(hasher, paths.len());
4359 for path in paths {
4360 hash_persisted_index_field_path(hasher, path);
4361 }
4362 }
4363 PersistedIndexKeySnapshot::Items(items) => {
4364 write_hash_tag_u8(hasher, 2);
4365 write_hash_len_u32(hasher, items.len());
4366 for item in items {
4367 match item {
4368 PersistedIndexKeyItemSnapshot::FieldPath(path) => {
4369 write_hash_tag_u8(hasher, 1);
4370 hash_persisted_index_field_path(hasher, path);
4371 }
4372 PersistedIndexKeyItemSnapshot::Expression(expression) => {
4373 write_hash_tag_u8(hasher, 2);
4374 write_hash_str_u32(hasher, persisted_expression_op_name(expression.op()));
4375 hash_persisted_index_field_path(hasher, expression.source());
4376 hash_accepted_field_kind(hasher, expression.input_kind());
4377 hash_accepted_field_kind(hasher, expression.output_kind());
4378 write_hash_str_u32(hasher, expression.canonical_text());
4379 }
4380 }
4381 }
4382 }
4383 }
4384}
4385
4386fn hash_persisted_index_field_path(
4387 hasher: &mut sha2::Sha256,
4388 path: &crate::db::schema::PersistedIndexFieldPathSnapshot,
4389) {
4390 write_hash_u32(hasher, path.field_id().get());
4391 write_hash_u32(hasher, u32::from(path.slot().get()));
4392 write_hash_len_u32(hasher, path.path().len());
4393 for segment in path.path() {
4394 write_hash_str_u32(hasher, segment);
4395 }
4396 hash_accepted_field_kind(hasher, path.kind());
4397 write_hash_tag_u8(hasher, u8::from(path.nullable()));
4398}
4399
4400fn hash_accepted_field_kind(hasher: &mut sha2::Sha256, kind: &AcceptedFieldKind) {
4401 match kind {
4402 AcceptedFieldKind::Account => write_hash_tag_u8(hasher, 1),
4403 AcceptedFieldKind::Blob { max_len } => {
4404 write_hash_tag_u8(hasher, 2);
4405 hash_optional_u32(hasher, *max_len);
4406 }
4407 AcceptedFieldKind::Bool => {
4408 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_BOOL);
4409 }
4410 AcceptedFieldKind::Date => write_hash_tag_u8(hasher, 4),
4411 AcceptedFieldKind::Decimal { scale } => {
4412 write_hash_tag_u8(hasher, 5);
4413 write_hash_u32(hasher, *scale);
4414 }
4415 AcceptedFieldKind::Duration => write_hash_tag_u8(hasher, 6),
4416 AcceptedFieldKind::Enum { type_id } => {
4417 write_hash_tag_u8(hasher, 7);
4418 write_hash_u32(hasher, type_id.get());
4419 }
4420 AcceptedFieldKind::Float32 => write_hash_tag_u8(hasher, 8),
4421 AcceptedFieldKind::Float64 => write_hash_tag_u8(hasher, 9),
4422 AcceptedFieldKind::Int8 => write_hash_tag_u8(hasher, 10),
4423 AcceptedFieldKind::Int16 => write_hash_tag_u8(hasher, 11),
4424 AcceptedFieldKind::Int32 => write_hash_tag_u8(hasher, 12),
4425 AcceptedFieldKind::Int64 => write_hash_tag_u8(hasher, 13),
4426 AcceptedFieldKind::Int128 => write_hash_tag_u8(hasher, 14),
4427 AcceptedFieldKind::IntBig { max_bytes } => {
4428 write_hash_tag_u8(hasher, 15);
4429 write_hash_u32(hasher, *max_bytes);
4430 }
4431 AcceptedFieldKind::Principal => write_hash_tag_u8(hasher, 16),
4432 AcceptedFieldKind::Subaccount => write_hash_tag_u8(hasher, 17),
4433 AcceptedFieldKind::Text { max_len } => {
4434 write_hash_tag_u8(hasher, 18);
4435 hash_optional_u32(hasher, *max_len);
4436 }
4437 AcceptedFieldKind::Timestamp => write_hash_tag_u8(hasher, 19),
4438 AcceptedFieldKind::Nat8 => write_hash_tag_u8(hasher, 20),
4439 AcceptedFieldKind::Nat16 => write_hash_tag_u8(hasher, 21),
4440 AcceptedFieldKind::Nat32 => write_hash_tag_u8(hasher, 22),
4441 AcceptedFieldKind::Nat64 => write_hash_tag_u8(hasher, 23),
4442 AcceptedFieldKind::Nat128 => write_hash_tag_u8(hasher, 24),
4443 AcceptedFieldKind::NatBig { max_bytes } => {
4444 write_hash_tag_u8(hasher, 25);
4445 write_hash_u32(hasher, *max_bytes);
4446 }
4447 AcceptedFieldKind::Ulid => write_hash_tag_u8(hasher, 26),
4448 AcceptedFieldKind::Unit => write_hash_tag_u8(hasher, 27),
4449 AcceptedFieldKind::Relation {
4450 target_path,
4451 target_entity_name,
4452 target_entity_tag,
4453 target_store_path,
4454 key_kind,
4455 } => {
4456 write_hash_tag_u8(hasher, 28);
4457 write_hash_str_u32(hasher, target_path);
4458 write_hash_str_u32(hasher, target_entity_name);
4459 write_hash_u64(hasher, target_entity_tag.value());
4460 write_hash_str_u32(hasher, target_store_path);
4461 hash_accepted_field_kind(hasher, key_kind);
4462 }
4463 AcceptedFieldKind::List(inner) => {
4464 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_LIST);
4465 hash_accepted_field_kind(hasher, inner);
4466 }
4467 AcceptedFieldKind::Set(inner) => {
4468 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_SET);
4469 hash_accepted_field_kind(hasher, inner);
4470 }
4471 AcceptedFieldKind::Map { key, value } => {
4472 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_MAP);
4473 hash_accepted_field_kind(hasher, key);
4474 hash_accepted_field_kind(hasher, value);
4475 }
4476 AcceptedFieldKind::Composite { type_id } => {
4477 write_hash_tag_u8(hasher, ACCEPTED_FIELD_KIND_FINGERPRINT_TAG_COMPOSITE);
4478 write_hash_u32(hasher, type_id.get());
4479 }
4480 AcceptedFieldKind::U256 => write_hash_tag_u8(hasher, 33),
4481 }
4482}
4483
4484fn hash_optional_u32(hasher: &mut sha2::Sha256, value: Option<u32>) {
4485 match value {
4486 Some(value) => {
4487 write_hash_tag_u8(hasher, 1);
4488 write_hash_u32(hasher, value);
4489 }
4490 None => write_hash_tag_u8(hasher, 0),
4491 }
4492}
4493
4494const fn persisted_index_origin_name(
4495 origin: crate::db::schema::PersistedIndexOrigin,
4496) -> &'static str {
4497 match origin {
4498 crate::db::schema::PersistedIndexOrigin::Generated => "generated",
4499 crate::db::schema::PersistedIndexOrigin::SqlDdl => "sql_ddl",
4500 }
4501}
4502
4503const fn persisted_expression_op_name(
4504 op: crate::db::schema::PersistedIndexExpressionOp,
4505) -> &'static str {
4506 match op {
4507 crate::db::schema::PersistedIndexExpressionOp::Lower => "lower",
4508 crate::db::schema::PersistedIndexExpressionOp::Upper => "upper",
4509 crate::db::schema::PersistedIndexExpressionOp::Trim => "trim",
4510 crate::db::schema::PersistedIndexExpressionOp::LowerTrim => "lower_trim",
4511 crate::db::schema::PersistedIndexExpressionOp::Date => "date",
4512 crate::db::schema::PersistedIndexExpressionOp::Year => "year",
4513 crate::db::schema::PersistedIndexExpressionOp::Month => "month",
4514 crate::db::schema::PersistedIndexExpressionOp::Day => "day",
4515 }
4516}
4517
4518#[cfg(test)]
4523mod tests;