1use std::{
2 collections::BTreeMap,
3 fs::{File, OpenOptions},
4 path::Path,
5 sync::{Arc, Mutex},
6};
7
8use fs2::FileExt as _;
9use iroh_db_core::{
10 AuthorId, BlobHash, CapabilityId, CollectionId, CommitCodecError, CommitEnvelope, CommitId,
11 DomainId, MigrationCodecError, MigrationId, SchemaActivation, SchemaDescriptor, SchemaError,
12 SchemaStage, VersionVector, staging_digest,
13};
14use iroh_db_security::{
15 AuthorityCheckpoint, CapabilityCertificate, ControlTransition, DomainDescriptor, Invitation,
16 Permission,
17};
18use redb::{Database, ReadableDatabase as _, ReadableTable as _, TableDefinition};
19
20use crate::{BlobStore, BlobStoreError, IndexValue, MaterializedRecord, StoredRecord};
21
22pub const DEFAULT_METADATA_CACHE_SIZE: usize = 8 * 1024 * 1024;
24
25const FORMAT_V5: &str = "iroh-db 5.0\n";
26const FORMAT_CONTENTS: &str = "iroh-db 6.0\n";
27const SNAPSHOT_BASELINE_VERSION: u16 = 1;
28const MAX_SNAPSHOT_BASELINE_BYTES: usize = 8 * 1024 * 1024;
29const MAX_SNAPSHOT_BASELINE_ENTRIES: usize = 65_536;
30const COMMIT_METADATA: TableDefinition<'static, &[u8], &[u8]> =
31 TableDefinition::new("commit_metadata_v1");
32const COMMIT_CONTEXTS: TableDefinition<'static, &[u8], &[u8]> =
33 TableDefinition::new("commit_contexts_v1");
34const RECORD_HEADS: TableDefinition<'static, &[u8], &[u8]> =
35 TableDefinition::new("record_heads_v1");
36const RECONCILIATION_READY: TableDefinition<'static, &[u8], u8> =
37 TableDefinition::new("reconciliation_ready_v1");
38const SNAPSHOT_COVERAGE: TableDefinition<'static, &[u8], &[u8]> =
39 TableDefinition::new("snapshot_coverage_v1");
40const SNAPSHOT_BASELINES: TableDefinition<'static, &[u8], &[u8]> =
41 TableDefinition::new("snapshot_baselines_v1");
42const BLOB_OWNERSHIP: TableDefinition<'static, &[u8], u8> =
43 TableDefinition::new("blob_ownership_v1");
44const AUTHOR_SEQUENCES: TableDefinition<'static, &[u8], u64> =
45 TableDefinition::new("author_sequences_v1");
46const AUTHOR_COMMITS: TableDefinition<'static, &[u8], &[u8]> =
47 TableDefinition::new("author_commits_v1");
48const AUTHOR_HEADS: TableDefinition<'static, &[u8], &[u8]> =
49 TableDefinition::new("author_heads_v1");
50const FRONTIER: TableDefinition<'static, &[u8], u8> = TableDefinition::new("frontier_v1");
51const QUARANTINED_COMMITS: TableDefinition<'static, &[u8], &[u8]> =
52 TableDefinition::new("quarantined_commits_v1");
53const QUARANTINED_AUTHORS: TableDefinition<'static, &[u8], u64> =
54 TableDefinition::new("quarantined_authors_v1");
55const SCHEMAS: TableDefinition<'static, &[u8], &[u8]> = TableDefinition::new("schemas_v1");
56const RECORDS: TableDefinition<'static, &[u8], &[u8]> = TableDefinition::new("records_v1");
57const RECORD_INDEX_KEYS: TableDefinition<'static, &[u8], &[u8]> =
58 TableDefinition::new("record_index_keys_v1");
59const SECONDARY_INDEX: TableDefinition<'static, &[u8], u8> =
60 TableDefinition::new("secondary_index_v1");
61const TRUSTED_AUTHORS: TableDefinition<'static, &[u8], u8> =
62 TableDefinition::new("trusted_authors_v1");
63const DOMAIN_DESCRIPTORS: TableDefinition<'static, &[u8], &[u8]> =
64 TableDefinition::new("domain_descriptors_v1");
65const CAPABILITIES: TableDefinition<'static, &[u8], &[u8]> =
66 TableDefinition::new("capabilities_v1");
67const ACTIVE_CAPABILITIES: TableDefinition<'static, &[u8], &[u8]> =
68 TableDefinition::new("active_capabilities_v1");
69const REVOKED_CAPABILITIES: TableDefinition<'static, &[u8], u64> =
70 TableDefinition::new("revoked_capabilities_v1");
71const SUBJECT_REVOCATIONS: TableDefinition<'static, &[u8], u64> =
72 TableDefinition::new("subject_revocations_v1");
73const AUTHORITY_CHECKPOINTS: TableDefinition<'static, &[u8], &[u8]> =
74 TableDefinition::new("authority_checkpoints_v1");
75const CONTROL_HEADS: TableDefinition<'static, &[u8], &[u8]> =
76 TableDefinition::new("control_heads_v1");
77const CONTROL_OBJECTS: TableDefinition<'static, &[u8], &[u8]> =
78 TableDefinition::new("control_objects_v1");
79const CONTROL_SLOTS: TableDefinition<'static, &[u8], &[u8]> =
80 TableDefinition::new("control_slots_v1");
81const CONTROL_FORKS: TableDefinition<'static, &[u8], &[u8]> =
82 TableDefinition::new("control_forks_v1");
83const DEVICE_GROUPS: TableDefinition<'static, &[u8], &[u8]> =
84 TableDefinition::new("device_groups_v1");
85const CONSUMED_INVITATIONS: TableDefinition<'static, &[u8], u8> =
86 TableDefinition::new("consumed_invitations_v1");
87const ACTIVE_SCHEMAS: TableDefinition<'static, &[u8], &[u8]> =
88 TableDefinition::new("active_schemas_v1");
89const MIGRATION_STAGES: TableDefinition<'static, &[u8], &[u8]> =
90 TableDefinition::new("migration_stages_v1");
91const MIGRATION_ACTIVATIONS: TableDefinition<'static, &[u8], &[u8]> =
92 TableDefinition::new("migration_activations_v1");
93const MIGRATION_CHECKPOINTS: TableDefinition<'static, &[u8], &[u8]> =
94 TableDefinition::new("migration_checkpoints_v1");
95const MIGRATION_FORKS: TableDefinition<'static, &[u8], &[u8]> =
96 TableDefinition::new("migration_forks_v1");
97const OPERATIONAL_METADATA: TableDefinition<'static, &[u8], u64> =
98 TableDefinition::new("operational_metadata_v1");
99const SNAPSHOT_OBSERVED_AT_TAG: u8 = 1;
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
103pub enum ApplyOutcome {
104 Applied(CommitId),
106 Duplicate(CommitId),
108}
109
110#[derive(Debug, Clone, PartialEq, Eq)]
112pub struct RecordHead {
113 commit_id: CommitId,
114 author: AuthorId,
115 sequence: u64,
116 mutation: MaterializedRecord,
117}
118
119impl RecordHead {
120 pub fn new(
122 commit_id: CommitId,
123 author: AuthorId,
124 sequence: u64,
125 mutation: MaterializedRecord,
126 ) -> Self {
127 Self {
128 commit_id,
129 author,
130 sequence,
131 mutation,
132 }
133 }
134
135 pub const fn commit_id(&self) -> CommitId {
136 self.commit_id
137 }
138
139 pub const fn author(&self) -> AuthorId {
140 self.author
141 }
142
143 pub const fn sequence(&self) -> u64 {
144 self.sequence
145 }
146
147 pub const fn mutation(&self) -> &MaterializedRecord {
148 &self.mutation
149 }
150}
151
152#[derive(Debug, Clone, PartialEq, Eq)]
154pub struct RecordHeadSet {
155 collection_id: CollectionId,
156 record_id: Vec<u8>,
157 heads: Vec<RecordHead>,
158}
159
160impl RecordHeadSet {
161 pub fn new(
162 collection_id: CollectionId,
163 record_id: Vec<u8>,
164 heads: Vec<RecordHead>,
165 ) -> Result<Self, StoreError> {
166 if heads.is_empty()
167 || heads.iter().any(|head| {
168 head.mutation.schema().collection_id() != collection_id
169 || head.mutation.record_id() != record_id
170 })
171 {
172 return Err(StoreError::CorruptDerivedMetadata(
173 "record head replacement has inconsistent targets".into(),
174 ));
175 }
176 Ok(Self {
177 collection_id,
178 record_id,
179 heads,
180 })
181 }
182
183 pub const fn collection_id(&self) -> CollectionId {
184 self.collection_id
185 }
186
187 pub fn record_id(&self) -> &[u8] {
188 &self.record_id
189 }
190
191 pub fn heads(&self) -> &[RecordHead] {
192 &self.heads
193 }
194}
195
196enum CommitPreflight {
197 Duplicate(CommitId),
198 Ready {
199 previous_author_head: Option<CommitId>,
200 },
201}
202
203#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
205pub enum StoreError {
206 #[error("store filesystem operation failed: {0}")]
208 Io(String),
209 #[error("metadata database operation failed: {0}")]
211 Database(String),
212 #[error(transparent)]
214 Blob(#[from] BlobStoreError),
215 #[error("blob store returned a mismatched content hash")]
217 BlobHashMismatch,
218 #[error("author sequence gap: expected {expected}, got {actual}")]
220 SequenceGap {
221 expected: u64,
223 actual: u64,
225 },
226 #[error("commit does not extend previous author commit {previous}")]
228 AuthorChainMissing {
229 previous: CommitId,
231 },
232 #[error("commit does not extend installed snapshot frontier commit {frontier}")]
234 SnapshotBaselineMissing {
235 frontier: CommitId,
237 },
238 #[error(
240 "author {author} equivocated at sequence {sequence}: conflicting commits {first} and {second}"
241 )]
242 AuthorEquivocation {
243 author: AuthorId,
245 sequence: u64,
247 first: CommitId,
249 second: CommitId,
251 },
252 #[error("commit is quarantined: {0}")]
254 QuarantinedCommit(CommitId),
255 #[error("author {author} is quarantined from sequence {from_sequence}")]
257 QuarantinedAuthor {
258 author: AuthorId,
260 from_sequence: u64,
262 },
263 #[error("commit dependency is missing: {0}")]
265 MissingDependency(CommitId),
266 #[error("unsupported store directory format")]
268 UnsupportedFormat,
269 #[error("database is already open for writing")]
271 AlreadyOpen,
272 #[error("commit blob is missing: {0}")]
274 MissingCommitBlob(CommitId),
275 #[error(transparent)]
277 CommitCodec(#[from] CommitCodecError),
278 #[error(transparent)]
280 Schema(#[from] SchemaError),
281 #[error(transparent)]
283 Migration(#[from] MigrationCodecError),
284 #[error("schema migration metadata is invalid: {0}")]
286 InvalidMigration(String),
287 #[error("collection {collection_id} is frozen by incompatible schema activations")]
289 SchemaMigrationFork { collection_id: CollectionId },
290 #[error("schema mismatch for collection {collection_id} version {version}")]
292 SchemaMismatch {
293 collection_id: CollectionId,
295 version: u32,
297 },
298 #[error("field {field_id} is not an index in collection {collection_id}")]
300 InvalidIndex {
301 collection_id: CollectionId,
303 field_id: u32,
305 },
306 #[error("derived metadata is corrupt: {0}")]
308 CorruptDerivedMetadata(String),
309 #[error("snapshot metadata is invalid: {0}")]
311 InvalidSnapshotMetadata(String),
312 #[error("domain authority metadata is invalid: {0}")]
314 InvalidAuthority(String),
315 #[error(
317 "domain {domain} is frozen by conflicting control transitions {first:?} and {second:?}"
318 )]
319 ControlFork {
320 domain: DomainId,
322 first: [u8; 32],
324 second: [u8; 32],
326 },
327}
328
329#[derive(Debug, Clone, PartialEq, Eq, minicbor::Encode, minicbor::Decode)]
331#[cbor(array)]
332pub struct SnapshotBaseline {
333 #[n(0)]
334 version: u16,
335 #[n(1)]
336 domain_id: DomainId,
337 #[n(2)]
338 epoch: u64,
339 #[n(3)]
340 manifest_hash: BlobHash,
341 #[n(4)]
342 state_root: [u8; 32],
343 #[n(5)]
344 frontier: Vec<CommitId>,
345 #[n(6)]
346 author_sequences: Vec<(AuthorId, u64)>,
347}
348
349impl SnapshotBaseline {
350 pub fn new(
351 domain_id: DomainId,
352 epoch: u64,
353 manifest_hash: BlobHash,
354 state_root: [u8; 32],
355 frontier: Vec<CommitId>,
356 author_sequences: Vec<(AuthorId, u64)>,
357 ) -> Result<Self, StoreError> {
358 let baseline = Self {
359 version: SNAPSHOT_BASELINE_VERSION,
360 domain_id,
361 epoch,
362 manifest_hash,
363 state_root,
364 frontier,
365 author_sequences,
366 };
367 baseline.validate()?;
368 Ok(baseline)
369 }
370
371 pub const fn domain_id(&self) -> DomainId {
372 self.domain_id
373 }
374 pub const fn epoch(&self) -> u64 {
375 self.epoch
376 }
377 pub const fn manifest_hash(&self) -> BlobHash {
378 self.manifest_hash
379 }
380 pub const fn state_root(&self) -> [u8; 32] {
381 self.state_root
382 }
383 pub fn frontier(&self) -> &[CommitId] {
384 &self.frontier
385 }
386 pub fn author_sequences(&self) -> &[(AuthorId, u64)] {
387 &self.author_sequences
388 }
389
390 fn encode_canonical(&self) -> Result<Vec<u8>, StoreError> {
391 minicbor::to_vec(self).map_err(|error| StoreError::Database(error.to_string()))
392 }
393
394 fn decode_canonical(bytes: &[u8]) -> Result<Self, StoreError> {
395 if bytes.len() > MAX_SNAPSHOT_BASELINE_BYTES {
396 return Err(StoreError::CorruptDerivedMetadata(
397 "snapshot baseline exceeds encoded size limit".into(),
398 ));
399 }
400 let mut decoder = minicbor::Decoder::new(bytes);
401 let baseline: Self = decoder.decode().map_err(|error| {
402 StoreError::CorruptDerivedMetadata(format!("snapshot baseline: {error}"))
403 })?;
404 baseline.validate()?;
405 if decoder.position() != bytes.len() || baseline.encode_canonical()? != bytes {
406 return Err(StoreError::CorruptDerivedMetadata(
407 "snapshot baseline is not canonical".into(),
408 ));
409 }
410 Ok(baseline)
411 }
412
413 fn validate(&self) -> Result<(), StoreError> {
414 if self.version != SNAPSHOT_BASELINE_VERSION
415 || self.frontier.len() > MAX_SNAPSHOT_BASELINE_ENTRIES
416 || self.author_sequences.len() > MAX_SNAPSHOT_BASELINE_ENTRIES
417 || !strictly_sorted(&self.frontier)
418 || !self
419 .author_sequences
420 .windows(2)
421 .all(|pair| pair[0].0 < pair[1].0)
422 || self.frontier.is_empty() != self.author_sequences.is_empty()
423 {
424 return Err(StoreError::InvalidSnapshotMetadata(
425 "snapshot recovery baseline is malformed".into(),
426 ));
427 }
428 Ok(())
429 }
430}
431
432#[derive(Clone)]
434pub struct Store {
435 inner: Arc<StoreInner>,
436}
437
438struct StoreInner {
439 database: Database,
440 blobs: Arc<dyn BlobStore>,
441 control_writer: Mutex<()>,
442 _lock: File,
443}
444
445impl Store {
446 pub fn open(root: impl AsRef<Path>, blobs: Arc<dyn BlobStore>) -> Result<Self, StoreError> {
448 Self::open_with_cache_size(root, blobs, DEFAULT_METADATA_CACHE_SIZE)
449 }
450
451 pub fn open_with_cache_size(
453 root: impl AsRef<Path>,
454 blobs: Arc<dyn BlobStore>,
455 cache_size: usize,
456 ) -> Result<Self, StoreError> {
457 let root = root.as_ref();
458 std::fs::create_dir_all(root).map_err(io_error)?;
459 std::fs::create_dir_all(root.join("locks")).map_err(io_error)?;
460 std::fs::create_dir_all(root.join("tmp")).map_err(io_error)?;
461
462 let lock = OpenOptions::new()
463 .create(true)
464 .read(true)
465 .write(true)
466 .truncate(false)
467 .open(root.join("locks/writer.lock"))
468 .map_err(io_error)?;
469 lock.try_lock_exclusive()
470 .map_err(|_| StoreError::AlreadyOpen)?;
471
472 prepare_store_format(root)?;
473
474 let mut builder = Database::builder();
475 builder.set_cache_size(cache_size);
476 let database = builder
477 .create(root.join("meta-v6.redb"))
478 .map_err(database_error)?;
479 initialize_tables(&database)?;
480 Ok(Self {
481 inner: Arc::new(StoreInner {
482 database,
483 blobs,
484 control_writer: Mutex::new(()),
485 _lock: lock,
486 }),
487 })
488 }
489
490 pub async fn persist_commit(
492 &self,
493 envelope: &CommitEnvelope,
494 ) -> Result<ApplyOutcome, StoreError> {
495 self.persist_commit_with_records(envelope, &[]).await
496 }
497
498 pub async fn persist_commit_with_records(
500 &self,
501 envelope: &CommitEnvelope,
502 records: &[MaterializedRecord],
503 ) -> Result<ApplyOutcome, StoreError> {
504 self.persist_commit_components(envelope, records, &[], None, None)
505 .await
506 }
507
508 pub async fn persist_commit_with_reconciliation(
510 &self,
511 envelope: &CommitEnvelope,
512 records: &[MaterializedRecord],
513 context: &VersionVector,
514 record_heads: &[RecordHeadSet],
515 ) -> Result<ApplyOutcome, StoreError> {
516 self.persist_commit_components(envelope, records, &[], None, Some((context, record_heads)))
517 .await
518 }
519
520 async fn persist_commit_components(
521 &self,
522 envelope: &CommitEnvelope,
523 records: &[MaterializedRecord],
524 stages: &[SchemaStage],
525 activation: Option<&SchemaActivation>,
526 reconciliation: Option<(&VersionVector, &[RecordHeadSet])>,
527 ) -> Result<ApplyOutcome, StoreError> {
528 let commit_id = envelope.commit_id();
529 let previous_author_head = match self.commit_preflight(envelope).await? {
530 CommitPreflight::Duplicate(commit_id) => {
531 return Ok(ApplyOutcome::Duplicate(commit_id));
532 }
533 CommitPreflight::Ready {
534 previous_author_head,
535 } => previous_author_head,
536 };
537 let expected_blob_hash = BlobHash::from_bytes(commit_id.to_bytes());
538 let actual_blob_hash = self.inner.blobs.put(envelope.encode_canonical()).await?;
539 if actual_blob_hash != expected_blob_hash {
540 return Err(StoreError::BlobHashMismatch);
541 }
542
543 self.publish_commit_components(
544 envelope,
545 records,
546 stages,
547 activation,
548 previous_author_head,
549 reconciliation,
550 )
551 }
552
553 pub async fn persist_schema_stages(
555 &self,
556 envelope: &CommitEnvelope,
557 stages: &[SchemaStage],
558 ) -> Result<ApplyOutcome, StoreError> {
559 if stages.is_empty() {
560 return Err(StoreError::InvalidMigration(
561 "a staging commit must contain at least one stage".into(),
562 ));
563 }
564 let commit_id = envelope.commit_id();
565 let previous_author_head = match self.commit_preflight(envelope).await? {
566 CommitPreflight::Duplicate(commit_id) => return Ok(ApplyOutcome::Duplicate(commit_id)),
567 CommitPreflight::Ready {
568 previous_author_head,
569 } => previous_author_head,
570 };
571 let expected = BlobHash::from_bytes(commit_id.to_bytes());
572 if self.inner.blobs.put(envelope.encode_canonical()).await? != expected {
573 return Err(StoreError::BlobHashMismatch);
574 }
575 self.publish_commit_components(envelope, &[], stages, None, previous_author_head, None)
576 }
577
578 pub async fn persist_schema_activation(
580 &self,
581 envelope: &CommitEnvelope,
582 activation: &SchemaActivation,
583 ) -> Result<ApplyOutcome, StoreError> {
584 let commit_id = envelope.commit_id();
585 let previous_author_head = match self.commit_preflight(envelope).await? {
586 CommitPreflight::Duplicate(commit_id) => return Ok(ApplyOutcome::Duplicate(commit_id)),
587 CommitPreflight::Ready {
588 previous_author_head,
589 } => previous_author_head,
590 };
591 let expected = BlobHash::from_bytes(commit_id.to_bytes());
592 if self.inner.blobs.put(envelope.encode_canonical()).await? != expected {
593 return Err(StoreError::BlobHashMismatch);
594 }
595 self.publish_commit_components(
596 envelope,
597 &[],
598 &[],
599 Some(activation),
600 previous_author_head,
601 None,
602 )
603 }
604
605 pub async fn persist_schema_stages_with_context(
607 &self,
608 envelope: &CommitEnvelope,
609 stages: &[SchemaStage],
610 context: &VersionVector,
611 ) -> Result<ApplyOutcome, StoreError> {
612 if stages.is_empty() {
613 return Err(StoreError::InvalidMigration(
614 "a staging commit must contain at least one stage".into(),
615 ));
616 }
617 self.persist_commit_components(envelope, &[], stages, None, Some((context, &[])))
618 .await
619 }
620
621 pub async fn persist_schema_activation_with_context(
623 &self,
624 envelope: &CommitEnvelope,
625 activation: &SchemaActivation,
626 context: &VersionVector,
627 ) -> Result<ApplyOutcome, StoreError> {
628 self.persist_commit_components(envelope, &[], &[], Some(activation), Some((context, &[])))
629 .await
630 }
631
632 #[allow(clippy::too_many_lines)]
633 fn publish_commit_components(
634 &self,
635 envelope: &CommitEnvelope,
636 records: &[MaterializedRecord],
637 stages: &[SchemaStage],
638 activation: Option<&SchemaActivation>,
639 previous_author_head: Option<CommitId>,
640 reconciliation: Option<(&VersionVector, &[RecordHeadSet])>,
641 ) -> Result<ApplyOutcome, StoreError> {
642 let commit_id = envelope.commit_id();
643 let write = self.inner.database.begin_write().map_err(database_error)?;
644 let mut migration_fork: Option<CollectionId>;
645 {
646 let mut commits = write.open_table(COMMIT_METADATA).map_err(database_error)?;
647 let coverage = write
648 .open_table(SNAPSHOT_COVERAGE)
649 .map_err(database_error)?;
650 let quarantined = write
651 .open_table(QUARANTINED_COMMITS)
652 .map_err(database_error)?;
653 if commits
654 .get(commit_id.as_bytes().as_slice())
655 .map_err(database_error)?
656 .is_some()
657 {
658 return Ok(ApplyOutcome::Duplicate(commit_id));
659 }
660
661 for dependency in envelope.header().dependencies() {
662 if !dependency_is_accepted(
663 &commits,
664 &coverage,
665 &quarantined,
666 envelope.header().domain_id(),
667 *dependency,
668 )? {
669 return Err(StoreError::MissingDependency(*dependency));
670 }
671 }
672
673 let author_key =
674 author_sequence_key(envelope.header().domain_id(), envelope.header().author());
675 let mut sequences = write.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
676 let last_sequence = sequences
677 .get(author_key.as_slice())
678 .map_err(database_error)?
679 .map(|value| value.value());
680 let expected_sequence = match last_sequence {
681 Some(sequence) => sequence.checked_add(1).ok_or(StoreError::SequenceGap {
682 expected: sequence,
683 actual: envelope.header().author_sequence(),
684 })?,
685 None => 0,
686 };
687 if envelope.header().author_sequence() != expected_sequence {
688 return Err(StoreError::SequenceGap {
689 expected: expected_sequence,
690 actual: envelope.header().author_sequence(),
691 });
692 }
693
694 let mut author_heads = write.open_table(AUTHOR_HEADS).map_err(database_error)?;
695 let current_author_head = author_heads
696 .get(author_key.as_slice())
697 .map_err(database_error)?
698 .map(|value| {
699 let bytes: [u8; 32] = value.value().try_into().map_err(|_| {
700 StoreError::CorruptDerivedMetadata(
701 "author head commit ID is invalid".into(),
702 )
703 })?;
704 Ok::<CommitId, StoreError>(CommitId::from_bytes(bytes))
705 })
706 .transpose()?;
707 if current_author_head != previous_author_head {
708 return Err(StoreError::SequenceGap {
709 expected: expected_sequence,
710 actual: envelope.header().author_sequence(),
711 });
712 }
713
714 let metadata = encode_metadata(envelope);
715 commits
716 .insert(commit_id.as_bytes().as_slice(), metadata.as_slice())
717 .map_err(database_error)?;
718 sequences
719 .insert(author_key.as_slice(), envelope.header().author_sequence())
720 .map_err(database_error)?;
721 write
722 .open_table(AUTHOR_COMMITS)
723 .map_err(database_error)?
724 .insert(
725 author_commit_key(
726 envelope.header().domain_id(),
727 envelope.header().author(),
728 envelope.header().author_sequence(),
729 )
730 .as_slice(),
731 commit_id.as_bytes().as_slice(),
732 )
733 .map_err(database_error)?;
734 author_heads
735 .insert(author_key.as_slice(), commit_id.as_bytes().as_slice())
736 .map_err(database_error)?;
737
738 publish_frontier(&write, envelope.header().domain_id(), envelope, commit_id)?;
739
740 apply_materialized_records(&write, envelope.header().domain_id(), records)?;
741 match reconciliation {
742 Some((context, record_heads)) => {
743 apply_reconciliation_update(&write, envelope, context, record_heads)?;
744 }
745 None => mark_reconciliation_stale(&write, envelope.header().domain_id())?,
746 }
747 migration_fork = apply_schema_stages(&write, envelope.header().domain_id(), stages)?;
748 if migration_fork.is_none()
749 && let Some(activation) = activation
750 && apply_schema_activation(&write, envelope.header().domain_id(), activation)?
751 {
752 migration_fork = Some(activation.to().collection_id());
753 }
754 }
755 write.commit().map_err(database_error)?;
756 if let Some(collection_id) = migration_fork {
757 return Err(StoreError::SchemaMigrationFork { collection_id });
758 }
759 Ok(ApplyOutcome::Applied(commit_id))
760 }
761
762 pub fn active_schema_activation(
764 &self,
765 domain_id: DomainId,
766 collection_id: CollectionId,
767 ) -> Result<Option<SchemaActivation>, StoreError> {
768 if self.migration_frozen(domain_id, collection_id)? {
769 return Err(StoreError::SchemaMigrationFork { collection_id });
770 }
771 let read = self.inner.database.begin_read().map_err(database_error)?;
772 let table = read.open_table(ACTIVE_SCHEMAS).map_err(database_error)?;
773 table
774 .get(active_schema_key(domain_id, collection_id).as_slice())
775 .map_err(database_error)?
776 .map(|value| {
777 SchemaActivation::decode_canonical(value.value()).map_err(StoreError::from)
778 })
779 .transpose()
780 }
781
782 pub fn active_schema_activations(
784 &self,
785 domain_id: DomainId,
786 ) -> Result<Vec<SchemaActivation>, StoreError> {
787 let read = self.inner.database.begin_read().map_err(database_error)?;
788 let forks = read.open_table(MIGRATION_FORKS).map_err(database_error)?;
789 if !table_keys_with_prefix(&forks, domain_id.as_bytes())?.is_empty() {
790 return Err(StoreError::InvalidMigration(
791 "domain contains a frozen schema activation fork".into(),
792 ));
793 }
794 let table = read.open_table(ACTIVE_SCHEMAS).map_err(database_error)?;
795 let mut activations = Vec::new();
796 for entry in table.iter().map_err(database_error)? {
797 let (key, value) = entry.map_err(database_error)?;
798 if key.value().starts_with(domain_id.as_bytes()) {
799 activations.push(SchemaActivation::decode_canonical(value.value())?);
800 }
801 }
802 activations.sort_unstable_by_key(|activation| activation.to().collection_id());
803 Ok(activations)
804 }
805
806 pub fn migration_frozen(
808 &self,
809 domain_id: DomainId,
810 collection_id: CollectionId,
811 ) -> Result<bool, StoreError> {
812 let read = self.inner.database.begin_read().map_err(database_error)?;
813 let table = read.open_table(MIGRATION_FORKS).map_err(database_error)?;
814 Ok(table
815 .get(active_schema_key(domain_id, collection_id).as_slice())
816 .map_err(database_error)?
817 .is_some())
818 }
819
820 pub fn migration_counts(
822 &self,
823 domain_id: DomainId,
824 ) -> Result<(usize, usize, usize), StoreError> {
825 let read = self.inner.database.begin_read().map_err(database_error)?;
826 let active = read.open_table(ACTIVE_SCHEMAS).map_err(database_error)?;
827 let stages = read.open_table(MIGRATION_STAGES).map_err(database_error)?;
828 let forks = read.open_table(MIGRATION_FORKS).map_err(database_error)?;
829 Ok((
830 table_keys_with_prefix(&active, domain_id.as_bytes())?.len(),
831 table_keys_with_prefix(&stages, domain_id.as_bytes())?.len(),
832 table_keys_with_prefix(&forks, domain_id.as_bytes())?.len(),
833 ))
834 }
835
836 pub fn record_snapshot_observed_at(
838 &self,
839 domain_id: DomainId,
840 unix_seconds: u64,
841 ) -> Result<(), StoreError> {
842 let write = self.inner.database.begin_write().map_err(database_error)?;
843 {
844 let mut table = write
845 .open_table(OPERATIONAL_METADATA)
846 .map_err(database_error)?;
847 table
848 .insert(
849 operational_metadata_key(domain_id, SNAPSHOT_OBSERVED_AT_TAG).as_slice(),
850 unix_seconds,
851 )
852 .map_err(database_error)?;
853 }
854 write.commit().map_err(database_error)
855 }
856
857 pub fn snapshot_observed_at(&self, domain_id: DomainId) -> Result<Option<u64>, StoreError> {
859 let read = self.inner.database.begin_read().map_err(database_error)?;
860 let table = read
861 .open_table(OPERATIONAL_METADATA)
862 .map_err(database_error)?;
863 Ok(table
864 .get(operational_metadata_key(domain_id, SNAPSHOT_OBSERVED_AT_TAG).as_slice())
865 .map_err(database_error)?
866 .map(|value| value.value()))
867 }
868
869 pub fn schema_stages(
871 &self,
872 domain_id: DomainId,
873 migration_id: MigrationId,
874 ) -> Result<Vec<SchemaStage>, StoreError> {
875 let read = self.inner.database.begin_read().map_err(database_error)?;
876 let table = read.open_table(MIGRATION_STAGES).map_err(database_error)?;
877 load_schema_stages(&table, domain_id, migration_id)
878 }
879
880 async fn commit_preflight(
881 &self,
882 envelope: &CommitEnvelope,
883 ) -> Result<CommitPreflight, StoreError> {
884 let commit_id = envelope.commit_id();
885 if self.is_commit_quarantined(commit_id)? {
886 return Err(StoreError::QuarantinedCommit(commit_id));
887 }
888 if let Some(from_sequence) =
889 self.author_quarantine(envelope.header().domain_id(), envelope.header().author())?
890 && envelope.header().author_sequence() >= from_sequence
891 {
892 return Err(StoreError::QuarantinedAuthor {
893 author: envelope.header().author(),
894 from_sequence,
895 });
896 }
897 if self.contains_commit(envelope.header().domain_id(), commit_id)? {
898 return Ok(CommitPreflight::Duplicate(commit_id));
899 }
900 if let Some(existing) = self.accepted_author_commit(
901 envelope.header().domain_id(),
902 envelope.header().author(),
903 envelope.header().author_sequence(),
904 )? && existing != commit_id
905 {
906 let (first, second) = if existing < commit_id {
907 (existing, commit_id)
908 } else {
909 (commit_id, existing)
910 };
911 return Err(StoreError::AuthorEquivocation {
912 author: envelope.header().author(),
913 sequence: envelope.header().author_sequence(),
914 first,
915 second,
916 });
917 }
918 let previous_author_head =
919 self.author_head(envelope.header().domain_id(), envelope.header().author())?;
920 if let Some(previous) = previous_author_head
921 && !self
922 .dependencies_include_ancestor(
923 envelope.header().domain_id(),
924 envelope.header().dependencies(),
925 previous,
926 )
927 .await?
928 {
929 return Err(StoreError::AuthorChainMissing { previous });
930 }
931 if let Some(baseline) = self.snapshot_baseline(envelope.header().domain_id())?
932 && baseline
933 .author_sequences()
934 .iter()
935 .any(|(author, _)| *author == envelope.header().author())
936 {
937 for &frontier in baseline.frontier() {
938 if !self
939 .dependencies_include_ancestor(
940 envelope.header().domain_id(),
941 envelope.header().dependencies(),
942 frontier,
943 )
944 .await?
945 {
946 return Err(StoreError::SnapshotBaselineMissing { frontier });
947 }
948 }
949 }
950 Ok(CommitPreflight::Ready {
951 previous_author_head,
952 })
953 }
954
955 pub fn contains_commit(
957 &self,
958 domain_id: DomainId,
959 commit_id: CommitId,
960 ) -> Result<bool, StoreError> {
961 let read = self.inner.database.begin_read().map_err(database_error)?;
962 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
963 let coverage = read.open_table(SNAPSHOT_COVERAGE).map_err(database_error)?;
964 let quarantined = read
965 .open_table(QUARANTINED_COMMITS)
966 .map_err(database_error)?;
967 let stored = commits
968 .get(commit_id.as_bytes().as_slice())
969 .map_err(database_error)?
970 .is_some_and(|metadata| metadata.value().starts_with(domain_id.as_bytes()))
971 && quarantined
972 .get(commit_id.as_bytes().as_slice())
973 .map_err(database_error)?
974 .is_none();
975 let snapshotted = coverage
976 .get(commit_id.as_bytes().as_slice())
977 .map_err(database_error)?
978 .is_some_and(|covered_domain| covered_domain.value() == domain_id.as_bytes());
979 Ok(stored || snapshotted)
980 }
981
982 pub fn is_commit_quarantined(&self, commit_id: CommitId) -> Result<bool, StoreError> {
984 let read = self.inner.database.begin_read().map_err(database_error)?;
985 let quarantined = read
986 .open_table(QUARANTINED_COMMITS)
987 .map_err(database_error)?;
988 Ok(quarantined
989 .get(commit_id.as_bytes().as_slice())
990 .map_err(database_error)?
991 .is_some())
992 }
993
994 pub fn author_quarantine(
996 &self,
997 domain_id: DomainId,
998 author: AuthorId,
999 ) -> Result<Option<u64>, StoreError> {
1000 let read = self.inner.database.begin_read().map_err(database_error)?;
1001 let quarantined = read
1002 .open_table(QUARANTINED_AUTHORS)
1003 .map_err(database_error)?;
1004 let key = author_sequence_key(domain_id, author);
1005 Ok(quarantined
1006 .get(key.as_slice())
1007 .map_err(database_error)?
1008 .map(|value| value.value()))
1009 }
1010
1011 pub fn last_sequence(
1013 &self,
1014 domain_id: DomainId,
1015 author: AuthorId,
1016 ) -> Result<Option<u64>, StoreError> {
1017 let read = self.inner.database.begin_read().map_err(database_error)?;
1018 let sequences = read.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
1019 let key = author_sequence_key(domain_id, author);
1020 Ok(sequences
1021 .get(key.as_slice())
1022 .map_err(database_error)?
1023 .map(|value| value.value()))
1024 }
1025
1026 pub fn author_head(
1028 &self,
1029 domain_id: DomainId,
1030 author: AuthorId,
1031 ) -> Result<Option<CommitId>, StoreError> {
1032 let read = self.inner.database.begin_read().map_err(database_error)?;
1033 let heads = read.open_table(AUTHOR_HEADS).map_err(database_error)?;
1034 let key = author_sequence_key(domain_id, author);
1035 heads
1036 .get(key.as_slice())
1037 .map_err(database_error)?
1038 .map(|value| {
1039 let bytes: [u8; 32] = value.value().try_into().map_err(|_| {
1040 StoreError::CorruptDerivedMetadata("author head commit ID is invalid".into())
1041 })?;
1042 Ok(CommitId::from_bytes(bytes))
1043 })
1044 .transpose()
1045 }
1046
1047 pub fn conflicting_author_commit(
1049 &self,
1050 envelope: &CommitEnvelope,
1051 ) -> Result<Option<CommitId>, StoreError> {
1052 let commit_id = envelope.commit_id();
1053 Ok(self
1054 .accepted_author_commit(
1055 envelope.header().domain_id(),
1056 envelope.header().author(),
1057 envelope.header().author_sequence(),
1058 )?
1059 .filter(|existing| *existing != commit_id))
1060 }
1061
1062 pub async fn quarantine_equivocation(
1064 &self,
1065 evidence: &CommitEnvelope,
1066 invalid_commits: &[CommitId],
1067 accepted_frontier: &[CommitId],
1068 author_sequences: &[(AuthorId, u64)],
1069 author_heads: &[(AuthorId, u64, CommitId)],
1070 records: &[MaterializedRecord],
1071 ) -> Result<(), StoreError> {
1072 let evidence_id = evidence.commit_id();
1073 let expected_hash = BlobHash::from_bytes(evidence_id.to_bytes());
1074 if self.inner.blobs.put(evidence.encode_canonical()).await? != expected_hash {
1075 return Err(StoreError::BlobHashMismatch);
1076 }
1077 let domain_id = evidence.header().domain_id();
1078 let write = self.inner.database.begin_write().map_err(database_error)?;
1079 clear_reconciliation_domain(&write, domain_id)?;
1080 clear_materialized_domain(&write, domain_id)?;
1081 apply_materialized_records(&write, domain_id, records)?;
1082 {
1083 let mut quarantined = write
1084 .open_table(QUARANTINED_COMMITS)
1085 .map_err(database_error)?;
1086 for commit_id in invalid_commits
1087 .iter()
1088 .copied()
1089 .chain(std::iter::once(evidence_id))
1090 {
1091 quarantined
1092 .insert(
1093 commit_id.as_bytes().as_slice(),
1094 domain_id.as_bytes().as_slice(),
1095 )
1096 .map_err(database_error)?;
1097 }
1098 }
1099 {
1100 let mut authors = write
1101 .open_table(QUARANTINED_AUTHORS)
1102 .map_err(database_error)?;
1103 let key = author_sequence_key(domain_id, evidence.header().author());
1104 let existing = authors
1105 .get(key.as_slice())
1106 .map_err(database_error)?
1107 .map(|value| value.value());
1108 let from_sequence = existing.map_or(evidence.header().author_sequence(), |current| {
1109 current.min(evidence.header().author_sequence())
1110 });
1111 authors
1112 .insert(key.as_slice(), from_sequence)
1113 .map_err(database_error)?;
1114 }
1115 {
1116 let mut sequences = write.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
1117 let keys = table_u64_keys_with_prefix(&sequences, domain_id.as_bytes())?;
1118 for key in keys {
1119 sequences.remove(key.as_slice()).map_err(database_error)?;
1120 }
1121 let mut heads = write.open_table(AUTHOR_HEADS).map_err(database_error)?;
1122 let keys = table_keys_with_prefix(&heads, domain_id.as_bytes())?;
1123 for key in keys {
1124 heads.remove(key.as_slice()).map_err(database_error)?;
1125 }
1126 for (author, sequence) in author_sequences {
1127 let key = author_sequence_key(domain_id, *author);
1128 sequences
1129 .insert(key.as_slice(), *sequence)
1130 .map_err(database_error)?;
1131 }
1132 for (author, _, commit_id) in author_heads {
1133 let key = author_sequence_key(domain_id, *author);
1134 heads
1135 .insert(key.as_slice(), commit_id.as_bytes().as_slice())
1136 .map_err(database_error)?;
1137 }
1138 }
1139 {
1140 let mut frontier = write.open_table(FRONTIER).map_err(database_error)?;
1141 let keys = table_u8_keys_with_prefix(&frontier, domain_id.as_bytes())?;
1142 for key in keys {
1143 frontier.remove(key.as_slice()).map_err(database_error)?;
1144 }
1145 for commit_id in accepted_frontier {
1146 let key = frontier_key(domain_id, *commit_id);
1147 frontier.insert(key.as_slice(), 0).map_err(database_error)?;
1148 }
1149 }
1150 write.commit().map_err(database_error)
1151 }
1152
1153 fn accepted_author_commit(
1154 &self,
1155 domain_id: DomainId,
1156 author: AuthorId,
1157 sequence: u64,
1158 ) -> Result<Option<CommitId>, StoreError> {
1159 if self
1160 .last_sequence(domain_id, author)?
1161 .is_none_or(|last| sequence > last)
1162 {
1163 return Ok(None);
1164 }
1165 let read = self.inner.database.begin_read().map_err(database_error)?;
1166 let author_commits = read.open_table(AUTHOR_COMMITS).map_err(database_error)?;
1167 let author_key = author_commit_key(domain_id, author, sequence);
1168 if let Some(value) = author_commits
1169 .get(author_key.as_slice())
1170 .map_err(database_error)?
1171 {
1172 let bytes: [u8; 32] = value.value().try_into().map_err(|_| {
1173 StoreError::CorruptDerivedMetadata("author commit ID is invalid".into())
1174 })?;
1175 let commit_id = CommitId::from_bytes(bytes);
1176 return if self.is_commit_quarantined(commit_id)? {
1177 Ok(None)
1178 } else {
1179 Ok(Some(commit_id))
1180 };
1181 }
1182 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1183 let quarantined = read
1184 .open_table(QUARANTINED_COMMITS)
1185 .map_err(database_error)?;
1186 for entry in commits.iter().map_err(database_error)? {
1187 let (key, value) = entry.map_err(database_error)?;
1188 if quarantined
1189 .get(key.value())
1190 .map_err(database_error)?
1191 .is_some()
1192 {
1193 continue;
1194 }
1195 let metadata = value.value();
1196 if metadata.len() != 112 {
1197 return Err(StoreError::CorruptDerivedMetadata(
1198 "commit metadata length is invalid".into(),
1199 ));
1200 }
1201 if &metadata[..32] != domain_id.as_bytes()
1202 || &metadata[32..64] != author.as_bytes()
1203 || metadata[72..80] != sequence.to_be_bytes()
1204 {
1205 continue;
1206 }
1207 let bytes: [u8; 32] = key.value().try_into().map_err(|_| {
1208 StoreError::CorruptDerivedMetadata("commit ID length is invalid".into())
1209 })?;
1210 return Ok(Some(CommitId::from_bytes(bytes)));
1211 }
1212 Ok(None)
1213 }
1214
1215 async fn dependencies_include_ancestor(
1216 &self,
1217 domain_id: DomainId,
1218 dependencies: &[CommitId],
1219 ancestor: CommitId,
1220 ) -> Result<bool, StoreError> {
1221 if self.reconciliation_index_ready(domain_id)? {
1222 if dependencies.contains(&ancestor) {
1223 return Ok(true);
1224 }
1225 let read = self.inner.database.begin_read().map_err(database_error)?;
1226 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1227 let Some(metadata) = commits
1228 .get(ancestor.as_bytes().as_slice())
1229 .map_err(database_error)?
1230 else {
1231 return Ok(false);
1232 };
1233 if metadata.value().len() != 112 || !metadata.value().starts_with(domain_id.as_bytes())
1234 {
1235 return Err(StoreError::CorruptDerivedMetadata(
1236 "ancestor commit metadata is invalid".into(),
1237 ));
1238 }
1239 let author: [u8; 32] = metadata.value()[32..64].try_into().map_err(|_| {
1240 StoreError::CorruptDerivedMetadata("ancestor author is invalid".into())
1241 })?;
1242 let sequence =
1243 u64::from_be_bytes(metadata.value()[72..80].try_into().map_err(|_| {
1244 StoreError::CorruptDerivedMetadata("ancestor sequence is invalid".into())
1245 })?);
1246 let dot = iroh_db_core::Dot::new(AuthorId::from_bytes(author), sequence, 0);
1247 drop(metadata);
1248 drop(commits);
1249 drop(read);
1250 let baseline = self.snapshot_baseline(domain_id)?;
1251 for dependency in dependencies {
1252 if self
1253 .commit_context(domain_id, *dependency)?
1254 .is_some_and(|context| context.covers(&dot))
1255 {
1256 return Ok(true);
1257 }
1258 if self.contains_commit(domain_id, *dependency)?
1259 && baseline.as_ref().is_some_and(|baseline| {
1260 baseline
1261 .author_sequences()
1262 .iter()
1263 .any(|(author, sequence)| {
1264 *author == dot.author() && *sequence >= dot.sequence()
1265 })
1266 })
1267 {
1268 return Ok(true);
1269 }
1270 }
1271 return Ok(false);
1272 }
1273 let mut pending = dependencies.to_vec();
1274 let mut visited = std::collections::BTreeSet::new();
1275 while let Some(commit_id) = pending.pop() {
1276 if commit_id == ancestor {
1277 return Ok(true);
1278 }
1279 if !visited.insert(commit_id) {
1280 continue;
1281 }
1282 match self.load_commit(commit_id).await {
1283 Ok(envelope) if envelope.header().domain_id() == domain_id => {
1284 pending.extend_from_slice(envelope.header().dependencies());
1285 }
1286 Ok(_) => return Ok(false),
1287 Err(StoreError::MissingCommitBlob(_))
1288 if self.contains_commit(domain_id, commit_id)? =>
1289 {
1290 }
1292 Err(error) => return Err(error),
1293 }
1294 }
1295 Ok(false)
1296 }
1297
1298 pub fn frontier(&self, domain_id: DomainId) -> Result<Vec<CommitId>, StoreError> {
1300 let read = self.inner.database.begin_read().map_err(database_error)?;
1301 let frontier = read.open_table(FRONTIER).map_err(database_error)?;
1302 let mut commits = Vec::new();
1303 for entry in frontier.iter().map_err(database_error)? {
1304 let (key, _) = entry.map_err(database_error)?;
1305 let key = key.value();
1306 if key.starts_with(domain_id.as_bytes()) {
1307 let bytes: [u8; 32] = key[32..]
1308 .try_into()
1309 .map_err(|_| StoreError::Database("invalid frontier key".into()))?;
1310 commits.push(CommitId::from_bytes(bytes));
1311 }
1312 }
1313 commits.sort_unstable();
1314 Ok(commits)
1315 }
1316
1317 pub fn reconciliation_index_ready(&self, domain_id: DomainId) -> Result<bool, StoreError> {
1319 let read = self.inner.database.begin_read().map_err(database_error)?;
1320 let ready = read
1321 .open_table(RECONCILIATION_READY)
1322 .map_err(database_error)?;
1323 Ok(ready
1324 .get(domain_id.as_bytes().as_slice())
1325 .map_err(database_error)?
1326 .is_some())
1327 }
1328
1329 pub fn commit_context(
1331 &self,
1332 domain_id: DomainId,
1333 commit_id: CommitId,
1334 ) -> Result<Option<VersionVector>, StoreError> {
1335 let read = self.inner.database.begin_read().map_err(database_error)?;
1336 let contexts = read.open_table(COMMIT_CONTEXTS).map_err(database_error)?;
1337 let key = commit_context_key(domain_id, commit_id);
1338 contexts
1339 .get(key.as_slice())
1340 .map_err(database_error)?
1341 .map(|value| decode_version_vector(value.value()))
1342 .transpose()
1343 }
1344
1345 pub fn record_heads(
1347 &self,
1348 domain_id: DomainId,
1349 collection_id: CollectionId,
1350 record_id: &[u8],
1351 ) -> Result<Vec<RecordHead>, StoreError> {
1352 let read = self.inner.database.begin_read().map_err(database_error)?;
1353 let heads = read.open_table(RECORD_HEADS).map_err(database_error)?;
1354 let prefix = record_head_prefix(domain_id, collection_id, record_id)?;
1355 let mut result = Vec::new();
1356 for entry in heads.range(prefix.as_slice()..).map_err(database_error)? {
1357 let (key, value) = entry.map_err(database_error)?;
1358 if !key.value().starts_with(&prefix) {
1359 break;
1360 }
1361 if key.value().len() != prefix.len() + 32 {
1362 return Err(StoreError::CorruptDerivedMetadata(
1363 "record head key has an invalid commit ID".into(),
1364 ));
1365 }
1366 let commit: [u8; 32] = key.value()[prefix.len()..].try_into().map_err(|_| {
1367 StoreError::CorruptDerivedMetadata(
1368 "record head commit ID has an invalid length".into(),
1369 )
1370 })?;
1371 result.push(decode_record_head(
1372 CommitId::from_bytes(commit),
1373 value.value(),
1374 )?);
1375 }
1376 result.sort_unstable_by_key(RecordHead::commit_id);
1377 Ok(result)
1378 }
1379
1380 pub fn replace_reconciliation_index(
1382 &self,
1383 domain_id: DomainId,
1384 contexts: &[(CommitId, VersionVector)],
1385 heads: &[RecordHead],
1386 ) -> Result<(), StoreError> {
1387 let write = self.inner.database.begin_write().map_err(database_error)?;
1388 clear_reconciliation_domain(&write, domain_id)?;
1389 {
1390 let mut table = write.open_table(COMMIT_CONTEXTS).map_err(database_error)?;
1391 for (commit_id, context) in contexts {
1392 let key = commit_context_key(domain_id, *commit_id);
1393 let encoded = encode_version_vector(context)?;
1394 table
1395 .insert(key.as_slice(), encoded.as_slice())
1396 .map_err(database_error)?;
1397 }
1398 }
1399 {
1400 let mut table = write.open_table(RECORD_HEADS).map_err(database_error)?;
1401 for head in heads {
1402 let key = record_head_key(domain_id, head)?;
1403 let value = encode_record_head(head)?;
1404 table
1405 .insert(key.as_slice(), value.as_slice())
1406 .map_err(database_error)?;
1407 }
1408 }
1409 write
1410 .open_table(RECONCILIATION_READY)
1411 .map_err(database_error)?
1412 .insert(domain_id.as_bytes().as_slice(), 1)
1413 .map_err(database_error)?;
1414 write.commit().map_err(database_error)
1415 }
1416
1417 pub fn get_record(
1419 &self,
1420 domain_id: DomainId,
1421 collection_id: CollectionId,
1422 record_id: &[u8],
1423 ) -> Result<Option<Vec<u8>>, StoreError> {
1424 let read = self.inner.database.begin_read().map_err(database_error)?;
1425 let records = read.open_table(RECORDS).map_err(database_error)?;
1426 let key = record_key(domain_id, collection_id, record_id);
1427 Ok(records
1428 .get(key.as_slice())
1429 .map_err(database_error)?
1430 .map(|value| value.value().to_vec()))
1431 }
1432
1433 pub fn list_records(
1435 &self,
1436 domain_id: DomainId,
1437 collection_id: CollectionId,
1438 ) -> Result<Vec<StoredRecord>, StoreError> {
1439 let read = self.inner.database.begin_read().map_err(database_error)?;
1440 let records = read.open_table(RECORDS).map_err(database_error)?;
1441 let prefix = record_prefix(domain_id, collection_id);
1442 let mut found = Vec::new();
1443 for entry in records.iter().map_err(database_error)? {
1444 let (key, value) = entry.map_err(database_error)?;
1445 let key = key.value();
1446 if key.starts_with(&prefix) {
1447 found.push(StoredRecord {
1448 record_id: key[prefix.len()..].to_vec(),
1449 state: value.value().to_vec(),
1450 });
1451 }
1452 }
1453 Ok(found)
1454 }
1455
1456 pub fn lookup_equal(
1458 &self,
1459 domain_id: DomainId,
1460 collection_id: CollectionId,
1461 field_id: u32,
1462 value: &[u8],
1463 ) -> Result<Vec<Vec<u8>>, StoreError> {
1464 let read = self.inner.database.begin_read().map_err(database_error)?;
1465 let index = read.open_table(SECONDARY_INDEX).map_err(database_error)?;
1466 let prefix = index_value_prefix(domain_id, collection_id, field_id, value)?;
1467 let mut record_ids = Vec::new();
1468 for entry in index.iter().map_err(database_error)? {
1469 let (key, _) = entry.map_err(database_error)?;
1470 let key = key.value();
1471 if !key.starts_with(&prefix) {
1472 continue;
1473 }
1474 let suffix = &key[prefix.len()..];
1475 if suffix.len() < 4 {
1476 return Err(StoreError::CorruptDerivedMetadata(
1477 "index record length is missing".into(),
1478 ));
1479 }
1480 let length = read_u32(&suffix[..4])?;
1481 if suffix.len() != 4 + length {
1482 return Err(StoreError::CorruptDerivedMetadata(
1483 "index record length is invalid".into(),
1484 ));
1485 }
1486 record_ids.push(suffix[4..].to_vec());
1487 }
1488 record_ids.sort();
1489 Ok(record_ids)
1490 }
1491
1492 pub fn registered_schema(
1494 &self,
1495 domain_id: DomainId,
1496 collection_id: CollectionId,
1497 version: u32,
1498 ) -> Result<Option<SchemaDescriptor>, StoreError> {
1499 let read = self.inner.database.begin_read().map_err(database_error)?;
1500 let schemas = read.open_table(SCHEMAS).map_err(database_error)?;
1501 let key = schema_key(domain_id, collection_id, version);
1502 schemas
1503 .get(key.as_slice())
1504 .map_err(database_error)?
1505 .map(|value| {
1506 SchemaDescriptor::decode_canonical(value.value()).map_err(StoreError::from)
1507 })
1508 .transpose()
1509 }
1510
1511 pub fn latest_registered_schema(
1513 &self,
1514 domain_id: DomainId,
1515 collection_id: CollectionId,
1516 ) -> Result<Option<SchemaDescriptor>, StoreError> {
1517 let read = self.inner.database.begin_read().map_err(database_error)?;
1518 let schemas = read.open_table(SCHEMAS).map_err(database_error)?;
1519 let mut prefix = [0_u8; 64];
1520 prefix[..32].copy_from_slice(domain_id.as_bytes());
1521 prefix[32..].copy_from_slice(collection_id.as_bytes());
1522 let mut latest: Option<SchemaDescriptor> = None;
1523 for entry in schemas.iter().map_err(database_error)? {
1524 let (key, value) = entry.map_err(database_error)?;
1525 if !key.value().starts_with(&prefix) {
1526 continue;
1527 }
1528 let schema = SchemaDescriptor::decode_canonical(value.value())?;
1529 if latest
1530 .as_ref()
1531 .is_none_or(|current| schema.version() > current.version())
1532 {
1533 latest = Some(schema);
1534 }
1535 }
1536 Ok(latest)
1537 }
1538
1539 pub fn list_commit_ids(&self, domain_id: DomainId) -> Result<Vec<CommitId>, StoreError> {
1541 let mut ids = self.list_stored_commit_ids(domain_id)?;
1542 let read = self.inner.database.begin_read().map_err(database_error)?;
1543 let coverage = read.open_table(SNAPSHOT_COVERAGE).map_err(database_error)?;
1544 for entry in coverage.iter().map_err(database_error)? {
1545 let (key, value) = entry.map_err(database_error)?;
1546 if value.value() == domain_id.as_bytes() {
1547 let bytes: [u8; 32] = key.value().try_into().map_err(|_| {
1548 StoreError::CorruptDerivedMetadata("coverage commit ID is invalid".into())
1549 })?;
1550 ids.push(CommitId::from_bytes(bytes));
1551 }
1552 }
1553 ids.sort_unstable();
1554 ids.dedup();
1555 Ok(ids)
1556 }
1557
1558 pub fn has_snapshot_coverage(&self, domain_id: DomainId) -> Result<bool, StoreError> {
1560 Ok(!self.snapshot_coverage(domain_id)?.is_empty())
1561 }
1562
1563 pub fn snapshot_coverage(&self, domain_id: DomainId) -> Result<Vec<CommitId>, StoreError> {
1565 let read = self.inner.database.begin_read().map_err(database_error)?;
1566 let coverage = read.open_table(SNAPSHOT_COVERAGE).map_err(database_error)?;
1567 let mut commits = Vec::new();
1568 for entry in coverage.iter().map_err(database_error)? {
1569 let (key, value) = entry.map_err(database_error)?;
1570 if value.value() == domain_id.as_bytes() {
1571 let commit_id = key.value().try_into().map_err(|_| {
1572 StoreError::CorruptDerivedMetadata("coverage commit ID is invalid".into())
1573 })?;
1574 commits.push(CommitId::from_bytes(commit_id));
1575 }
1576 }
1577 commits.sort_unstable();
1578 if commits.windows(2).any(|pair| pair[0] == pair[1]) {
1579 return Err(StoreError::CorruptDerivedMetadata(
1580 "snapshot coverage contains duplicate commits".into(),
1581 ));
1582 }
1583 Ok(commits)
1584 }
1585
1586 pub fn snapshot_baseline(
1588 &self,
1589 domain_id: DomainId,
1590 ) -> Result<Option<SnapshotBaseline>, StoreError> {
1591 let read = self.inner.database.begin_read().map_err(database_error)?;
1592 let baselines = read
1593 .open_table(SNAPSHOT_BASELINES)
1594 .map_err(database_error)?;
1595 let baseline = baselines
1596 .get(domain_id.as_bytes().as_slice())
1597 .map_err(database_error)?
1598 .map(|bytes| SnapshotBaseline::decode_canonical(bytes.value()))
1599 .transpose()?;
1600 if baseline
1601 .as_ref()
1602 .is_some_and(|baseline| baseline.domain_id() != domain_id)
1603 {
1604 return Err(StoreError::CorruptDerivedMetadata(
1605 "snapshot baseline crosses domains".into(),
1606 ));
1607 }
1608 Ok(baseline)
1609 }
1610
1611 pub fn claim_blob(&self, domain_id: DomainId, hash: BlobHash) -> Result<(), StoreError> {
1613 let write = self.inner.database.begin_write().map_err(database_error)?;
1614 write
1615 .open_table(BLOB_OWNERSHIP)
1616 .map_err(database_error)?
1617 .insert(blob_ownership_key(domain_id, hash).as_slice(), 0)
1618 .map_err(database_error)?;
1619 write.commit().map_err(database_error)
1620 }
1621
1622 pub fn owns_blob(&self, domain_id: DomainId, hash: BlobHash) -> Result<bool, StoreError> {
1624 let read = self.inner.database.begin_read().map_err(database_error)?;
1625 let ownership = read.open_table(BLOB_OWNERSHIP).map_err(database_error)?;
1626 if ownership
1627 .get(blob_ownership_key(domain_id, hash).as_slice())
1628 .map_err(database_error)?
1629 .is_some()
1630 {
1631 return Ok(true);
1632 }
1633 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1634 let quarantined = read
1635 .open_table(QUARANTINED_COMMITS)
1636 .map_err(database_error)?;
1637 Ok(commits
1638 .get(hash.as_bytes().as_slice())
1639 .map_err(database_error)?
1640 .is_some_and(|metadata| metadata.value().starts_with(domain_id.as_bytes()))
1641 && quarantined
1642 .get(hash.as_bytes().as_slice())
1643 .map_err(database_error)?
1644 .is_none())
1645 }
1646
1647 pub fn list_stored_commit_ids(&self, domain_id: DomainId) -> Result<Vec<CommitId>, StoreError> {
1649 let read = self.inner.database.begin_read().map_err(database_error)?;
1650 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1651 let quarantined = read
1652 .open_table(QUARANTINED_COMMITS)
1653 .map_err(database_error)?;
1654 let mut ids = Vec::new();
1655 for entry in commits.iter().map_err(database_error)? {
1656 let (key, value) = entry.map_err(database_error)?;
1657 if value.value().starts_with(domain_id.as_bytes())
1658 && quarantined
1659 .get(key.value())
1660 .map_err(database_error)?
1661 .is_none()
1662 {
1663 let bytes: [u8; 32] = key.value().try_into().map_err(|_| {
1664 StoreError::CorruptDerivedMetadata("commit ID length is invalid".into())
1665 })?;
1666 ids.push(CommitId::from_bytes(bytes));
1667 }
1668 }
1669 ids.sort_unstable();
1670 Ok(ids)
1671 }
1672
1673 pub fn list_stored_commit_ids_through(
1675 &self,
1676 domain_id: DomainId,
1677 max_epoch: u64,
1678 ) -> Result<Vec<CommitId>, StoreError> {
1679 let read = self.inner.database.begin_read().map_err(database_error)?;
1680 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1681 let quarantined = read
1682 .open_table(QUARANTINED_COMMITS)
1683 .map_err(database_error)?;
1684 let mut ids = Vec::new();
1685 for entry in commits.iter().map_err(database_error)? {
1686 let (key, value) = entry.map_err(database_error)?;
1687 let metadata = value.value();
1688 if metadata.len() != 112 || !metadata.starts_with(domain_id.as_bytes()) {
1689 continue;
1690 }
1691 let epoch = u64::from_be_bytes(metadata[64..72].try_into().map_err(|_| {
1692 StoreError::CorruptDerivedMetadata("commit epoch metadata is invalid".into())
1693 })?);
1694 if epoch <= max_epoch
1695 && quarantined
1696 .get(key.value())
1697 .map_err(database_error)?
1698 .is_none()
1699 {
1700 let bytes: [u8; 32] = key.value().try_into().map_err(|_| {
1701 StoreError::CorruptDerivedMetadata("commit ID length is invalid".into())
1702 })?;
1703 ids.push(CommitId::from_bytes(bytes));
1704 }
1705 }
1706 ids.sort_unstable();
1707 Ok(ids)
1708 }
1709
1710 pub fn author_sequences(
1712 &self,
1713 domain_id: DomainId,
1714 ) -> Result<Vec<(AuthorId, u64)>, StoreError> {
1715 let read = self.inner.database.begin_read().map_err(database_error)?;
1716 let sequences = read.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
1717 let mut values = Vec::new();
1718 for entry in sequences.iter().map_err(database_error)? {
1719 let (key, sequence) = entry.map_err(database_error)?;
1720 let key = key.value();
1721 if key.starts_with(domain_id.as_bytes()) {
1722 let author: [u8; 32] = key[32..].try_into().map_err(|_| {
1723 StoreError::CorruptDerivedMetadata("author sequence key is invalid".into())
1724 })?;
1725 values.push((AuthorId::from_bytes(author), sequence.value()));
1726 }
1727 }
1728 values.sort_unstable_by_key(|(author, _)| *author);
1729 Ok(values)
1730 }
1731
1732 pub fn export_materialized(
1734 &self,
1735 domain_id: DomainId,
1736 ) -> Result<Vec<MaterializedRecord>, StoreError> {
1737 let read = self.inner.database.begin_read().map_err(database_error)?;
1738 let schemas = read.open_table(SCHEMAS).map_err(database_error)?;
1739 let records = read.open_table(RECORDS).map_err(database_error)?;
1740 let record_indexes = read.open_table(RECORD_INDEX_KEYS).map_err(database_error)?;
1741 let mut schema_values = BTreeMap::<CollectionId, SchemaDescriptor>::new();
1742 for entry in schemas.iter().map_err(database_error)? {
1743 let (key, value) = entry.map_err(database_error)?;
1744 if key.value().starts_with(domain_id.as_bytes()) {
1745 let schema = SchemaDescriptor::decode_canonical(value.value())?;
1746 schema_values
1747 .entry(schema.collection_id())
1748 .and_modify(|current| {
1749 if schema.version() > current.version() {
1750 *current = schema.clone();
1751 }
1752 })
1753 .or_insert(schema);
1754 }
1755 }
1756
1757 let mut exported = Vec::new();
1758 for schema in schema_values.into_values() {
1759 let prefix = record_prefix(domain_id, schema.collection_id());
1760 for entry in records.iter().map_err(database_error)? {
1761 let (key, value) = entry.map_err(database_error)?;
1762 let record_key = key.value();
1763 if !record_key.starts_with(&prefix) {
1764 continue;
1765 }
1766 let record_id = record_key[prefix.len()..].to_vec();
1767 let index_keys = record_indexes
1768 .get(record_key)
1769 .map_err(database_error)?
1770 .map(|encoded| decode_key_list(encoded.value()))
1771 .transpose()?
1772 .unwrap_or_default();
1773 let indexes = index_keys
1774 .iter()
1775 .map(|key| {
1776 decode_index_value(key, domain_id, schema.collection_id(), &record_id)
1777 })
1778 .collect::<Result<Vec<_>, _>>()?;
1779 exported.push(
1780 MaterializedRecord::upsert(
1781 schema.clone(),
1782 record_id,
1783 value.value().to_vec(),
1784 indexes,
1785 )
1786 .map_err(|error| StoreError::CorruptDerivedMetadata(error.to_string()))?,
1787 );
1788 }
1789 }
1790 Ok(exported)
1791 }
1792
1793 pub fn replace_materialized(
1795 &self,
1796 domain_id: DomainId,
1797 records: &[MaterializedRecord],
1798 ) -> Result<(), StoreError> {
1799 let write = self.inner.database.begin_write().map_err(database_error)?;
1800 clear_materialized_domain(&write, domain_id)?;
1801 apply_materialized_records(&write, domain_id, records)?;
1802 write.commit().map_err(database_error)
1803 }
1804
1805 #[allow(clippy::too_many_arguments)]
1807 pub fn install_snapshot(
1808 &self,
1809 domain_id: DomainId,
1810 records: &[MaterializedRecord],
1811 covered_commits: &[CommitId],
1812 frontier_commits: &[CommitId],
1813 author_sequences: &[(AuthorId, u64)],
1814 activations: &[SchemaActivation],
1815 observed_at_unix_seconds: u64,
1816 baseline: &SnapshotBaseline,
1817 ) -> Result<(), StoreError> {
1818 validate_snapshot_metadata(records, covered_commits, frontier_commits, author_sequences)?;
1819 if baseline.domain_id() != domain_id
1820 || baseline.frontier() != frontier_commits
1821 || baseline.author_sequences() != author_sequences
1822 {
1823 return Err(StoreError::InvalidSnapshotMetadata(
1824 "snapshot recovery baseline differs from signed metadata".into(),
1825 ));
1826 }
1827 let baseline_bytes = baseline.encode_canonical()?;
1828 let write = self.inner.database.begin_write().map_err(database_error)?;
1829 clear_reconciliation_domain(&write, domain_id)?;
1830 clear_materialized_domain(&write, domain_id)?;
1831 clear_domain_migration_metadata(&write, domain_id)?;
1832 clear_domain_snapshot_metadata(&write, domain_id)?;
1833 apply_materialized_records(&write, domain_id, records)?;
1834 install_snapshot_activations(&write, domain_id, records, activations)?;
1835 {
1836 write
1837 .open_table(SNAPSHOT_BASELINES)
1838 .map_err(database_error)?
1839 .insert(domain_id.as_bytes().as_slice(), baseline_bytes.as_slice())
1840 .map_err(database_error)?;
1841 let mut coverage = write
1842 .open_table(SNAPSHOT_COVERAGE)
1843 .map_err(database_error)?;
1844 for commit_id in covered_commits {
1845 coverage
1846 .insert(
1847 commit_id.as_bytes().as_slice(),
1848 domain_id.as_bytes().as_slice(),
1849 )
1850 .map_err(database_error)?;
1851 }
1852 }
1853 {
1854 let mut frontier = write.open_table(FRONTIER).map_err(database_error)?;
1855 for commit_id in frontier_commits {
1856 let key = frontier_key(domain_id, *commit_id);
1857 frontier.insert(key.as_slice(), 0).map_err(database_error)?;
1858 }
1859 }
1860 {
1861 let mut sequences = write.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
1862 for (author, sequence) in author_sequences {
1863 let key = author_sequence_key(domain_id, *author);
1864 sequences
1865 .insert(key.as_slice(), *sequence)
1866 .map_err(database_error)?;
1867 }
1868 }
1869 write
1870 .open_table(OPERATIONAL_METADATA)
1871 .map_err(database_error)?
1872 .insert(
1873 operational_metadata_key(domain_id, SNAPSHOT_OBSERVED_AT_TAG).as_slice(),
1874 observed_at_unix_seconds,
1875 )
1876 .map_err(database_error)?;
1877 write
1878 .open_table(RECONCILIATION_READY)
1879 .map_err(database_error)?
1880 .insert(domain_id.as_bytes().as_slice(), 1)
1881 .map_err(database_error)?;
1882 write.commit().map_err(database_error)
1883 }
1884
1885 pub async fn load_commit(&self, commit_id: CommitId) -> Result<CommitEnvelope, StoreError> {
1887 let read = self.inner.database.begin_read().map_err(database_error)?;
1888 let commits = read.open_table(COMMIT_METADATA).map_err(database_error)?;
1889 if commits
1890 .get(commit_id.as_bytes().as_slice())
1891 .map_err(database_error)?
1892 .is_none()
1893 {
1894 return Err(StoreError::MissingCommitBlob(commit_id));
1895 }
1896 drop(commits);
1897 drop(read);
1898 let hash = BlobHash::from_bytes(commit_id.to_bytes());
1899 let bytes = self
1900 .inner
1901 .blobs
1902 .get(hash)
1903 .await?
1904 .ok_or(StoreError::MissingCommitBlob(commit_id))?;
1905 let envelope = CommitEnvelope::decode_canonical(&bytes)?;
1906 if envelope.commit_id() != commit_id {
1907 return Err(StoreError::BlobHashMismatch);
1908 }
1909 Ok(envelope)
1910 }
1911
1912 pub async fn scan_commit_orphans(&self) -> Result<Vec<CommitId>, StoreError> {
1917 let mut orphans = Vec::new();
1918 for hash in self.inner.blobs.list().await? {
1919 let Some(bytes) = self.inner.blobs.get(hash).await? else {
1920 continue;
1921 };
1922 let Ok(envelope) = CommitEnvelope::decode_canonical(&bytes) else {
1923 continue;
1924 };
1925 let commit_id = envelope.commit_id();
1926 if commit_id.to_bytes() == hash.to_bytes()
1927 && !self.contains_commit(envelope.header().domain_id(), commit_id)?
1928 {
1929 orphans.push(commit_id);
1930 }
1931 }
1932 orphans.sort_unstable();
1933 orphans.dedup();
1934 Ok(orphans)
1935 }
1936
1937 pub fn bootstrap_domain(
1939 &self,
1940 descriptor: &DomainDescriptor,
1941 root: &CapabilityCertificate,
1942 ) -> Result<(), StoreError> {
1943 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
1944 StoreError::InvalidAuthority("control writer lock is unavailable".into())
1945 })?;
1946 descriptor.verify().map_err(authority_error)?;
1947 root.verify_signature().map_err(authority_error)?;
1948 if root.domain_id() != descriptor.domain_id()
1949 || root.subject() != descriptor.owner()
1950 || root.issuer_capability().is_some()
1951 || !root.permissions().contains(Permission::Admin)
1952 {
1953 return Err(StoreError::InvalidAuthority(
1954 "owner capability does not match descriptor".into(),
1955 ));
1956 }
1957 let descriptor_bytes = descriptor.encode_canonical().map_err(authority_error)?;
1958 let capability_bytes = root.encode_canonical().map_err(authority_error)?;
1959 let capability_id = root.id().map_err(authority_error)?;
1960 let write = self.inner.database.begin_write().map_err(database_error)?;
1961 {
1962 let mut descriptors = write
1963 .open_table(DOMAIN_DESCRIPTORS)
1964 .map_err(database_error)?;
1965 let existing = descriptors
1966 .get(descriptor.domain_id().as_bytes().as_slice())
1967 .map_err(database_error)?
1968 .map(|value| value.value().to_vec());
1969 if existing
1970 .as_deref()
1971 .is_some_and(|bytes| bytes != descriptor_bytes)
1972 {
1973 return Err(StoreError::InvalidAuthority(
1974 "domain descriptor conflicts with genesis".into(),
1975 ));
1976 }
1977 descriptors
1978 .insert(
1979 descriptor.domain_id().as_bytes().as_slice(),
1980 descriptor_bytes.as_slice(),
1981 )
1982 .map_err(database_error)?;
1983 }
1984 insert_capability_tables(&write, root, capability_id, &capability_bytes)?;
1985 write.commit().map_err(database_error)
1986 }
1987
1988 pub fn import_domain_authority_atomic(
1990 &self,
1991 descriptor: &DomainDescriptor,
1992 chain: &[CapabilityCertificate],
1993 ) -> Result<CapabilityId, StoreError> {
1994 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
1995 StoreError::InvalidAuthority("control writer lock is unavailable".into())
1996 })?;
1997 descriptor.verify().map_err(authority_error)?;
1998 if chain.is_empty() || chain.len() > 9 {
1999 return Err(StoreError::InvalidAuthority(
2000 "capability chain has invalid depth".into(),
2001 ));
2002 }
2003 let root = &chain[0];
2004 root.verify_signature().map_err(authority_error)?;
2005 if root.domain_id() != descriptor.domain_id()
2006 || root.subject() != descriptor.owner()
2007 || root.issuer_capability().is_some()
2008 || !root.permissions().contains(Permission::Admin)
2009 {
2010 return Err(StoreError::InvalidAuthority(
2011 "capability root does not match descriptor".into(),
2012 ));
2013 }
2014 for pair in chain.windows(2) {
2015 pair[1]
2016 .validate_delegation(&pair[0])
2017 .map_err(authority_error)?;
2018 }
2019 if let Some(existing) = self.domain_descriptor(descriptor.domain_id())?
2020 && existing != *descriptor
2021 {
2022 return Err(StoreError::InvalidAuthority(
2023 "domain descriptor conflicts with genesis".into(),
2024 ));
2025 }
2026 let mut validated = Vec::with_capacity(chain.len());
2027 for capability in chain {
2028 if capability.domain_id() != descriptor.domain_id() {
2029 return Err(StoreError::InvalidAuthority(
2030 "capability chain crosses domains".into(),
2031 ));
2032 }
2033 let capability_id = capability.id().map_err(authority_error)?;
2034 let bytes = capability.encode_canonical().map_err(authority_error)?;
2035 if let Some(stored) = self.capability(capability_id)?
2036 && stored != *capability
2037 {
2038 return Err(StoreError::InvalidAuthority(
2039 "capability ID collision".into(),
2040 ));
2041 }
2042 validated.push((capability.clone(), capability_id, bytes));
2043 }
2044 if !self.capability_anchors_are_current(descriptor.domain_id(), chain)? {
2045 return Err(StoreError::InvalidAuthority(
2046 "capability chain is revoked or has an invalid control anchor".into(),
2047 ));
2048 }
2049
2050 let descriptor_bytes = descriptor.encode_canonical().map_err(authority_error)?;
2051 let (leaf, leaf_id, _) = validated
2052 .last()
2053 .ok_or_else(|| StoreError::InvalidAuthority("capability chain is empty".into()))?;
2054 let write = self.inner.database.begin_write().map_err(database_error)?;
2055 write
2056 .open_table(DOMAIN_DESCRIPTORS)
2057 .map_err(database_error)?
2058 .insert(
2059 descriptor.domain_id().as_bytes().as_slice(),
2060 descriptor_bytes.as_slice(),
2061 )
2062 .map_err(database_error)?;
2063 for (_, capability_id, bytes) in &validated {
2064 insert_capability_record(&write, *capability_id, bytes)?;
2065 }
2066 let (root, root_id, _) = validated
2067 .first()
2068 .ok_or_else(|| StoreError::InvalidAuthority("capability chain is empty".into()))?;
2069 set_active_capability(&write, root, *root_id)?;
2070 set_active_capability(&write, leaf, *leaf_id)?;
2071 {
2072 let mut trusted = write.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
2073 for (capability, _, _) in &validated {
2074 let key = trusted_author_key(descriptor.domain_id(), capability.subject());
2075 trusted.insert(key.as_slice(), 1).map_err(database_error)?;
2076 }
2077 }
2078 write.commit().map_err(database_error)?;
2079 Ok(*leaf_id)
2080 }
2081
2082 #[allow(clippy::too_many_lines)]
2084 pub fn import_authority_checkpoint_atomic(
2085 &self,
2086 invitation: &Invitation,
2087 ) -> Result<CapabilityId, StoreError> {
2088 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2089 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2090 })?;
2091 let encoded_invitation = invitation.encode_canonical().map_err(authority_error)?;
2092 Invitation::decode_canonical(&encoded_invitation).map_err(authority_error)?;
2093 let descriptor = invitation.descriptor();
2094 if self.domain_descriptor(descriptor.domain_id())?.is_some()
2095 || self.control_head(descriptor.domain_id())?.is_some()
2096 {
2097 return Err(StoreError::InvalidAuthority(
2098 "authority checkpoint cannot replace existing domain authority".into(),
2099 ));
2100 }
2101 let checkpoint = invitation.authority_checkpoint();
2102 let mut evidence = BTreeMap::<CapabilityId, CapabilityCertificate>::new();
2103 for capability in invitation.authority_evidence() {
2104 let id = capability.id().map_err(authority_error)?;
2105 evidence.insert(id, capability.clone());
2106 }
2107 if evidence.keys().copied().collect::<Vec<_>>() != checkpoint.capability_ids() {
2108 return Err(StoreError::InvalidAuthority(
2109 "checkpoint capability evidence is incomplete".into(),
2110 ));
2111 }
2112 let leaf = invitation
2113 .capability_chain()
2114 .last()
2115 .ok_or_else(|| StoreError::InvalidAuthority("checkpoint leaf is missing".into()))?;
2116 let leaf_id = leaf.id().map_err(authority_error)?;
2117 let checkpoint_bytes = checkpoint.encode_canonical().map_err(authority_error)?;
2118 let descriptor_bytes = descriptor.encode_canonical().map_err(authority_error)?;
2119 let write = self.inner.database.begin_write().map_err(database_error)?;
2120 write
2121 .open_table(DOMAIN_DESCRIPTORS)
2122 .map_err(database_error)?
2123 .insert(
2124 descriptor.domain_id().as_bytes().as_slice(),
2125 descriptor_bytes.as_slice(),
2126 )
2127 .map_err(database_error)?;
2128 for (id, capability) in &evidence {
2129 let bytes = capability.encode_canonical().map_err(authority_error)?;
2130 insert_capability_record(&write, *id, &bytes)?;
2131 }
2132 for id in checkpoint.active_capability_ids() {
2133 let capability = evidence.get(id).ok_or_else(|| {
2134 StoreError::InvalidAuthority("checkpoint active capability is missing".into())
2135 })?;
2136 set_active_capability(&write, capability, *id)?;
2137 }
2138 write
2139 .open_table(AUTHORITY_CHECKPOINTS)
2140 .map_err(database_error)?
2141 .insert(
2142 descriptor.domain_id().as_bytes().as_slice(),
2143 checkpoint_bytes.as_slice(),
2144 )
2145 .map_err(database_error)?;
2146 for (subject, sequence) in checkpoint.subject_revocations() {
2147 let key = subject_revocation_key(descriptor.domain_id(), *subject);
2148 write
2149 .open_table(SUBJECT_REVOCATIONS)
2150 .map_err(database_error)?
2151 .insert(key.as_slice(), *sequence)
2152 .map_err(database_error)?;
2153 }
2154 if let Some(head) = checkpoint.head() {
2155 let control_id = head.id().map_err(authority_error)?;
2156 insert_control_observation(&write, head, control_id)?;
2157 let bytes = head.encode_canonical().map_err(authority_error)?;
2158 write
2159 .open_table(CONTROL_HEADS)
2160 .map_err(database_error)?
2161 .insert(
2162 descriptor.domain_id().as_bytes().as_slice(),
2163 bytes.as_slice(),
2164 )
2165 .map_err(database_error)?;
2166 }
2167 let mut trusted = write.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
2168 for id in checkpoint.active_capability_ids() {
2169 let capability = evidence.get(id).ok_or_else(|| {
2170 StoreError::InvalidAuthority("checkpoint active capability is missing".into())
2171 })?;
2172 let key = trusted_author_key(descriptor.domain_id(), capability.subject());
2173 trusted.insert(key.as_slice(), 1).map_err(database_error)?;
2174 }
2175 drop(trusted);
2176 write.commit().map_err(database_error)?;
2177 Ok(leaf_id)
2178 }
2179
2180 pub fn install_capability(
2182 &self,
2183 capability: &CapabilityCertificate,
2184 ) -> Result<CapabilityId, StoreError> {
2185 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2186 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2187 })?;
2188 let issuer_id = capability.issuer_capability().ok_or_else(|| {
2189 StoreError::InvalidAuthority("delegated capability lacks issuer".into())
2190 })?;
2191 let capability_id = capability.id().map_err(authority_error)?;
2192 let mut chain = self.capability_chain(issuer_id)?;
2193 chain.push(capability.clone());
2194 let validated =
2195 self.validate_capability_chain_for_install(&chain, Some(capability_id), false)?;
2196 let write = self.inner.database.begin_write().map_err(database_error)?;
2197 for (_, id, bytes) in &validated {
2198 insert_capability_record(&write, *id, bytes)?;
2199 }
2200 set_active_capability(&write, capability, capability_id)?;
2201 write.commit().map_err(database_error)?;
2202 Ok(capability_id)
2203 }
2204
2205 pub fn install_capability_chain_atomic(
2211 &self,
2212 chain: &[CapabilityCertificate],
2213 ) -> Result<CapabilityId, StoreError> {
2214 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2215 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2216 })?;
2217 let validated = self.validate_capability_chain_for_install(chain, None, false)?;
2218 let (leaf, leaf_id, _) = validated
2219 .last()
2220 .ok_or_else(|| StoreError::InvalidAuthority("capability chain is empty".into()))?;
2221 let write = self.inner.database.begin_write().map_err(database_error)?;
2222 for (_, capability_id, bytes) in &validated {
2223 insert_capability_record(&write, *capability_id, bytes)?;
2224 }
2225 set_active_capability(&write, leaf, *leaf_id)?;
2226 write.commit().map_err(database_error)?;
2227 Ok(*leaf_id)
2228 }
2229
2230 pub fn capability(
2232 &self,
2233 capability_id: CapabilityId,
2234 ) -> Result<Option<CapabilityCertificate>, StoreError> {
2235 let read = self.inner.database.begin_read().map_err(database_error)?;
2236 let capabilities = read.open_table(CAPABILITIES).map_err(database_error)?;
2237 capabilities
2238 .get(capability_id.as_bytes().as_slice())
2239 .map_err(database_error)?
2240 .map(|value| {
2241 CapabilityCertificate::decode_canonical(value.value()).map_err(authority_error)
2242 })
2243 .transpose()
2244 }
2245
2246 pub fn capability_chain(
2248 &self,
2249 leaf: CapabilityId,
2250 ) -> Result<Vec<CapabilityCertificate>, StoreError> {
2251 let mut reversed = Vec::new();
2252 let mut next = Some(leaf);
2253 while let Some(capability_id) = next {
2254 if reversed.len() >= 9 {
2255 return Err(StoreError::InvalidAuthority(
2256 "capability chain exceeds depth bound".into(),
2257 ));
2258 }
2259 let capability = self.capability(capability_id)?.ok_or_else(|| {
2260 StoreError::InvalidAuthority("capability chain link is missing".into())
2261 })?;
2262 next = capability.issuer_capability();
2263 reversed.push(capability);
2264 }
2265 reversed.reverse();
2266 for pair in reversed.windows(2) {
2267 pair[1]
2268 .validate_delegation(&pair[0])
2269 .map_err(authority_error)?;
2270 }
2271 Ok(reversed)
2272 }
2273
2274 pub fn domain_descriptor(
2276 &self,
2277 domain_id: DomainId,
2278 ) -> Result<Option<DomainDescriptor>, StoreError> {
2279 let read = self.inner.database.begin_read().map_err(database_error)?;
2280 let descriptors = read
2281 .open_table(DOMAIN_DESCRIPTORS)
2282 .map_err(database_error)?;
2283 descriptors
2284 .get(domain_id.as_bytes().as_slice())
2285 .map_err(database_error)?
2286 .map(|value| DomainDescriptor::decode_canonical(value.value()).map_err(authority_error))
2287 .transpose()
2288 }
2289
2290 pub fn active_capability(
2292 &self,
2293 domain_id: DomainId,
2294 subject: AuthorId,
2295 ) -> Result<Option<CapabilityId>, StoreError> {
2296 let read = self.inner.database.begin_read().map_err(database_error)?;
2297 let active = read
2298 .open_table(ACTIVE_CAPABILITIES)
2299 .map_err(database_error)?;
2300 let key = active_capability_key(domain_id, subject);
2301 active
2302 .get(key.as_slice())
2303 .map_err(database_error)?
2304 .map(|value| {
2305 let bytes: [u8; 32] = value.value().try_into().map_err(|_| {
2306 StoreError::CorruptDerivedMetadata("active capability ID is invalid".into())
2307 })?;
2308 Ok(CapabilityId::from_bytes(bytes))
2309 })
2310 .transpose()
2311 }
2312
2313 pub fn active_capabilities(
2315 &self,
2316 domain_id: DomainId,
2317 ) -> Result<Vec<CapabilityId>, StoreError> {
2318 let read = self.inner.database.begin_read().map_err(database_error)?;
2319 let active = read
2320 .open_table(ACTIVE_CAPABILITIES)
2321 .map_err(database_error)?;
2322 let mut ids = Vec::new();
2323 for entry in active.iter().map_err(database_error)? {
2324 let (key, value) = entry.map_err(database_error)?;
2325 if !key.value().starts_with(domain_id.as_bytes()) {
2326 continue;
2327 }
2328 let bytes: [u8; 32] = value.value().try_into().map_err(|_| {
2329 StoreError::CorruptDerivedMetadata("active capability ID is invalid".into())
2330 })?;
2331 ids.push(CapabilityId::from_bytes(bytes));
2332 }
2333 ids.sort_unstable();
2334 ids.dedup();
2335 Ok(ids)
2336 }
2337
2338 pub fn authorize_capability(
2340 &self,
2341 domain_id: DomainId,
2342 subject: AuthorId,
2343 capability_id: CapabilityId,
2344 permission: Permission,
2345 epoch: u64,
2346 ) -> Result<bool, StoreError> {
2347 self.ensure_control_not_forked(domain_id)?;
2348 if self.active_capability(domain_id, subject)? != Some(capability_id) {
2349 return Ok(false);
2350 }
2351 let chain = match self.capability_chain(capability_id) {
2352 Ok(chain) => chain,
2353 Err(StoreError::InvalidAuthority(_)) => return Ok(false),
2354 Err(error) => return Err(error),
2355 };
2356 if !self.capability_anchors_are_current(domain_id, &chain)? {
2357 return Ok(false);
2358 }
2359 let Some(capability) = chain.last() else {
2360 return Ok(false);
2361 };
2362 Ok(capability.domain_id() == domain_id
2363 && capability.subject() == subject
2364 && capability.permits(permission, epoch))
2365 }
2366
2367 pub fn control_head(
2369 &self,
2370 domain_id: DomainId,
2371 ) -> Result<Option<ControlTransition>, StoreError> {
2372 self.ensure_control_not_forked(domain_id)?;
2373 let read = self.inner.database.begin_read().map_err(database_error)?;
2374 let heads = read.open_table(CONTROL_HEADS).map_err(database_error)?;
2375 heads
2376 .get(domain_id.as_bytes().as_slice())
2377 .map_err(database_error)?
2378 .map(|value| {
2379 ControlTransition::decode_canonical(value.value()).map_err(authority_error)
2380 })
2381 .transpose()
2382 }
2383
2384 pub fn accepted_control_chain(
2386 &self,
2387 domain_id: DomainId,
2388 ) -> Result<Vec<ControlTransition>, StoreError> {
2389 self.ensure_control_not_forked(domain_id)?;
2390 let baseline = self.authority_checkpoint(domain_id)?;
2391 let baseline_id = baseline
2392 .as_ref()
2393 .and_then(AuthorityCheckpoint::head)
2394 .map(ControlTransition::id)
2395 .transpose()
2396 .map_err(authority_error)?;
2397 let Some(mut current) = self.control_head(domain_id)? else {
2398 return Ok(Vec::new());
2399 };
2400 let mut reversed = Vec::new();
2401 loop {
2402 let current_id = current.id().map_err(authority_error)?;
2403 let predecessor = current.predecessor();
2404 reversed.push(current);
2405 if Some(current_id) == baseline_id {
2406 break;
2407 }
2408 let Some(predecessor) = predecessor else {
2409 break;
2410 };
2411 current = self
2412 .control_object(domain_id, predecessor)?
2413 .ok_or_else(|| {
2414 StoreError::CorruptDerivedMetadata(
2415 "accepted control predecessor is missing".into(),
2416 )
2417 })?;
2418 }
2419 reversed.reverse();
2420 let first_sequence = baseline
2421 .as_ref()
2422 .and_then(AuthorityCheckpoint::head)
2423 .map_or(1, ControlTransition::sequence);
2424 for (index, control) in reversed.iter().enumerate() {
2425 let expected = u64::try_from(index)
2426 .map_err(|_| {
2427 StoreError::CorruptDerivedMetadata("control chain is too long".into())
2428 })?
2429 .checked_add(first_sequence)
2430 .ok_or_else(|| {
2431 StoreError::CorruptDerivedMetadata("control sequence overflow".into())
2432 })?;
2433 if control.domain_id() != domain_id || control.sequence() != expected {
2434 return Err(StoreError::CorruptDerivedMetadata(
2435 "accepted control chain is not contiguous".into(),
2436 ));
2437 }
2438 }
2439 Ok(reversed)
2440 }
2441
2442 pub fn authority_checkpoint(
2444 &self,
2445 domain_id: DomainId,
2446 ) -> Result<Option<AuthorityCheckpoint>, StoreError> {
2447 let read = self.inner.database.begin_read().map_err(database_error)?;
2448 let table = read
2449 .open_table(AUTHORITY_CHECKPOINTS)
2450 .map_err(database_error)?;
2451 let checkpoint = table
2452 .get(domain_id.as_bytes().as_slice())
2453 .map_err(database_error)?
2454 .map(|value| {
2455 AuthorityCheckpoint::decode_canonical(value.value()).map_err(authority_error)
2456 })
2457 .transpose()?;
2458 if checkpoint
2459 .as_ref()
2460 .is_some_and(|checkpoint| checkpoint.domain_id() != domain_id)
2461 {
2462 return Err(StoreError::CorruptDerivedMetadata(
2463 "authority checkpoint crosses domains".into(),
2464 ));
2465 }
2466 Ok(checkpoint)
2467 }
2468
2469 pub fn authority_checkpoint_consistent(&self, domain_id: DomainId) -> Result<bool, StoreError> {
2471 let Some(checkpoint) = self.authority_checkpoint(domain_id)? else {
2472 return Ok(true);
2473 };
2474 if self.domain_descriptor(domain_id)?.is_none() {
2475 return Ok(false);
2476 }
2477 if let Some(head) = checkpoint.head() {
2478 let id = head.id().map_err(authority_error)?;
2479 if self.control_transition(domain_id, id)?.as_ref() != Some(head) {
2480 return Ok(false);
2481 }
2482 }
2483 for id in checkpoint.capability_ids() {
2484 if self
2485 .capability(*id)?
2486 .is_none_or(|capability| capability.domain_id() != domain_id)
2487 {
2488 return Ok(false);
2489 }
2490 }
2491 Ok(true)
2492 }
2493
2494 pub fn control_transition(
2496 &self,
2497 domain_id: DomainId,
2498 control_id: [u8; 32],
2499 ) -> Result<Option<ControlTransition>, StoreError> {
2500 self.control_object(domain_id, control_id)
2501 }
2502
2503 pub fn accepted_control_at_sequence(
2505 &self,
2506 domain_id: DomainId,
2507 sequence: u64,
2508 ) -> Result<Option<ControlTransition>, StoreError> {
2509 Ok(self
2510 .accepted_control_chain(domain_id)?
2511 .into_iter()
2512 .find(|control| control.sequence() == sequence))
2513 }
2514
2515 pub fn current_epoch(&self, domain_id: DomainId) -> Result<u64, StoreError> {
2517 self.ensure_control_not_forked(domain_id)?;
2518 Ok(self
2519 .control_head(domain_id)?
2520 .map_or(0, |head| head.new_epoch()))
2521 }
2522
2523 pub fn current_control_sequence(&self, domain_id: DomainId) -> Result<u64, StoreError> {
2525 Ok(self
2526 .control_head(domain_id)?
2527 .as_ref()
2528 .map_or(0, ControlTransition::sequence))
2529 }
2530
2531 pub fn subject_revocations_consistent(&self, domain_id: DomainId) -> Result<bool, StoreError> {
2533 let expected = self.expected_subject_revocations(domain_id)?;
2534 let read = self.inner.database.begin_read().map_err(database_error)?;
2535 let table = read
2536 .open_table(SUBJECT_REVOCATIONS)
2537 .map_err(database_error)?;
2538 let mut actual = BTreeMap::new();
2539 for entry in table.iter().map_err(database_error)? {
2540 let (key, value) = entry.map_err(database_error)?;
2541 let key = key.value();
2542 if key.len() != 64 || !key.starts_with(domain_id.as_bytes()) {
2543 continue;
2544 }
2545 let subject = AuthorId::from_bytes(key[32..].try_into().map_err(|_| {
2546 StoreError::CorruptDerivedMetadata("subject revocation endpoint is invalid".into())
2547 })?);
2548 actual.insert(subject, value.value());
2549 }
2550 Ok(actual == expected)
2551 }
2552
2553 pub fn subject_revocations(
2555 &self,
2556 domain_id: DomainId,
2557 ) -> Result<Vec<(AuthorId, u64)>, StoreError> {
2558 Ok(self
2559 .expected_subject_revocations(domain_id)?
2560 .into_iter()
2561 .collect())
2562 }
2563
2564 pub fn repair_subject_revocations(&self, domain_id: DomainId) -> Result<(), StoreError> {
2566 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2567 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2568 })?;
2569 let expected = self.expected_subject_revocations(domain_id)?;
2570 let write = self.inner.database.begin_write().map_err(database_error)?;
2571 {
2572 let mut table = write
2573 .open_table(SUBJECT_REVOCATIONS)
2574 .map_err(database_error)?;
2575 table
2576 .retain(|key, _| !key.starts_with(domain_id.as_bytes()))
2577 .map_err(database_error)?;
2578 for (subject, sequence) in expected {
2579 let key = subject_revocation_key(domain_id, subject);
2580 table
2581 .insert(key.as_slice(), sequence)
2582 .map_err(database_error)?;
2583 }
2584 }
2585 write.commit().map_err(database_error)
2586 }
2587
2588 pub fn current_controller(&self, domain_id: DomainId) -> Result<AuthorId, StoreError> {
2590 self.ensure_control_not_forked(domain_id)?;
2591 if let Some(head) = self.control_head(domain_id)? {
2592 return Ok(head.controller());
2593 }
2594 self.domain_descriptor(domain_id)?
2595 .map(|descriptor| descriptor.owner())
2596 .ok_or_else(|| StoreError::InvalidAuthority("domain descriptor is missing".into()))
2597 }
2598
2599 pub fn current_writer(&self, domain_id: DomainId) -> Result<Option<AuthorId>, StoreError> {
2601 self.ensure_control_not_forked(domain_id)?;
2602 if let Some(head) = self.control_head(domain_id)? {
2603 return Ok(Some(head.writer()));
2604 }
2605 Ok(self
2606 .domain_descriptor(domain_id)?
2607 .map(|descriptor| descriptor.owner()))
2608 }
2609
2610 pub fn observe_control(&self, control: &ControlTransition) -> Result<(), StoreError> {
2612 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2613 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2614 })?;
2615 control.verify().map_err(authority_error)?;
2616 self.ensure_control_not_forked(control.domain_id())?;
2617 let current = self.control_head(control.domain_id())?;
2618 let control_id = control.id().map_err(authority_error)?;
2619 if current
2620 .as_ref()
2621 .map(ControlTransition::id)
2622 .transpose()
2623 .map_err(authority_error)?
2624 .is_some_and(|current_id| current_id == control_id)
2625 {
2626 return Ok(());
2627 }
2628 if self.control_slot_is_observed(control, control_id)? {
2629 return Ok(());
2630 }
2631 if let Some(head) = current.as_ref()
2632 && control.sequence() == head.sequence()
2633 && control.predecessor() == head.predecessor()
2634 && control.previous_epoch() == head.previous_epoch()
2635 {
2636 self.validate_control_sibling(control, head)?;
2637 return self.freeze_control_fork(head, control);
2638 }
2639 let (sequence, predecessor, epoch) = if let Some(head) = current.as_ref() {
2640 (
2641 head.sequence().checked_add(1).ok_or_else(|| {
2642 StoreError::InvalidAuthority("control sequence exhausted".into())
2643 })?,
2644 Some(head.id().map_err(authority_error)?),
2645 head.new_epoch(),
2646 )
2647 } else {
2648 (1, None, 0)
2649 };
2650 if control.sequence() != sequence
2651 || control.predecessor() != predecessor
2652 || control.previous_epoch() != epoch
2653 || self.frontier(control.domain_id())? != control.cut()
2654 {
2655 return Err(StoreError::InvalidAuthority(
2656 "control transition does not extend the active head".into(),
2657 ));
2658 }
2659 self.validate_control_authority(
2660 control,
2661 self.current_controller(control.domain_id())?,
2662 epoch,
2663 None,
2664 )?;
2665 let write = self.inner.database.begin_write().map_err(database_error)?;
2666 insert_control_observation(&write, control, control_id)?;
2667 write.commit().map_err(database_error)
2668 }
2669
2670 pub fn apply_control(&self, control: &ControlTransition) -> Result<(), StoreError> {
2672 self.apply_control_bundle(control, &[])
2673 }
2674
2675 pub fn apply_control_bundle(
2677 &self,
2678 control: &ControlTransition,
2679 successor_chain: &[CapabilityCertificate],
2680 ) -> Result<(), StoreError> {
2681 let _control_writer = self.inner.control_writer.lock().map_err(|_| {
2682 StoreError::InvalidAuthority("control writer lock is unavailable".into())
2683 })?;
2684 control.verify().map_err(authority_error)?;
2685 self.ensure_control_not_forked(control.domain_id())?;
2686 let current = self.control_head(control.domain_id())?;
2687 let control_id = control.id().map_err(authority_error)?;
2688 if current
2689 .as_ref()
2690 .map(ControlTransition::id)
2691 .transpose()
2692 .map_err(authority_error)?
2693 .is_some_and(|current_id| current_id == control_id)
2694 {
2695 return Ok(());
2696 }
2697 self.control_slot_is_observed(control, control_id)?;
2698 if let Some(head) = current.as_ref()
2699 && control.sequence() == head.sequence()
2700 && control.predecessor() == head.predecessor()
2701 && control.previous_epoch() == head.previous_epoch()
2702 {
2703 self.validate_control_sibling(control, head)?;
2704 return self.freeze_control_fork(head, control);
2705 }
2706 let (sequence, predecessor, epoch) = if let Some(head) = current.as_ref() {
2707 (
2708 head.sequence().checked_add(1).ok_or_else(|| {
2709 StoreError::InvalidAuthority("control sequence exhausted".into())
2710 })?,
2711 Some(head.id().map_err(authority_error)?),
2712 head.new_epoch(),
2713 )
2714 } else {
2715 (1, None, 0)
2716 };
2717 if control.sequence() != sequence
2718 || control.predecessor() != predecessor
2719 || control.previous_epoch() != epoch
2720 || self.frontier(control.domain_id())? != control.cut()
2721 {
2722 return Err(StoreError::InvalidAuthority(
2723 "control transition does not extend the active head".into(),
2724 ));
2725 }
2726 let validated_successor = if successor_chain.is_empty() {
2727 Vec::new()
2728 } else {
2729 self.validate_capability_chain_for_install(
2730 successor_chain,
2731 Some(control.controller_capability()),
2732 true,
2733 )?
2734 };
2735 let candidate_successor_chain = validated_successor
2736 .iter()
2737 .map(|(capability, _, _)| capability.clone())
2738 .collect::<Vec<_>>();
2739 self.validate_control_authority(
2740 control,
2741 self.current_controller(control.domain_id())?,
2742 epoch,
2743 (!candidate_successor_chain.is_empty()).then_some(candidate_successor_chain.as_slice()),
2744 )?;
2745 if control.revoked_subjects().contains(&control.author()) {
2746 return Err(StoreError::InvalidAuthority(
2747 "control author cannot revoke itself".into(),
2748 ));
2749 }
2750 let bytes = control.encode_canonical().map_err(authority_error)?;
2751 let write = self.inner.database.begin_write().map_err(database_error)?;
2752 {
2753 for (_, capability_id, capability_bytes) in &validated_successor {
2754 insert_capability_record(&write, *capability_id, capability_bytes)?;
2755 }
2756 if let Some((leaf, leaf_id, _)) = validated_successor.last() {
2757 set_active_capability(&write, leaf, *leaf_id)?;
2758 }
2759 insert_control_observation(&write, control, control_id)?;
2760 let mut heads = write.open_table(CONTROL_HEADS).map_err(database_error)?;
2761 heads
2762 .insert(control.domain_id().as_bytes().as_slice(), bytes.as_slice())
2763 .map_err(database_error)?;
2764 }
2765 apply_control_revocations(&write, control)?;
2766 write.commit().map_err(database_error)
2767 }
2768
2769 fn validate_control_authority(
2770 &self,
2771 control: &ControlTransition,
2772 expected_controller: AuthorId,
2773 epoch: u64,
2774 candidate_successor_chain: Option<&[CapabilityCertificate]>,
2775 ) -> Result<(), StoreError> {
2776 let can_control = self.authorize_capability(
2777 control.domain_id(),
2778 control.author(),
2779 control.author_capability(),
2780 Permission::Revoke,
2781 epoch,
2782 )? || self.authorize_capability(
2783 control.domain_id(),
2784 control.author(),
2785 control.author_capability(),
2786 Permission::Admin,
2787 epoch,
2788 )?;
2789 let next_chain = if let Some(candidate) = candidate_successor_chain {
2790 candidate.to_vec()
2791 } else {
2792 self.capability_chain(control.controller_capability())?
2793 };
2794 let next_capability = next_chain.last().ok_or_else(|| {
2795 StoreError::InvalidAuthority("successor capability chain is empty".into())
2796 })?;
2797 let next_can_admin = next_capability.id().map_err(authority_error)?
2798 == control.controller_capability()
2799 && self.capability_anchors_are_current_with_historical_nomination(
2800 control.domain_id(),
2801 &next_chain,
2802 candidate_successor_chain.is_some(),
2803 )?
2804 && !next_chain
2805 .iter()
2806 .any(|capability| control.revoked_subjects().contains(&capability.subject()))
2807 && {
2808 let capability = next_capability;
2809 capability.domain_id() == control.domain_id()
2810 && capability.subject() == control.controller()
2811 && capability.permits(Permission::Admin, control.new_epoch())
2812 };
2813 if control.author() != expected_controller || !can_control || !next_can_admin {
2814 return Err(StoreError::InvalidAuthority(
2815 "control transition lacks designated controller authority".into(),
2816 ));
2817 }
2818 Ok(())
2819 }
2820
2821 fn validate_capability_chain_for_install(
2822 &self,
2823 chain: &[CapabilityCertificate],
2824 replace_leaf: Option<CapabilityId>,
2825 allow_historical_nomination: bool,
2826 ) -> Result<Vec<(CapabilityCertificate, CapabilityId, Vec<u8>)>, StoreError> {
2827 if chain.is_empty() || chain.len() > 9 {
2828 return Err(StoreError::InvalidAuthority(
2829 "capability chain has invalid depth".into(),
2830 ));
2831 }
2832 let root = &chain[0];
2833 let descriptor = self
2834 .domain_descriptor(root.domain_id())?
2835 .ok_or_else(|| StoreError::InvalidAuthority("domain descriptor is missing".into()))?;
2836 let root_id = root.id().map_err(authority_error)?;
2837 let stored_root = self
2838 .capability(root_id)?
2839 .ok_or_else(|| StoreError::InvalidAuthority("capability root is unknown".into()))?;
2840 if root.subject() != descriptor.owner()
2841 || root.issuer_capability().is_some()
2842 || stored_root != *root
2843 {
2844 return Err(StoreError::InvalidAuthority(
2845 "capability chain has an invalid root".into(),
2846 ));
2847 }
2848 for pair in chain.windows(2) {
2849 pair[1]
2850 .validate_delegation(&pair[0])
2851 .map_err(authority_error)?;
2852 }
2853 let mut validated = Vec::with_capacity(chain.len());
2854 for capability in chain {
2855 if capability.domain_id() != root.domain_id() {
2856 return Err(StoreError::InvalidAuthority(
2857 "capability chain crosses domains".into(),
2858 ));
2859 }
2860 let capability_id = capability.id().map_err(authority_error)?;
2861 let bytes = capability.encode_canonical().map_err(authority_error)?;
2862 if let Some(stored) = self.capability(capability_id)?
2863 && stored != *capability
2864 {
2865 return Err(StoreError::InvalidAuthority(
2866 "capability ID collision".into(),
2867 ));
2868 }
2869 validated.push((capability.clone(), capability_id, bytes));
2870 }
2871 let (leaf, leaf_id, _) = validated
2872 .last()
2873 .ok_or_else(|| StoreError::InvalidAuthority("capability chain is empty".into()))?;
2874 if replace_leaf.is_some_and(|expected| expected != *leaf_id) {
2875 return Err(StoreError::InvalidAuthority(
2876 "capability chain leaf does not match signed successor".into(),
2877 ));
2878 }
2879 if replace_leaf.is_none()
2880 && let Some(active) = self.active_capability(leaf.domain_id(), leaf.subject())?
2881 && active != *leaf_id
2882 {
2883 let current = self.capability(active)?.ok_or_else(|| {
2884 StoreError::CorruptDerivedMetadata("active capability record is missing".into())
2885 })?;
2886 if leaf.issued_control_sequence() <= current.issued_control_sequence() {
2887 return Err(StoreError::InvalidAuthority(
2888 "capability chain would replace an active leaf without advancing control"
2889 .into(),
2890 ));
2891 }
2892 }
2893 let capabilities = validated
2894 .iter()
2895 .map(|(capability, _, _)| capability.clone())
2896 .collect::<Vec<_>>();
2897 if !self.capability_anchors_are_current_with_historical_nomination(
2898 root.domain_id(),
2899 &capabilities,
2900 allow_historical_nomination,
2901 )? {
2902 return Err(StoreError::InvalidAuthority(
2903 "capability chain is revoked or has an invalid control anchor".into(),
2904 ));
2905 }
2906 Ok(validated)
2907 }
2908
2909 fn capability_anchors_are_current(
2910 &self,
2911 domain_id: DomainId,
2912 chain: &[CapabilityCertificate],
2913 ) -> Result<bool, StoreError> {
2914 self.capability_anchors_are_current_with_historical_nomination(domain_id, chain, false)
2915 }
2916
2917 fn capability_anchors_are_current_with_historical_nomination(
2918 &self,
2919 domain_id: DomainId,
2920 chain: &[CapabilityCertificate],
2921 allow_historical_nomination: bool,
2922 ) -> Result<bool, StoreError> {
2923 let controls = self.accepted_control_chain(domain_id)?;
2924 let checkpoint = self.authority_checkpoint(domain_id)?;
2925 let expected_revocations = self.expected_subject_revocations(domain_id)?;
2926 let read = self.inner.database.begin_read().map_err(database_error)?;
2927 let revoked = read
2928 .open_table(REVOKED_CAPABILITIES)
2929 .map_err(database_error)?;
2930 let subject_revocations = read
2931 .open_table(SUBJECT_REVOCATIONS)
2932 .map_err(database_error)?;
2933 for capability in chain {
2934 if capability.domain_id() != domain_id {
2935 return Ok(false);
2936 }
2937 let capability_id = capability.id().map_err(authority_error)?;
2938 let baseline_sequence = checkpoint
2939 .as_ref()
2940 .and_then(AuthorityCheckpoint::head)
2941 .map(ControlTransition::sequence);
2942 let before_truncated_baseline = baseline_sequence
2943 .is_some_and(|sequence| capability.issued_control_sequence() < sequence);
2944 let checkpoint_allows = before_truncated_baseline
2945 && checkpoint
2946 .as_ref()
2947 .is_some_and(|checkpoint| checkpoint.capability_ids().contains(&capability_id));
2948 let anchor_is_accepted = if before_truncated_baseline {
2949 checkpoint_allows || allow_historical_nomination
2950 } else if capability.issued_control_sequence() == 0 {
2951 capability.issued_control_head().is_none()
2952 } else {
2953 controls
2954 .iter()
2955 .find(|control| control.sequence() == capability.issued_control_sequence())
2956 .map(ControlTransition::id)
2957 .transpose()
2958 .map_err(authority_error)?
2959 == capability.issued_control_head()
2960 };
2961 if !anchor_is_accepted {
2962 return Ok(false);
2963 }
2964 if revoked
2965 .get(capability_id.as_bytes().as_slice())
2966 .map_err(database_error)?
2967 .is_some()
2968 {
2969 return Ok(false);
2970 }
2971 let subject_key = subject_revocation_key(domain_id, capability.subject());
2972 let stored_revocation = subject_revocations
2973 .get(subject_key.as_slice())
2974 .map_err(database_error)?
2975 .map(|sequence| sequence.value());
2976 let expected_revocation = expected_revocations.get(&capability.subject()).copied();
2977 if stored_revocation != expected_revocation {
2978 return Err(StoreError::CorruptDerivedMetadata(
2979 "subject revocation index differs from signed control".into(),
2980 ));
2981 }
2982 if expected_revocation
2983 .is_some_and(|sequence| sequence > capability.issued_control_sequence())
2984 {
2985 return Ok(false);
2986 }
2987 }
2988 Ok(true)
2989 }
2990
2991 fn expected_subject_revocations(
2992 &self,
2993 domain_id: DomainId,
2994 ) -> Result<BTreeMap<AuthorId, u64>, StoreError> {
2995 let checkpoint = self.authority_checkpoint(domain_id)?;
2996 let mut revocations: BTreeMap<AuthorId, u64> = checkpoint
2997 .as_ref()
2998 .map(|checkpoint| checkpoint.subject_revocations().iter().copied().collect())
2999 .unwrap_or_default();
3000 let baseline_sequence = checkpoint
3001 .as_ref()
3002 .and_then(AuthorityCheckpoint::head)
3003 .map_or(0, ControlTransition::sequence);
3004 for control in self.accepted_control_chain(domain_id)? {
3005 if control.sequence() <= baseline_sequence {
3006 continue;
3007 }
3008 for subject in control.revoked_subjects() {
3009 revocations
3010 .entry(*subject)
3011 .and_modify(|sequence| *sequence = (*sequence).max(control.sequence()))
3012 .or_insert(control.sequence());
3013 }
3014 }
3015 Ok(revocations)
3016 }
3017
3018 fn validate_control_sibling(
3019 &self,
3020 conflicting: &ControlTransition,
3021 accepted: &ControlTransition,
3022 ) -> Result<(), StoreError> {
3023 if conflicting.author() != self.controller_before(accepted)? {
3024 return Err(StoreError::InvalidAuthority(
3025 "control sibling was not signed by the predecessor controller".into(),
3026 ));
3027 }
3028 Ok(())
3029 }
3030
3031 fn controller_before(&self, control: &ControlTransition) -> Result<AuthorId, StoreError> {
3032 if let Some(checkpoint_head) = self
3033 .authority_checkpoint(control.domain_id())?
3034 .as_ref()
3035 .and_then(AuthorityCheckpoint::head)
3036 && checkpoint_head.id().map_err(authority_error)?
3037 == control.id().map_err(authority_error)?
3038 {
3039 return Ok(control.author());
3040 }
3041 let Some(predecessor) = control.predecessor() else {
3042 return self
3043 .domain_descriptor(control.domain_id())?
3044 .map(|descriptor| descriptor.owner())
3045 .ok_or_else(|| {
3046 StoreError::InvalidAuthority("domain descriptor is missing".into())
3047 });
3048 };
3049 self.control_object(control.domain_id(), predecessor)?
3050 .map(|value| value.controller())
3051 .ok_or_else(|| {
3052 StoreError::InvalidAuthority("control predecessor object is missing".into())
3053 })
3054 }
3055
3056 fn control_object(
3057 &self,
3058 domain_id: DomainId,
3059 control_id: [u8; 32],
3060 ) -> Result<Option<ControlTransition>, StoreError> {
3061 let read = self.inner.database.begin_read().map_err(database_error)?;
3062 let objects = read.open_table(CONTROL_OBJECTS).map_err(database_error)?;
3063 let key = control_object_key(domain_id, control_id);
3064 objects
3065 .get(key.as_slice())
3066 .map_err(database_error)?
3067 .map(|value| {
3068 ControlTransition::decode_canonical(value.value()).map_err(authority_error)
3069 })
3070 .transpose()
3071 }
3072
3073 fn control_slot(&self, control: &ControlTransition) -> Result<Option<[u8; 32]>, StoreError> {
3074 let read = self.inner.database.begin_read().map_err(database_error)?;
3075 let slots = read.open_table(CONTROL_SLOTS).map_err(database_error)?;
3076 slots
3077 .get(control_slot_key(control).as_slice())
3078 .map_err(database_error)?
3079 .map(|value| {
3080 value.value().try_into().map_err(|_| {
3081 StoreError::CorruptDerivedMetadata("control slot ID is invalid".into())
3082 })
3083 })
3084 .transpose()
3085 }
3086
3087 fn control_slot_is_observed(
3088 &self,
3089 control: &ControlTransition,
3090 control_id: [u8; 32],
3091 ) -> Result<bool, StoreError> {
3092 let Some(observed_id) = self.control_slot(control)? else {
3093 return Ok(false);
3094 };
3095 if observed_id == control_id {
3096 return Ok(true);
3097 }
3098 let observed = self
3099 .control_object(control.domain_id(), observed_id)?
3100 .ok_or_else(|| {
3101 StoreError::CorruptDerivedMetadata("observed control slot object is missing".into())
3102 })?;
3103 self.validate_control_sibling(control, &observed)?;
3104 self.freeze_control_fork(&observed, control)?;
3105 Ok(false)
3106 }
3107
3108 fn freeze_control_fork(
3109 &self,
3110 accepted: &ControlTransition,
3111 conflicting: &ControlTransition,
3112 ) -> Result<(), StoreError> {
3113 let domain = accepted.domain_id();
3114 let accepted_id = accepted.id().map_err(authority_error)?;
3115 let conflicting_id = conflicting.id().map_err(authority_error)?;
3116 let (first, second) = if accepted_id < conflicting_id {
3117 (accepted_id, conflicting_id)
3118 } else {
3119 (conflicting_id, accepted_id)
3120 };
3121 let bytes = conflicting.encode_canonical().map_err(authority_error)?;
3122 let mut evidence = [0_u8; 64];
3123 evidence[..32].copy_from_slice(&first);
3124 evidence[32..].copy_from_slice(&second);
3125 let write = self.inner.database.begin_write().map_err(database_error)?;
3126 write
3127 .open_table(CONTROL_OBJECTS)
3128 .map_err(database_error)?
3129 .insert(
3130 control_object_key(domain, conflicting_id).as_slice(),
3131 bytes.as_slice(),
3132 )
3133 .map_err(database_error)?;
3134 write
3135 .open_table(CONTROL_FORKS)
3136 .map_err(database_error)?
3137 .insert(domain.as_bytes().as_slice(), evidence.as_slice())
3138 .map_err(database_error)?;
3139 write.commit().map_err(database_error)?;
3140 Err(StoreError::ControlFork {
3141 domain,
3142 first,
3143 second,
3144 })
3145 }
3146
3147 fn ensure_control_not_forked(&self, domain: DomainId) -> Result<(), StoreError> {
3148 let read = self.inner.database.begin_read().map_err(database_error)?;
3149 let forks = read.open_table(CONTROL_FORKS).map_err(database_error)?;
3150 let Some(value) = forks
3151 .get(domain.as_bytes().as_slice())
3152 .map_err(database_error)?
3153 else {
3154 return Ok(());
3155 };
3156 let evidence: [u8; 64] = value.value().try_into().map_err(|_| {
3157 StoreError::CorruptDerivedMetadata("control fork evidence is invalid".into())
3158 })?;
3159 let first = evidence[..32].try_into().expect("exact evidence length");
3160 let second = evidence[32..].try_into().expect("exact evidence length");
3161 Err(StoreError::ControlFork {
3162 domain,
3163 first,
3164 second,
3165 })
3166 }
3167
3168 pub fn trust_author(&self, domain_id: DomainId, author: AuthorId) -> Result<(), StoreError> {
3170 let write = self.inner.database.begin_write().map_err(database_error)?;
3171 {
3172 let mut trusted = write.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
3173 let key = trusted_author_key(domain_id, author);
3174 trusted.insert(key.as_slice(), 0).map_err(database_error)?;
3175 }
3176 write.commit().map_err(database_error)
3177 }
3178
3179 pub fn is_author_trusted(
3181 &self,
3182 domain_id: DomainId,
3183 author: AuthorId,
3184 ) -> Result<bool, StoreError> {
3185 let read = self.inner.database.begin_read().map_err(database_error)?;
3186 let trusted = read.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
3187 let key = trusted_author_key(domain_id, author);
3188 Ok(trusted
3189 .get(key.as_slice())
3190 .map_err(database_error)?
3191 .is_some())
3192 }
3193
3194 pub fn trusted_authors(&self, domain_id: DomainId) -> Result<Vec<AuthorId>, StoreError> {
3196 let read = self.inner.database.begin_read().map_err(database_error)?;
3197 let trusted = read.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
3198 let mut authors = Vec::new();
3199 for entry in trusted.iter().map_err(database_error)? {
3200 let (key, _) = entry.map_err(database_error)?;
3201 let key = key.value();
3202 if key.starts_with(domain_id.as_bytes()) {
3203 let bytes: [u8; 32] = key[32..].try_into().map_err(|_| {
3204 StoreError::CorruptDerivedMetadata("trusted author length is invalid".into())
3205 })?;
3206 authors.push(AuthorId::from_bytes(bytes));
3207 }
3208 }
3209 authors.sort_unstable();
3210 Ok(authors)
3211 }
3212
3213 pub fn set_device_group(
3215 &self,
3216 domain_id: DomainId,
3217 label: &str,
3218 members: &[AuthorId],
3219 ) -> Result<(), StoreError> {
3220 if label.is_empty() || label.len() > 128 || members.len() > 256 {
3221 return Err(StoreError::InvalidAuthority(
3222 "device group exceeds label or member bounds".into(),
3223 ));
3224 }
3225 let mut canonical = members.to_vec();
3226 canonical.sort_unstable();
3227 canonical.dedup();
3228 let key = device_group_key(domain_id, label);
3229 let value: Vec<u8> = canonical
3230 .iter()
3231 .flat_map(|author| author.as_bytes().iter().copied())
3232 .collect();
3233 let write = self.inner.database.begin_write().map_err(database_error)?;
3234 write
3235 .open_table(DEVICE_GROUPS)
3236 .map_err(database_error)?
3237 .insert(key.as_slice(), value.as_slice())
3238 .map_err(database_error)?;
3239 write.commit().map_err(database_error)
3240 }
3241
3242 pub fn device_groups(
3244 &self,
3245 domain_id: DomainId,
3246 ) -> Result<Vec<(String, Vec<AuthorId>)>, StoreError> {
3247 let read = self.inner.database.begin_read().map_err(database_error)?;
3248 let groups = read.open_table(DEVICE_GROUPS).map_err(database_error)?;
3249 let mut result = Vec::new();
3250 for entry in groups.iter().map_err(database_error)? {
3251 let (key, value) = entry.map_err(database_error)?;
3252 let key = key.value();
3253 if !key.starts_with(domain_id.as_bytes()) {
3254 continue;
3255 }
3256 let label = std::str::from_utf8(&key[32..])
3257 .map_err(|_| StoreError::CorruptDerivedMetadata("group label is invalid".into()))?
3258 .to_owned();
3259 let chunks = value.value().chunks_exact(32);
3260 if !chunks.remainder().is_empty() {
3261 return Err(StoreError::CorruptDerivedMetadata(
3262 "group member data is invalid".into(),
3263 ));
3264 }
3265 let members = chunks
3266 .map(|chunk| {
3267 let bytes: [u8; 32] = chunk.try_into().expect("exact chunk length");
3268 AuthorId::from_bytes(bytes)
3269 })
3270 .collect();
3271 result.push((label, members));
3272 }
3273 result.sort_by(|left, right| left.0.cmp(&right.0));
3274 Ok(result)
3275 }
3276
3277 pub fn consume_invitation(&self, invitation_id: [u8; 32]) -> Result<bool, StoreError> {
3279 let write = self.inner.database.begin_write().map_err(database_error)?;
3280 let fresh = {
3281 let mut consumed = write
3282 .open_table(CONSUMED_INVITATIONS)
3283 .map_err(database_error)?;
3284 if consumed
3285 .get(invitation_id.as_slice())
3286 .map_err(database_error)?
3287 .is_some()
3288 {
3289 false
3290 } else {
3291 consumed
3292 .insert(invitation_id.as_slice(), 0)
3293 .map_err(database_error)?;
3294 true
3295 }
3296 };
3297 write.commit().map_err(database_error)?;
3298 Ok(fresh)
3299 }
3300}
3301
3302fn dependency_is_accepted(
3303 commits: &impl redb::ReadableTable<&'static [u8], &'static [u8]>,
3304 coverage: &impl redb::ReadableTable<&'static [u8], &'static [u8]>,
3305 quarantined: &impl redb::ReadableTable<&'static [u8], &'static [u8]>,
3306 domain_id: DomainId,
3307 dependency: CommitId,
3308) -> Result<bool, StoreError> {
3309 let key = dependency.as_bytes().as_slice();
3310 let stored = commits
3311 .get(key)
3312 .map_err(database_error)?
3313 .is_some_and(|metadata| metadata.value().starts_with(domain_id.as_bytes()))
3314 && quarantined.get(key).map_err(database_error)?.is_none();
3315 let snapshotted = coverage
3316 .get(key)
3317 .map_err(database_error)?
3318 .is_some_and(|domain| domain.value() == domain_id.as_bytes());
3319 Ok(stored || snapshotted)
3320}
3321
3322fn publish_frontier(
3323 write: &redb::WriteTransaction,
3324 domain_id: DomainId,
3325 envelope: &CommitEnvelope,
3326 commit_id: CommitId,
3327) -> Result<(), StoreError> {
3328 let mut frontier = write.open_table(FRONTIER).map_err(database_error)?;
3329 for dependency in envelope.header().dependencies() {
3330 let dependency_key = frontier_key(domain_id, *dependency);
3331 frontier
3332 .remove(dependency_key.as_slice())
3333 .map_err(database_error)?;
3334 }
3335 let commit_key = frontier_key(domain_id, commit_id);
3336 frontier
3337 .insert(commit_key.as_slice(), 0)
3338 .map_err(database_error)?;
3339 Ok(())
3340}
3341
3342fn apply_reconciliation_update(
3343 write: &redb::WriteTransaction,
3344 envelope: &CommitEnvelope,
3345 context: &VersionVector,
3346 record_heads: &[RecordHeadSet],
3347) -> Result<(), StoreError> {
3348 let own_dot = iroh_db_core::Dot::new(
3349 envelope.header().author(),
3350 envelope.header().author_sequence(),
3351 0,
3352 );
3353 if !context.covers(&own_dot) {
3354 return Err(StoreError::CorruptDerivedMetadata(
3355 "commit context does not include its own author sequence".into(),
3356 ));
3357 }
3358 let domain_id = envelope.header().domain_id();
3359 {
3360 let key = commit_context_key(domain_id, envelope.commit_id());
3361 let encoded = encode_version_vector(context)?;
3362 write
3363 .open_table(COMMIT_CONTEXTS)
3364 .map_err(database_error)?
3365 .insert(key.as_slice(), encoded.as_slice())
3366 .map_err(database_error)?;
3367 }
3368 {
3369 let mut table = write.open_table(RECORD_HEADS).map_err(database_error)?;
3370 for replacement in record_heads {
3371 let prefix = record_head_prefix(
3372 domain_id,
3373 replacement.collection_id(),
3374 replacement.record_id(),
3375 )?;
3376 for key in table_keys_with_prefix(&table, &prefix)? {
3377 table.remove(key.as_slice()).map_err(database_error)?;
3378 }
3379 for head in replacement.heads() {
3380 let key = record_head_key(domain_id, head)?;
3381 let value = encode_record_head(head)?;
3382 table
3383 .insert(key.as_slice(), value.as_slice())
3384 .map_err(database_error)?;
3385 }
3386 }
3387 }
3388 write
3389 .open_table(RECONCILIATION_READY)
3390 .map_err(database_error)?
3391 .insert(domain_id.as_bytes().as_slice(), 1)
3392 .map_err(database_error)?;
3393 Ok(())
3394}
3395
3396fn mark_reconciliation_stale(
3397 write: &redb::WriteTransaction,
3398 domain_id: DomainId,
3399) -> Result<(), StoreError> {
3400 write
3401 .open_table(RECONCILIATION_READY)
3402 .map_err(database_error)?
3403 .remove(domain_id.as_bytes().as_slice())
3404 .map_err(database_error)?;
3405 Ok(())
3406}
3407
3408fn validate_snapshot_metadata(
3409 records: &[MaterializedRecord],
3410 covered_commits: &[CommitId],
3411 frontier_commits: &[CommitId],
3412 author_sequences: &[(AuthorId, u64)],
3413) -> Result<(), StoreError> {
3414 let records_are_ordered = records.windows(2).all(|pair| {
3415 (pair[0].schema().collection_id(), pair[0].record_id())
3416 < (pair[1].schema().collection_id(), pair[1].record_id())
3417 });
3418 if !records_are_ordered || records.iter().any(|record| record.state().is_none()) {
3419 return Err(StoreError::InvalidSnapshotMetadata(
3420 "materialized records are not unique canonical upserts".into(),
3421 ));
3422 }
3423 if !strictly_sorted(covered_commits) || !strictly_sorted(frontier_commits) {
3424 return Err(StoreError::InvalidSnapshotMetadata(
3425 "commit coverage and frontier must be strictly ordered".into(),
3426 ));
3427 }
3428 if !author_sequences
3429 .windows(2)
3430 .all(|pair| pair[0].0 < pair[1].0)
3431 {
3432 return Err(StoreError::InvalidSnapshotMetadata(
3433 "author sequence entries must be strictly ordered".into(),
3434 ));
3435 }
3436 if frontier_commits
3437 .iter()
3438 .any(|commit_id| covered_commits.binary_search(commit_id).is_err())
3439 {
3440 return Err(StoreError::InvalidSnapshotMetadata(
3441 "snapshot frontier is not covered by the snapshot".into(),
3442 ));
3443 }
3444 if covered_commits.is_empty() != frontier_commits.is_empty()
3445 || covered_commits.is_empty() != author_sequences.is_empty()
3446 {
3447 return Err(StoreError::InvalidSnapshotMetadata(
3448 "empty history must have an empty frontier and version vector".into(),
3449 ));
3450 }
3451 Ok(())
3452}
3453
3454fn strictly_sorted<T: Ord>(values: &[T]) -> bool {
3455 values.windows(2).all(|pair| pair[0] < pair[1])
3456}
3457
3458fn initialize_tables(database: &Database) -> Result<(), StoreError> {
3459 let write = database.begin_write().map_err(database_error)?;
3460 write.open_table(COMMIT_METADATA).map_err(database_error)?;
3461 write.open_table(COMMIT_CONTEXTS).map_err(database_error)?;
3462 write.open_table(RECORD_HEADS).map_err(database_error)?;
3463 write
3464 .open_table(RECONCILIATION_READY)
3465 .map_err(database_error)?;
3466 write
3467 .open_table(SNAPSHOT_COVERAGE)
3468 .map_err(database_error)?;
3469 write
3470 .open_table(SNAPSHOT_BASELINES)
3471 .map_err(database_error)?;
3472 write.open_table(BLOB_OWNERSHIP).map_err(database_error)?;
3473 write.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
3474 write.open_table(AUTHOR_COMMITS).map_err(database_error)?;
3475 write.open_table(AUTHOR_HEADS).map_err(database_error)?;
3476 write.open_table(FRONTIER).map_err(database_error)?;
3477 write
3478 .open_table(QUARANTINED_COMMITS)
3479 .map_err(database_error)?;
3480 write
3481 .open_table(QUARANTINED_AUTHORS)
3482 .map_err(database_error)?;
3483 write.open_table(SCHEMAS).map_err(database_error)?;
3484 write.open_table(RECORDS).map_err(database_error)?;
3485 write
3486 .open_table(RECORD_INDEX_KEYS)
3487 .map_err(database_error)?;
3488 write.open_table(SECONDARY_INDEX).map_err(database_error)?;
3489 write.open_table(TRUSTED_AUTHORS).map_err(database_error)?;
3490 write
3491 .open_table(DOMAIN_DESCRIPTORS)
3492 .map_err(database_error)?;
3493 write.open_table(CAPABILITIES).map_err(database_error)?;
3494 write
3495 .open_table(ACTIVE_CAPABILITIES)
3496 .map_err(database_error)?;
3497 write
3498 .open_table(REVOKED_CAPABILITIES)
3499 .map_err(database_error)?;
3500 write
3501 .open_table(SUBJECT_REVOCATIONS)
3502 .map_err(database_error)?;
3503 write
3504 .open_table(AUTHORITY_CHECKPOINTS)
3505 .map_err(database_error)?;
3506 write.open_table(CONTROL_HEADS).map_err(database_error)?;
3507 write.open_table(CONTROL_OBJECTS).map_err(database_error)?;
3508 write.open_table(CONTROL_SLOTS).map_err(database_error)?;
3509 write.open_table(CONTROL_FORKS).map_err(database_error)?;
3510 write.open_table(DEVICE_GROUPS).map_err(database_error)?;
3511 write
3512 .open_table(CONSUMED_INVITATIONS)
3513 .map_err(database_error)?;
3514 write.open_table(ACTIVE_SCHEMAS).map_err(database_error)?;
3515 write.open_table(MIGRATION_STAGES).map_err(database_error)?;
3516 write
3517 .open_table(MIGRATION_ACTIVATIONS)
3518 .map_err(database_error)?;
3519 write
3520 .open_table(MIGRATION_CHECKPOINTS)
3521 .map_err(database_error)?;
3522 write.open_table(MIGRATION_FORKS).map_err(database_error)?;
3523 write
3524 .open_table(OPERATIONAL_METADATA)
3525 .map_err(database_error)?;
3526 write.commit().map_err(database_error)
3527}
3528
3529fn prepare_store_format(root: &Path) -> Result<(), StoreError> {
3530 let format_path = root.join("FORMAT");
3531 if !format_path.exists() {
3532 return write_format_marker(root, FORMAT_CONTENTS);
3533 }
3534 match std::fs::read_to_string(&format_path)
3535 .map_err(io_error)?
3536 .as_str()
3537 {
3538 FORMAT_CONTENTS => Ok(()),
3539 FORMAT_V5 => shadow_migrate_v5(root),
3540 _ => Err(StoreError::UnsupportedFormat),
3541 }
3542}
3543
3544fn shadow_migrate_v5(root: &Path) -> Result<(), StoreError> {
3545 let legacy = root.join("meta.redb");
3546 let preserved = root.join("meta-v5.redb");
3547 let candidate = root.join("tmp/meta-v6.redb");
3548 let live = root.join("meta-v6.redb");
3549 if !preserved.exists() {
3550 if !legacy.is_file() {
3551 return Err(StoreError::Io("format v5 metadata file is missing".into()));
3552 }
3553 copy_and_sync(&legacy, &preserved)?;
3554 }
3555 if candidate.exists() {
3556 std::fs::remove_file(&candidate).map_err(io_error)?;
3557 }
3558 copy_and_sync(&preserved, &candidate)?;
3559 {
3560 let mut database = Database::create(&candidate).map_err(database_error)?;
3561 initialize_tables(&database)?;
3562 if !database.check_integrity().map_err(database_error)? {
3563 return Err(StoreError::Database(
3564 "v6 shadow metadata required integrity repair".into(),
3565 ));
3566 }
3567 }
3568 if live.exists() {
3569 std::fs::remove_file(&live).map_err(io_error)?;
3570 }
3571 std::fs::rename(&candidate, &live).map_err(io_error)?;
3572 sync_directory(root)?;
3573 write_format_marker(root, FORMAT_CONTENTS)
3574}
3575
3576fn copy_and_sync(source: &Path, destination: &Path) -> Result<(), StoreError> {
3577 std::fs::copy(source, destination).map_err(io_error)?;
3578 OpenOptions::new()
3579 .read(true)
3580 .write(true)
3581 .open(destination)
3582 .map_err(io_error)?
3583 .sync_all()
3584 .map_err(io_error)
3585}
3586
3587fn write_format_marker(root: &Path, contents: &str) -> Result<(), StoreError> {
3588 let candidate = root.join("tmp/FORMAT.next");
3589 std::fs::write(&candidate, contents).map_err(io_error)?;
3590 OpenOptions::new()
3591 .read(true)
3592 .write(true)
3593 .open(&candidate)
3594 .map_err(io_error)?
3595 .sync_all()
3596 .map_err(io_error)?;
3597 std::fs::rename(candidate, root.join("FORMAT")).map_err(io_error)?;
3598 sync_directory(root)
3599}
3600
3601fn sync_directory(path: &Path) -> Result<(), StoreError> {
3602 File::open(path)
3603 .map_err(io_error)?
3604 .sync_all()
3605 .map_err(io_error)
3606}
3607
3608fn trusted_author_key(domain_id: DomainId, author: AuthorId) -> [u8; 64] {
3609 let mut key = [0_u8; 64];
3610 key[..32].copy_from_slice(domain_id.as_bytes());
3611 key[32..].copy_from_slice(author.as_bytes());
3612 key
3613}
3614
3615fn active_capability_key(domain_id: DomainId, subject: AuthorId) -> [u8; 64] {
3616 let mut key = [0_u8; 64];
3617 key[..32].copy_from_slice(domain_id.as_bytes());
3618 key[32..].copy_from_slice(subject.as_bytes());
3619 key
3620}
3621
3622fn subject_revocation_key(domain_id: DomainId, subject: AuthorId) -> [u8; 64] {
3623 active_capability_key(domain_id, subject)
3624}
3625
3626fn control_object_key(domain_id: DomainId, control_id: [u8; 32]) -> [u8; 64] {
3627 let mut key = [0_u8; 64];
3628 key[..32].copy_from_slice(domain_id.as_bytes());
3629 key[32..].copy_from_slice(&control_id);
3630 key
3631}
3632
3633fn control_slot_key(control: &ControlTransition) -> [u8; 81] {
3634 let mut key = [0_u8; 81];
3635 key[..32].copy_from_slice(control.domain_id().as_bytes());
3636 key[32..40].copy_from_slice(&control.sequence().to_be_bytes());
3637 if let Some(predecessor) = control.predecessor() {
3638 key[40] = 1;
3639 key[41..73].copy_from_slice(&predecessor);
3640 }
3641 key[73..].copy_from_slice(&control.previous_epoch().to_be_bytes());
3642 key
3643}
3644
3645fn insert_control_observation(
3646 write: &redb::WriteTransaction,
3647 control: &ControlTransition,
3648 control_id: [u8; 32],
3649) -> Result<(), StoreError> {
3650 let bytes = control.encode_canonical().map_err(authority_error)?;
3651 write
3652 .open_table(CONTROL_OBJECTS)
3653 .map_err(database_error)?
3654 .insert(
3655 control_object_key(control.domain_id(), control_id).as_slice(),
3656 bytes.as_slice(),
3657 )
3658 .map_err(database_error)?;
3659 write
3660 .open_table(CONTROL_SLOTS)
3661 .map_err(database_error)?
3662 .insert(control_slot_key(control).as_slice(), control_id.as_slice())
3663 .map_err(database_error)?;
3664 Ok(())
3665}
3666
3667fn apply_control_revocations(
3668 write: &redb::WriteTransaction,
3669 control: &ControlTransition,
3670) -> Result<(), StoreError> {
3671 for subject in control.revoked_subjects() {
3672 let subject_key = subject_revocation_key(control.domain_id(), *subject);
3673 write
3674 .open_table(SUBJECT_REVOCATIONS)
3675 .map_err(database_error)?
3676 .insert(subject_key.as_slice(), control.sequence())
3677 .map_err(database_error)?;
3678 let active_key = active_capability_key(control.domain_id(), *subject);
3679 let capability_id = {
3680 let mut active = write
3681 .open_table(ACTIVE_CAPABILITIES)
3682 .map_err(database_error)?;
3683 let id = active
3684 .get(active_key.as_slice())
3685 .map_err(database_error)?
3686 .map(|value| value.value().to_vec());
3687 active
3688 .remove(active_key.as_slice())
3689 .map_err(database_error)?;
3690 id
3691 };
3692 if let Some(capability_id) = capability_id {
3693 write
3694 .open_table(REVOKED_CAPABILITIES)
3695 .map_err(database_error)?
3696 .insert(capability_id.as_slice(), control.sequence())
3697 .map_err(database_error)?;
3698 }
3699 let trusted_key = trusted_author_key(control.domain_id(), *subject);
3700 write
3701 .open_table(TRUSTED_AUTHORS)
3702 .map_err(database_error)?
3703 .remove(trusted_key.as_slice())
3704 .map_err(database_error)?;
3705 }
3706 Ok(())
3707}
3708
3709fn device_group_key(domain_id: DomainId, label: &str) -> Vec<u8> {
3710 let mut key = Vec::with_capacity(32 + label.len());
3711 key.extend_from_slice(domain_id.as_bytes());
3712 key.extend_from_slice(label.as_bytes());
3713 key
3714}
3715
3716fn insert_capability_tables(
3717 write: &redb::WriteTransaction,
3718 capability: &CapabilityCertificate,
3719 capability_id: CapabilityId,
3720 bytes: &[u8],
3721) -> Result<(), StoreError> {
3722 insert_capability_record(write, capability_id, bytes)?;
3723 set_active_capability(write, capability, capability_id)
3724}
3725
3726fn insert_capability_record(
3727 write: &redb::WriteTransaction,
3728 capability_id: CapabilityId,
3729 bytes: &[u8],
3730) -> Result<(), StoreError> {
3731 {
3732 let mut capabilities = write.open_table(CAPABILITIES).map_err(database_error)?;
3733 let existing = capabilities
3734 .get(capability_id.as_bytes().as_slice())
3735 .map_err(database_error)?
3736 .map(|value| value.value().to_vec());
3737 if existing.as_deref().is_some_and(|stored| stored != bytes) {
3738 return Err(StoreError::InvalidAuthority(
3739 "capability ID collision".into(),
3740 ));
3741 }
3742 capabilities
3743 .insert(capability_id.as_bytes().as_slice(), bytes)
3744 .map_err(database_error)?;
3745 }
3746 Ok(())
3747}
3748
3749fn set_active_capability(
3750 write: &redb::WriteTransaction,
3751 capability: &CapabilityCertificate,
3752 capability_id: CapabilityId,
3753) -> Result<(), StoreError> {
3754 let mut active = write
3755 .open_table(ACTIVE_CAPABILITIES)
3756 .map_err(database_error)?;
3757 let key = active_capability_key(capability.domain_id(), capability.subject());
3758 active
3759 .insert(key.as_slice(), capability_id.as_bytes().as_slice())
3760 .map_err(database_error)?;
3761 Ok(())
3762}
3763
3764#[allow(clippy::needless_pass_by_value)]
3765fn authority_error(error: impl std::fmt::Display) -> StoreError {
3766 StoreError::InvalidAuthority(error.to_string())
3767}
3768
3769fn apply_materialized_records(
3770 write: &redb::WriteTransaction,
3771 domain_id: DomainId,
3772 mutations: &[MaterializedRecord],
3773) -> Result<(), StoreError> {
3774 for mutation in mutations {
3775 let collection_id = mutation.schema.collection_id();
3776 let schema_key = schema_key(domain_id, collection_id, mutation.schema.version());
3777 let schema_bytes = mutation.schema.encode_canonical();
3778 {
3779 let mut schemas = write.open_table(SCHEMAS).map_err(database_error)?;
3780 let existing = schemas
3781 .get(schema_key.as_slice())
3782 .map_err(database_error)?
3783 .map(|value| value.value().to_vec());
3784 if let Some(existing) = existing {
3785 if existing != schema_bytes {
3786 return Err(StoreError::SchemaMismatch {
3787 collection_id,
3788 version: mutation.schema.version(),
3789 });
3790 }
3791 } else {
3792 schemas
3793 .insert(schema_key.as_slice(), schema_bytes.as_slice())
3794 .map_err(database_error)?;
3795 }
3796 }
3797
3798 for index in &mutation.indexes {
3799 let declared = mutation
3800 .schema
3801 .fields()
3802 .iter()
3803 .any(|field| field.id() == index.field_id && field.is_indexed());
3804 if !declared {
3805 return Err(StoreError::InvalidIndex {
3806 collection_id,
3807 field_id: index.field_id,
3808 });
3809 }
3810 }
3811
3812 let record_key = record_key(domain_id, collection_id, &mutation.record_id);
3813 remove_record_indexes(write, &record_key)?;
3814 let mut records = write.open_table(RECORDS).map_err(database_error)?;
3815 match &mutation.state {
3816 Some(state) => {
3817 records
3818 .insert(record_key.as_slice(), state.as_slice())
3819 .map_err(database_error)?;
3820 let index_keys = mutation
3821 .indexes
3822 .iter()
3823 .map(|index| {
3824 index_key(
3825 domain_id,
3826 collection_id,
3827 index.field_id,
3828 &index.value,
3829 &mutation.record_id,
3830 )
3831 })
3832 .collect::<Result<Vec<_>, _>>()?;
3833 let mut secondary = write.open_table(SECONDARY_INDEX).map_err(database_error)?;
3834 for key in &index_keys {
3835 secondary
3836 .insert(key.as_slice(), 0)
3837 .map_err(database_error)?;
3838 }
3839 let encoded_keys = encode_key_list(&index_keys)?;
3840 let mut record_indexes = write
3841 .open_table(RECORD_INDEX_KEYS)
3842 .map_err(database_error)?;
3843 record_indexes
3844 .insert(record_key.as_slice(), encoded_keys.as_slice())
3845 .map_err(database_error)?;
3846 }
3847 None => {
3848 records
3849 .remove(record_key.as_slice())
3850 .map_err(database_error)?;
3851 }
3852 }
3853 }
3854 Ok(())
3855}
3856
3857fn apply_schema_stages(
3858 write: &redb::WriteTransaction,
3859 domain_id: DomainId,
3860 stages: &[SchemaStage],
3861) -> Result<Option<CollectionId>, StoreError> {
3862 if stages.is_empty() {
3863 return Ok(None);
3864 }
3865 staging_digest(stages)?;
3866 let mut table = write.open_table(MIGRATION_STAGES).map_err(database_error)?;
3867 for stage in stages {
3868 let target = MaterializedRecord::decode_canonical(stage.target_record())
3869 .map_err(|error| StoreError::InvalidMigration(error.to_string()))?;
3870 if target.state().is_none()
3871 || target.record_id() != stage.source_record_id()
3872 || target.schema().schema_id() != stage.to().schema_id()
3873 {
3874 return Err(StoreError::InvalidMigration(
3875 "stage target does not match its target schema and source key".into(),
3876 ));
3877 }
3878 let key = migration_stage_key(domain_id, stage.migration_id(), stage.source_record_id());
3879 let bytes = stage.encode_canonical();
3880 let existing = table
3881 .get(key.as_slice())
3882 .map_err(database_error)?
3883 .map(|value| value.value().to_vec());
3884 if let Some(existing) = existing.filter(|existing| existing != &bytes) {
3885 let evidence = minicbor::to_vec((existing, bytes.clone()))
3886 .map_err(|error| StoreError::InvalidMigration(error.to_string()))?;
3887 drop(table);
3888 write
3889 .open_table(MIGRATION_FORKS)
3890 .map_err(database_error)?
3891 .insert(
3892 active_schema_key(domain_id, stage.to().collection_id()).as_slice(),
3893 evidence.as_slice(),
3894 )
3895 .map_err(database_error)?;
3896 return Ok(Some(stage.to().collection_id()));
3897 }
3898 table
3899 .insert(key.as_slice(), bytes.as_slice())
3900 .map_err(database_error)?;
3901 }
3902 Ok(None)
3903}
3904
3905fn apply_schema_activation(
3906 write: &redb::WriteTransaction,
3907 domain_id: DomainId,
3908 activation: &SchemaActivation,
3909) -> Result<bool, StoreError> {
3910 let collection_id = activation.to().collection_id();
3911 let history_key =
3912 migration_activation_key(domain_id, collection_id, activation.from().schema_id());
3913 let activation_bytes = activation.encode_canonical();
3914 let conflict = {
3915 let history = write
3916 .open_table(MIGRATION_ACTIVATIONS)
3917 .map_err(database_error)?;
3918 history
3919 .get(history_key.as_slice())
3920 .map_err(database_error)?
3921 .map(|existing| existing.value().to_vec())
3922 .filter(|existing| existing != &activation_bytes)
3923 };
3924 if let Some(existing) = conflict {
3925 let evidence = minicbor::to_vec((existing, activation_bytes.clone()))
3926 .map_err(|error| StoreError::InvalidMigration(error.to_string()))?;
3927 write
3928 .open_table(MIGRATION_FORKS)
3929 .map_err(database_error)?
3930 .insert(
3931 active_schema_key(domain_id, collection_id).as_slice(),
3932 evidence.as_slice(),
3933 )
3934 .map_err(database_error)?;
3935 return Ok(true);
3936 }
3937
3938 let stages = {
3939 let table = write.open_table(MIGRATION_STAGES).map_err(database_error)?;
3940 load_schema_stages(&table, domain_id, activation.migration_id())?
3941 };
3942 if u64::try_from(stages.len()).ok() != Some(activation.staged_records())
3943 || staging_digest(&stages)? != activation.staging_digest()
3944 || stages.iter().any(|stage| {
3945 stage.from().schema_id() != activation.from().schema_id()
3946 || stage.to().schema_id() != activation.to().schema_id()
3947 })
3948 {
3949 return Err(StoreError::InvalidMigration(
3950 "activation does not commit to the complete durable staging set".into(),
3951 ));
3952 }
3953
3954 let mut targets = Vec::with_capacity(stages.len());
3955 for stage in &stages {
3956 targets.push(
3957 MaterializedRecord::decode_canonical(stage.target_record())
3958 .map_err(|error| StoreError::InvalidMigration(error.to_string()))?,
3959 );
3960 }
3961 clear_materialized_collection(write, domain_id, collection_id)?;
3962 register_schema(write, domain_id, activation.from())?;
3963 register_schema(write, domain_id, activation.to())?;
3964 apply_materialized_records(write, domain_id, &targets)?;
3965
3966 write
3967 .open_table(MIGRATION_ACTIVATIONS)
3968 .map_err(database_error)?
3969 .insert(history_key.as_slice(), activation_bytes.as_slice())
3970 .map_err(database_error)?;
3971 write
3972 .open_table(ACTIVE_SCHEMAS)
3973 .map_err(database_error)?
3974 .insert(
3975 active_schema_key(domain_id, collection_id).as_slice(),
3976 activation_bytes.as_slice(),
3977 )
3978 .map_err(database_error)?;
3979 Ok(false)
3980}
3981
3982fn load_schema_stages(
3983 table: &impl redb::ReadableTable<&'static [u8], &'static [u8]>,
3984 domain_id: DomainId,
3985 migration_id: MigrationId,
3986) -> Result<Vec<SchemaStage>, StoreError> {
3987 let prefix = migration_stage_prefix(domain_id, migration_id);
3988 let mut stages = Vec::new();
3989 for entry in table.iter().map_err(database_error)? {
3990 let (key, value) = entry.map_err(database_error)?;
3991 if key.value().starts_with(&prefix) {
3992 stages.push(SchemaStage::decode_canonical(value.value())?);
3993 }
3994 }
3995 stages.sort_unstable_by(|left, right| left.source_record_id().cmp(right.source_record_id()));
3996 Ok(stages)
3997}
3998
3999fn register_schema(
4000 write: &redb::WriteTransaction,
4001 domain_id: DomainId,
4002 schema: &SchemaDescriptor,
4003) -> Result<(), StoreError> {
4004 let key = schema_key(domain_id, schema.collection_id(), schema.version());
4005 let bytes = schema.encode_canonical();
4006 let mut schemas = write.open_table(SCHEMAS).map_err(database_error)?;
4007 if let Some(existing) = schemas.get(key.as_slice()).map_err(database_error)?
4008 && existing.value() != bytes
4009 {
4010 return Err(StoreError::SchemaMismatch {
4011 collection_id: schema.collection_id(),
4012 version: schema.version(),
4013 });
4014 }
4015 schemas
4016 .insert(key.as_slice(), bytes.as_slice())
4017 .map_err(database_error)?;
4018 Ok(())
4019}
4020
4021fn clear_materialized_collection(
4022 write: &redb::WriteTransaction,
4023 domain_id: DomainId,
4024 collection_id: CollectionId,
4025) -> Result<(), StoreError> {
4026 let prefix = record_prefix(domain_id, collection_id);
4027 {
4028 let mut table = write.open_table(RECORDS).map_err(database_error)?;
4029 for key in table_keys_with_prefix(&table, &prefix)? {
4030 table.remove(key.as_slice()).map_err(database_error)?;
4031 }
4032 }
4033 {
4034 let mut table = write
4035 .open_table(RECORD_INDEX_KEYS)
4036 .map_err(database_error)?;
4037 for key in table_keys_with_prefix(&table, &prefix)? {
4038 table.remove(key.as_slice()).map_err(database_error)?;
4039 }
4040 }
4041 {
4042 let mut table = write.open_table(SECONDARY_INDEX).map_err(database_error)?;
4043 let mut keys = Vec::new();
4044 for entry in table.iter().map_err(database_error)? {
4045 let (key, _) = entry.map_err(database_error)?;
4046 if key.value().starts_with(&prefix) {
4047 keys.push(key.value().to_vec());
4048 }
4049 }
4050 for key in keys {
4051 table.remove(key.as_slice()).map_err(database_error)?;
4052 }
4053 }
4054 Ok(())
4055}
4056
4057fn install_snapshot_activations(
4058 write: &redb::WriteTransaction,
4059 domain_id: DomainId,
4060 records: &[MaterializedRecord],
4061 activations: &[SchemaActivation],
4062) -> Result<(), StoreError> {
4063 if activations
4064 .windows(2)
4065 .any(|pair| pair[0].to().collection_id() >= pair[1].to().collection_id())
4066 {
4067 return Err(StoreError::InvalidSnapshotMetadata(
4068 "schema activations are not unique and ordered".into(),
4069 ));
4070 }
4071 for activation in activations {
4072 let collection_id = activation.to().collection_id();
4073 if records.iter().any(|record| {
4074 record.schema().collection_id() == collection_id
4075 && record.schema().schema_id() != activation.to().schema_id()
4076 }) {
4077 return Err(StoreError::InvalidSnapshotMetadata(
4078 "record schema disagrees with active migration".into(),
4079 ));
4080 }
4081 register_schema(write, domain_id, activation.from())?;
4082 register_schema(write, domain_id, activation.to())?;
4083 let bytes = activation.encode_canonical();
4084 write
4085 .open_table(ACTIVE_SCHEMAS)
4086 .map_err(database_error)?
4087 .insert(
4088 active_schema_key(domain_id, collection_id).as_slice(),
4089 bytes.as_slice(),
4090 )
4091 .map_err(database_error)?;
4092 write
4093 .open_table(MIGRATION_ACTIVATIONS)
4094 .map_err(database_error)?
4095 .insert(
4096 migration_activation_key(domain_id, collection_id, activation.from().schema_id())
4097 .as_slice(),
4098 bytes.as_slice(),
4099 )
4100 .map_err(database_error)?;
4101 }
4102 Ok(())
4103}
4104
4105fn clear_domain_migration_metadata(
4106 write: &redb::WriteTransaction,
4107 domain_id: DomainId,
4108) -> Result<(), StoreError> {
4109 for table_definition in [
4110 ACTIVE_SCHEMAS,
4111 MIGRATION_STAGES,
4112 MIGRATION_ACTIVATIONS,
4113 MIGRATION_FORKS,
4114 ] {
4115 let mut table = write.open_table(table_definition).map_err(database_error)?;
4116 for key in table_keys_with_prefix(&table, domain_id.as_bytes())? {
4117 table.remove(key.as_slice()).map_err(database_error)?;
4118 }
4119 }
4120 let mut checkpoints = write
4121 .open_table(MIGRATION_CHECKPOINTS)
4122 .map_err(database_error)?;
4123 for key in table_keys_with_prefix(&checkpoints, domain_id.as_bytes())? {
4124 checkpoints.remove(key.as_slice()).map_err(database_error)?;
4125 }
4126 Ok(())
4127}
4128
4129fn remove_record_indexes(
4130 write: &redb::WriteTransaction,
4131 record_key: &[u8],
4132) -> Result<(), StoreError> {
4133 let old_keys = {
4134 let record_indexes = write
4135 .open_table(RECORD_INDEX_KEYS)
4136 .map_err(database_error)?;
4137 record_indexes
4138 .get(record_key)
4139 .map_err(database_error)?
4140 .map(|value| value.value().to_vec())
4141 };
4142 if let Some(old_keys) = old_keys {
4143 let old_keys = decode_key_list(&old_keys)?;
4144 let mut secondary = write.open_table(SECONDARY_INDEX).map_err(database_error)?;
4145 for key in old_keys {
4146 secondary.remove(key.as_slice()).map_err(database_error)?;
4147 }
4148 let mut record_indexes = write
4149 .open_table(RECORD_INDEX_KEYS)
4150 .map_err(database_error)?;
4151 record_indexes.remove(record_key).map_err(database_error)?;
4152 }
4153 Ok(())
4154}
4155
4156fn clear_materialized_domain(
4157 write: &redb::WriteTransaction,
4158 domain_id: DomainId,
4159) -> Result<(), StoreError> {
4160 {
4161 let mut table = write.open_table(SCHEMAS).map_err(database_error)?;
4162 let keys = table_keys_with_prefix(&table, domain_id.as_bytes())?;
4163 for key in keys {
4164 table.remove(key.as_slice()).map_err(database_error)?;
4165 }
4166 }
4167 {
4168 let mut table = write.open_table(RECORDS).map_err(database_error)?;
4169 let keys = table_keys_with_prefix(&table, domain_id.as_bytes())?;
4170 for key in keys {
4171 table.remove(key.as_slice()).map_err(database_error)?;
4172 }
4173 }
4174 {
4175 let mut table = write
4176 .open_table(RECORD_INDEX_KEYS)
4177 .map_err(database_error)?;
4178 let keys = table_keys_with_prefix(&table, domain_id.as_bytes())?;
4179 for key in keys {
4180 table.remove(key.as_slice()).map_err(database_error)?;
4181 }
4182 }
4183 {
4184 let mut table = write.open_table(SECONDARY_INDEX).map_err(database_error)?;
4185 let mut keys = Vec::new();
4186 for entry in table.iter().map_err(database_error)? {
4187 let (key, _) = entry.map_err(database_error)?;
4188 if key.value().starts_with(domain_id.as_bytes()) {
4189 keys.push(key.value().to_vec());
4190 }
4191 }
4192 for key in keys {
4193 table.remove(key.as_slice()).map_err(database_error)?;
4194 }
4195 }
4196 Ok(())
4197}
4198
4199fn clear_reconciliation_domain(
4200 write: &redb::WriteTransaction,
4201 domain_id: DomainId,
4202) -> Result<(), StoreError> {
4203 {
4204 let mut table = write.open_table(COMMIT_CONTEXTS).map_err(database_error)?;
4205 for key in table_keys_with_prefix(&table, domain_id.as_bytes())? {
4206 table.remove(key.as_slice()).map_err(database_error)?;
4207 }
4208 }
4209 {
4210 let mut table = write.open_table(RECORD_HEADS).map_err(database_error)?;
4211 for key in table_keys_with_prefix(&table, domain_id.as_bytes())? {
4212 table.remove(key.as_slice()).map_err(database_error)?;
4213 }
4214 }
4215 mark_reconciliation_stale(write, domain_id)
4216}
4217
4218fn clear_domain_snapshot_metadata(
4219 write: &redb::WriteTransaction,
4220 domain_id: DomainId,
4221) -> Result<(), StoreError> {
4222 write
4223 .open_table(SNAPSHOT_BASELINES)
4224 .map_err(database_error)?
4225 .remove(domain_id.as_bytes().as_slice())
4226 .map_err(database_error)?;
4227 {
4228 let mut coverage = write
4229 .open_table(SNAPSHOT_COVERAGE)
4230 .map_err(database_error)?;
4231 let mut keys = Vec::new();
4232 for entry in coverage.iter().map_err(database_error)? {
4233 let (key, value) = entry.map_err(database_error)?;
4234 if value.value() == domain_id.as_bytes() {
4235 keys.push(key.value().to_vec());
4236 }
4237 }
4238 for key in keys {
4239 coverage.remove(key.as_slice()).map_err(database_error)?;
4240 }
4241 }
4242 {
4243 let mut frontier = write.open_table(FRONTIER).map_err(database_error)?;
4244 let keys = table_u8_keys_with_prefix(&frontier, domain_id.as_bytes())?;
4245 for key in keys {
4246 frontier.remove(key.as_slice()).map_err(database_error)?;
4247 }
4248 }
4249 {
4250 let mut sequences = write.open_table(AUTHOR_SEQUENCES).map_err(database_error)?;
4251 let mut keys = Vec::new();
4252 for entry in sequences.iter().map_err(database_error)? {
4253 let (key, _) = entry.map_err(database_error)?;
4254 if key.value().starts_with(domain_id.as_bytes()) {
4255 keys.push(key.value().to_vec());
4256 }
4257 }
4258 for key in keys {
4259 sequences.remove(key.as_slice()).map_err(database_error)?;
4260 }
4261 }
4262 {
4263 let mut heads = write.open_table(AUTHOR_HEADS).map_err(database_error)?;
4264 let keys = table_keys_with_prefix(&heads, domain_id.as_bytes())?;
4265 for key in keys {
4266 heads.remove(key.as_slice()).map_err(database_error)?;
4267 }
4268 }
4269 Ok(())
4270}
4271
4272fn table_keys_with_prefix(
4273 table: &impl redb::ReadableTable<&'static [u8], &'static [u8]>,
4274 prefix: &[u8],
4275) -> Result<Vec<Vec<u8>>, StoreError> {
4276 let mut keys = Vec::new();
4277 for entry in table.range(prefix..).map_err(database_error)? {
4278 let (key, _) = entry.map_err(database_error)?;
4279 if !key.value().starts_with(prefix) {
4280 break;
4281 }
4282 keys.push(key.value().to_vec());
4283 }
4284 Ok(keys)
4285}
4286
4287fn table_u8_keys_with_prefix(
4288 table: &impl redb::ReadableTable<&'static [u8], u8>,
4289 prefix: &[u8],
4290) -> Result<Vec<Vec<u8>>, StoreError> {
4291 let mut keys = Vec::new();
4292 for entry in table.range(prefix..).map_err(database_error)? {
4293 let (key, _) = entry.map_err(database_error)?;
4294 if !key.value().starts_with(prefix) {
4295 break;
4296 }
4297 keys.push(key.value().to_vec());
4298 }
4299 Ok(keys)
4300}
4301
4302fn table_u64_keys_with_prefix(
4303 table: &impl redb::ReadableTable<&'static [u8], u64>,
4304 prefix: &[u8],
4305) -> Result<Vec<Vec<u8>>, StoreError> {
4306 let mut keys = Vec::new();
4307 for entry in table.range(prefix..).map_err(database_error)? {
4308 let (key, _) = entry.map_err(database_error)?;
4309 if !key.value().starts_with(prefix) {
4310 break;
4311 }
4312 keys.push(key.value().to_vec());
4313 }
4314 Ok(keys)
4315}
4316
4317fn decode_index_value(
4318 key: &[u8],
4319 domain_id: DomainId,
4320 collection_id: CollectionId,
4321 record_id: &[u8],
4322) -> Result<IndexValue, StoreError> {
4323 if key.len() < 76
4324 || !key.starts_with(domain_id.as_bytes())
4325 || key[32..64] != *collection_id.as_bytes()
4326 {
4327 return Err(StoreError::CorruptDerivedMetadata(
4328 "secondary index prefix is invalid".into(),
4329 ));
4330 }
4331 let field_id = u32::from_be_bytes(
4332 key[64..68]
4333 .try_into()
4334 .map_err(|_| StoreError::CorruptDerivedMetadata("index field is invalid".into()))?,
4335 );
4336 let value_len = read_u32(&key[68..72])?;
4337 let value_end = 72_usize
4338 .checked_add(value_len)
4339 .ok_or_else(|| StoreError::CorruptDerivedMetadata("index value overflow".into()))?;
4340 if key.len() < value_end + 4 {
4341 return Err(StoreError::CorruptDerivedMetadata(
4342 "secondary index value is truncated".into(),
4343 ));
4344 }
4345 let record_len = read_u32(&key[value_end..value_end + 4])?;
4346 if key.len() != value_end + 4 + record_len || &key[value_end + 4..] != record_id {
4347 return Err(StoreError::CorruptDerivedMetadata(
4348 "secondary index record ID is invalid".into(),
4349 ));
4350 }
4351 Ok(IndexValue::new(field_id, key[72..value_end].to_vec()))
4352}
4353
4354fn schema_key(domain_id: DomainId, collection_id: CollectionId, version: u32) -> [u8; 68] {
4355 let mut key = [0_u8; 68];
4356 key[..32].copy_from_slice(domain_id.as_bytes());
4357 key[32..64].copy_from_slice(collection_id.as_bytes());
4358 key[64..].copy_from_slice(&version.to_be_bytes());
4359 key
4360}
4361
4362fn active_schema_key(domain_id: DomainId, collection_id: CollectionId) -> [u8; 64] {
4363 let mut key = [0_u8; 64];
4364 key[..32].copy_from_slice(domain_id.as_bytes());
4365 key[32..].copy_from_slice(collection_id.as_bytes());
4366 key
4367}
4368
4369fn operational_metadata_key(domain_id: DomainId, tag: u8) -> [u8; 33] {
4370 let mut key = [0_u8; 33];
4371 key[..32].copy_from_slice(domain_id.as_bytes());
4372 key[32] = tag;
4373 key
4374}
4375
4376fn migration_stage_prefix(domain_id: DomainId, migration_id: MigrationId) -> [u8; 64] {
4377 let mut prefix = [0_u8; 64];
4378 prefix[..32].copy_from_slice(domain_id.as_bytes());
4379 prefix[32..].copy_from_slice(migration_id.as_bytes());
4380 prefix
4381}
4382
4383fn migration_stage_key(
4384 domain_id: DomainId,
4385 migration_id: MigrationId,
4386 record_id: &[u8],
4387) -> Vec<u8> {
4388 let prefix = migration_stage_prefix(domain_id, migration_id);
4389 let mut key = Vec::with_capacity(prefix.len() + record_id.len());
4390 key.extend_from_slice(&prefix);
4391 key.extend_from_slice(record_id);
4392 key
4393}
4394
4395fn migration_activation_key(
4396 domain_id: DomainId,
4397 collection_id: CollectionId,
4398 from_schema: iroh_db_core::SchemaId,
4399) -> [u8; 96] {
4400 let mut key = [0_u8; 96];
4401 key[..32].copy_from_slice(domain_id.as_bytes());
4402 key[32..64].copy_from_slice(collection_id.as_bytes());
4403 key[64..].copy_from_slice(from_schema.as_bytes());
4404 key
4405}
4406
4407fn record_prefix(domain_id: DomainId, collection_id: CollectionId) -> [u8; 64] {
4408 let mut prefix = [0_u8; 64];
4409 prefix[..32].copy_from_slice(domain_id.as_bytes());
4410 prefix[32..].copy_from_slice(collection_id.as_bytes());
4411 prefix
4412}
4413
4414fn record_key(domain_id: DomainId, collection_id: CollectionId, record_id: &[u8]) -> Vec<u8> {
4415 let prefix = record_prefix(domain_id, collection_id);
4416 let mut key = Vec::with_capacity(prefix.len() + record_id.len());
4417 key.extend_from_slice(&prefix);
4418 key.extend_from_slice(record_id);
4419 key
4420}
4421
4422fn index_value_prefix(
4423 domain_id: DomainId,
4424 collection_id: CollectionId,
4425 field_id: u32,
4426 value: &[u8],
4427) -> Result<Vec<u8>, StoreError> {
4428 let value_len = u32::try_from(value.len())
4429 .map_err(|_| StoreError::CorruptDerivedMetadata("index value length exceeds u32".into()))?;
4430 let mut key = Vec::with_capacity(72 + value.len());
4431 key.extend_from_slice(domain_id.as_bytes());
4432 key.extend_from_slice(collection_id.as_bytes());
4433 key.extend_from_slice(&field_id.to_be_bytes());
4434 key.extend_from_slice(&value_len.to_be_bytes());
4435 key.extend_from_slice(value);
4436 Ok(key)
4437}
4438
4439fn index_key(
4440 domain_id: DomainId,
4441 collection_id: CollectionId,
4442 field_id: u32,
4443 value: &[u8],
4444 record_id: &[u8],
4445) -> Result<Vec<u8>, StoreError> {
4446 let mut key = index_value_prefix(domain_id, collection_id, field_id, value)?;
4447 let record_len = u32::try_from(record_id.len())
4448 .map_err(|_| StoreError::CorruptDerivedMetadata("record ID length exceeds u32".into()))?;
4449 key.extend_from_slice(&record_len.to_be_bytes());
4450 key.extend_from_slice(record_id);
4451 Ok(key)
4452}
4453
4454fn encode_key_list(keys: &[Vec<u8>]) -> Result<Vec<u8>, StoreError> {
4455 let count = u32::try_from(keys.len())
4456 .map_err(|_| StoreError::CorruptDerivedMetadata("too many index keys".into()))?;
4457 let mut encoded = Vec::new();
4458 encoded.extend_from_slice(&count.to_be_bytes());
4459 for key in keys {
4460 let length = u32::try_from(key.len()).map_err(|_| {
4461 StoreError::CorruptDerivedMetadata("index key length exceeds u32".into())
4462 })?;
4463 encoded.extend_from_slice(&length.to_be_bytes());
4464 encoded.extend_from_slice(key);
4465 }
4466 Ok(encoded)
4467}
4468
4469fn decode_key_list(bytes: &[u8]) -> Result<Vec<Vec<u8>>, StoreError> {
4470 if bytes.len() < 4 {
4471 return Err(StoreError::CorruptDerivedMetadata(
4472 "index key list count is missing".into(),
4473 ));
4474 }
4475 let count = read_u32(&bytes[..4])?;
4476 let mut offset = 4;
4477 let mut keys = Vec::with_capacity(count);
4478 for _ in 0..count {
4479 if bytes.len().saturating_sub(offset) < 4 {
4480 return Err(StoreError::CorruptDerivedMetadata(
4481 "index key length is missing".into(),
4482 ));
4483 }
4484 let length = read_u32(&bytes[offset..offset + 4])?;
4485 offset += 4;
4486 if bytes.len().saturating_sub(offset) < length {
4487 return Err(StoreError::CorruptDerivedMetadata(
4488 "index key bytes are truncated".into(),
4489 ));
4490 }
4491 keys.push(bytes[offset..offset + length].to_vec());
4492 offset += length;
4493 }
4494 if offset != bytes.len() {
4495 return Err(StoreError::CorruptDerivedMetadata(
4496 "index key list has trailing data".into(),
4497 ));
4498 }
4499 Ok(keys)
4500}
4501
4502fn read_u32(bytes: &[u8]) -> Result<usize, StoreError> {
4503 let bytes: [u8; 4] = bytes
4504 .try_into()
4505 .map_err(|_| StoreError::CorruptDerivedMetadata("expected a four-byte length".into()))?;
4506 usize::try_from(u32::from_be_bytes(bytes)).map_err(|error| {
4507 StoreError::CorruptDerivedMetadata(format!("invalid platform length: {error}"))
4508 })
4509}
4510
4511fn author_sequence_key(domain_id: DomainId, author: AuthorId) -> [u8; 64] {
4512 let mut key = [0_u8; 64];
4513 key[..32].copy_from_slice(domain_id.as_bytes());
4514 key[32..].copy_from_slice(author.as_bytes());
4515 key
4516}
4517
4518fn author_commit_key(domain_id: DomainId, author: AuthorId, sequence: u64) -> [u8; 72] {
4519 let mut key = [0_u8; 72];
4520 key[..64].copy_from_slice(&author_sequence_key(domain_id, author));
4521 key[64..].copy_from_slice(&sequence.to_be_bytes());
4522 key
4523}
4524
4525fn blob_ownership_key(domain_id: DomainId, hash: BlobHash) -> [u8; 64] {
4526 let mut key = [0_u8; 64];
4527 key[..32].copy_from_slice(domain_id.as_bytes());
4528 key[32..].copy_from_slice(hash.as_bytes());
4529 key
4530}
4531
4532fn frontier_key(domain_id: DomainId, commit_id: CommitId) -> [u8; 64] {
4533 let mut key = [0_u8; 64];
4534 key[..32].copy_from_slice(domain_id.as_bytes());
4535 key[32..].copy_from_slice(commit_id.as_bytes());
4536 key
4537}
4538
4539fn commit_context_key(domain_id: DomainId, commit_id: CommitId) -> [u8; 64] {
4540 frontier_key(domain_id, commit_id)
4541}
4542
4543fn record_head_prefix(
4544 domain_id: DomainId,
4545 collection_id: CollectionId,
4546 record_id: &[u8],
4547) -> Result<Vec<u8>, StoreError> {
4548 let record_length = u32::try_from(record_id.len()).map_err(|_| {
4549 StoreError::CorruptDerivedMetadata("record head ID length exceeds u32".into())
4550 })?;
4551 let mut key = Vec::with_capacity(68 + record_id.len());
4552 key.extend_from_slice(domain_id.as_bytes());
4553 key.extend_from_slice(collection_id.as_bytes());
4554 key.extend_from_slice(&record_length.to_be_bytes());
4555 key.extend_from_slice(record_id);
4556 Ok(key)
4557}
4558
4559fn record_head_key(domain_id: DomainId, head: &RecordHead) -> Result<Vec<u8>, StoreError> {
4560 let mut key = record_head_prefix(
4561 domain_id,
4562 head.mutation().schema().collection_id(),
4563 head.mutation().record_id(),
4564 )?;
4565 key.extend_from_slice(head.commit_id().as_bytes());
4566 Ok(key)
4567}
4568
4569fn encode_version_vector(context: &VersionVector) -> Result<Vec<u8>, StoreError> {
4570 minicbor::to_vec(context).map_err(|error| {
4571 StoreError::CorruptDerivedMetadata(format!("commit context encoding failed: {error}"))
4572 })
4573}
4574
4575fn decode_version_vector(bytes: &[u8]) -> Result<VersionVector, StoreError> {
4576 let mut decoder = minicbor::Decoder::new(bytes);
4577 let context: VersionVector = decoder.decode().map_err(|error| {
4578 StoreError::CorruptDerivedMetadata(format!("commit context is invalid: {error}"))
4579 })?;
4580 if decoder.position() != bytes.len() || encode_version_vector(&context)? != bytes {
4581 return Err(StoreError::CorruptDerivedMetadata(
4582 "commit context is not canonical".into(),
4583 ));
4584 }
4585 Ok(context)
4586}
4587
4588fn encode_record_head(head: &RecordHead) -> Result<Vec<u8>, StoreError> {
4589 let mutation = head.mutation().encode_canonical().map_err(|error| {
4590 StoreError::CorruptDerivedMetadata(format!("record head mutation is invalid: {error}"))
4591 })?;
4592 let mut encoded = Vec::with_capacity(40 + mutation.len());
4593 encoded.extend_from_slice(head.author().as_bytes());
4594 encoded.extend_from_slice(&head.sequence().to_be_bytes());
4595 encoded.extend_from_slice(&mutation);
4596 Ok(encoded)
4597}
4598
4599fn decode_record_head(commit_id: CommitId, bytes: &[u8]) -> Result<RecordHead, StoreError> {
4600 if bytes.len() < 40 {
4601 return Err(StoreError::CorruptDerivedMetadata(
4602 "record head value is truncated".into(),
4603 ));
4604 }
4605 let author: [u8; 32] = bytes[..32].try_into().expect("exact author length");
4606 let sequence: [u8; 8] = bytes[32..40].try_into().expect("exact sequence length");
4607 let mutation = MaterializedRecord::decode_canonical(&bytes[40..]).map_err(|error| {
4608 StoreError::CorruptDerivedMetadata(format!("record head mutation is invalid: {error}"))
4609 })?;
4610 Ok(RecordHead::new(
4611 commit_id,
4612 AuthorId::from_bytes(author),
4613 u64::from_be_bytes(sequence),
4614 mutation,
4615 ))
4616}
4617
4618fn encode_metadata(envelope: &CommitEnvelope) -> Vec<u8> {
4619 let mut metadata = Vec::with_capacity(112);
4620 metadata.extend_from_slice(envelope.header().domain_id().as_bytes());
4621 metadata.extend_from_slice(envelope.header().author().as_bytes());
4622 metadata.extend_from_slice(&envelope.header().epoch().to_be_bytes());
4623 metadata.extend_from_slice(&envelope.header().author_sequence().to_be_bytes());
4624 metadata.extend_from_slice(envelope.header().capability_id().as_bytes());
4625 metadata
4626}
4627
4628#[allow(clippy::needless_pass_by_value)]
4630fn io_error(error: std::io::Error) -> StoreError {
4631 StoreError::Io(error.to_string())
4632}
4633
4634#[allow(clippy::needless_pass_by_value)]
4636fn database_error(error: impl std::fmt::Display) -> StoreError {
4637 StoreError::Database(error.to_string())
4638}