iroh_db_store/
store.rs

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
22/// Default upper bound for the metadata redb userspace page cache.
23pub 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/// The result of durably recording an immutable commit.
102#[derive(Debug, Clone, Copy, PartialEq, Eq)]
103pub enum ApplyOutcome {
104    /// The commit blob and metadata became durable.
105    Applied(CommitId),
106    /// The exact commit was already durable; no metadata changed.
107    Duplicate(CommitId),
108}
109
110/// One immutable operation that remains a causal head for a record.
111#[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    /// Constructs a durable record-head entry from an authenticated commit operation.
121    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/// Complete replacement set for one record's rebuildable causal heads.
153#[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/// A local durable-store failure.
204#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
205pub enum StoreError {
206    /// A filesystem operation failed.
207    #[error("store filesystem operation failed: {0}")]
208    Io(String),
209    /// The redb metadata operation failed.
210    #[error("metadata database operation failed: {0}")]
211    Database(String),
212    /// The immutable blob repository failed.
213    #[error(transparent)]
214    Blob(#[from] BlobStoreError),
215    /// The blob repository returned a hash other than the bytes' expected hash.
216    #[error("blob store returned a mismatched content hash")]
217    BlobHashMismatch,
218    /// A commit skipped or repeated an author's required next sequence.
219    #[error("author sequence gap: expected {expected}, got {actual}")]
220    SequenceGap {
221        /// The only sequence accepted next.
222        expected: u64,
223        /// The sequence supplied by the commit.
224        actual: u64,
225    },
226    /// The next author sequence did not causally extend that author's accepted head.
227    #[error("commit does not extend previous author commit {previous}")]
228    AuthorChainMissing {
229        /// The accepted author commit that must be in dependency ancestry.
230        previous: CommitId,
231    },
232    /// A post-snapshot author commit did not extend the complete installed frontier.
233    #[error("commit does not extend installed snapshot frontier commit {frontier}")]
234    SnapshotBaselineMissing {
235        /// A current frontier commit absent from incoming dependency ancestry.
236        frontier: CommitId,
237    },
238    /// Two different immutable commits claim one author-local sequence.
239    #[error(
240        "author {author} equivocated at sequence {sequence}: conflicting commits {first} and {second}"
241    )]
242    AuthorEquivocation {
243        /// The signing author that reused a sequence.
244        author: AuthorId,
245        /// The reused author-local sequence.
246        sequence: u64,
247        /// Lexicographically first conflicting commit ID.
248        first: CommitId,
249        /// Lexicographically second conflicting commit ID.
250        second: CommitId,
251    },
252    /// A commit is immutable evidence in a quarantined equivocation branch.
253    #[error("commit is quarantined: {0}")]
254    QuarantinedCommit(CommitId),
255    /// An author's signing key equivocated and cannot extend that branch.
256    #[error("author {author} is quarantined from sequence {from_sequence}")]
257    QuarantinedAuthor {
258        /// The equivocated author identity.
259        author: AuthorId,
260        /// The first sequence no longer accepted from this author.
261        from_sequence: u64,
262    },
263    /// A declared causal dependency is not accepted locally.
264    #[error("commit dependency is missing: {0}")]
265    MissingDependency(CommitId),
266    /// The directory has an incompatible format marker.
267    #[error("unsupported store directory format")]
268    UnsupportedFormat,
269    /// Another process already holds the exclusive writer lock.
270    #[error("database is already open for writing")]
271    AlreadyOpen,
272    /// A stored commit blob was missing.
273    #[error("commit blob is missing: {0}")]
274    MissingCommitBlob(CommitId),
275    /// Stored immutable bytes were not a canonical commit.
276    #[error(transparent)]
277    CommitCodec(#[from] CommitCodecError),
278    /// A stored schema descriptor could not be decoded.
279    #[error(transparent)]
280    Schema(#[from] SchemaError),
281    /// A canonical schema migration object was invalid.
282    #[error(transparent)]
283    Migration(#[from] MigrationCodecError),
284    /// Migration-derived state did not match its immutable commit history.
285    #[error("schema migration metadata is invalid: {0}")]
286    InvalidMigration(String),
287    /// Two activation commits attempted different successors for one schema.
288    #[error("collection {collection_id} is frozen by incompatible schema activations")]
289    SchemaMigrationFork { collection_id: CollectionId },
290    /// The same domain, collection and version was registered with different bytes.
291    #[error("schema mismatch for collection {collection_id} version {version}")]
292    SchemaMismatch {
293        /// Stable collection identifier.
294        collection_id: CollectionId,
295        /// Application schema version.
296        version: u32,
297    },
298    /// A supplied index was absent from the schema or was not declared indexed.
299    #[error("field {field_id} is not an index in collection {collection_id}")]
300    InvalidIndex {
301        /// Stable collection identifier.
302        collection_id: CollectionId,
303        /// Stable field identifier.
304        field_id: u32,
305    },
306    /// Stored derived metadata was structurally corrupt.
307    #[error("derived metadata is corrupt: {0}")]
308    CorruptDerivedMetadata(String),
309    /// A snapshot supplied a noncanonical or inconsistent causal baseline.
310    #[error("snapshot metadata is invalid: {0}")]
311    InvalidSnapshotMetadata(String),
312    /// A signed domain authority object failed validation or conflicted with durable state.
313    #[error("domain authority metadata is invalid: {0}")]
314    InvalidAuthority(String),
315    /// The designated controller signed two different successors for one control slot.
316    #[error(
317        "domain {domain} is frozen by conflicting control transitions {first:?} and {second:?}"
318    )]
319    ControlFork {
320        /// Security domain whose controller equivocated.
321        domain: DomainId,
322        /// Lexicographically first transition ID.
323        first: [u8; 32],
324        /// Lexicographically second transition ID.
325        second: [u8; 32],
326    },
327}
328
329/// Canonical durable reference to the immutable signed snapshot used as a history baseline.
330#[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/// Transactional metadata paired with an immutable content-addressed blob repository.
433#[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    /// Opens or creates a local store directory and acquires its writer lock.
447    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    /// Opens a local store with an explicit redb userspace page-cache bound.
452    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    /// Makes immutable bytes durable before atomically publishing their metadata.
491    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    /// Makes a commit blob durable, then atomically publishes metadata and derived records.
499    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    /// Persists a record commit and atomically advances its incremental causal index.
509    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    /// Persists an invisible bounded set of migration stages with its immutable commit.
554    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    /// Atomically publishes a completely staged schema transition with its commit.
579    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    /// Persists schema stages while atomically advancing the causal commit index.
606    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    /// Persists schema activation while atomically advancing the causal commit index.
622    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    /// Returns the active schema activation for a collection, if it has migrated.
763    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    /// Lists active schema cuts for a domain in collection order.
783    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    /// Returns whether incompatible activation siblings froze a collection.
807    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    /// Returns low-cardinality active, staged-record, and frozen migration counts.
821    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    /// Records when this device last created or installed a verified snapshot.
837    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    /// Returns the durable local observation time of the newest verified snapshot.
858    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    /// Loads every durable stage for one migration in source-key order.
870    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    /// Returns whether commit metadata is accepted locally.
956    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    /// Returns whether a commit ID is retained only as equivocation evidence.
983    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    /// Returns the first rejected sequence for an equivocated author, if any.
995    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    /// Returns the latest accepted author sequence in a domain.
1012    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    /// Returns the latest accepted immutable commit for an author in a domain.
1027    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    /// Returns an accepted conflicting commit for the envelope's author sequence.
1048    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    /// Atomically quarantines an equivocated branch and publishes rebuilt accepted state.
1063    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                    // Snapshot-covered history is trusted as an opaque causal baseline.
1291                }
1292                Err(error) => return Err(error),
1293            }
1294        }
1295        Ok(false)
1296    }
1297
1298    /// Returns the canonical sorted commit frontier for a domain.
1299    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    /// Returns whether the rebuildable causal and record-head indexes are complete.
1318    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    /// Loads the inclusive causal context cached for one accepted commit.
1330    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    /// Loads only the current causal heads for one materialized record.
1346    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    /// Atomically replaces the complete rebuildable reconciliation index for a domain.
1381    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    /// Loads one canonical materialized record state by primary key.
1418    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    /// Lists canonical record IDs and states in stable primary-key order.
1434    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    /// Resolves an equality index to stable record IDs.
1457    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    /// Loads a registered schema version, if present.
1493    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    /// Loads the highest registered schema version for one collection.
1512    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    /// Lists accepted commit IDs, including snapshot-covered history, in deterministic order.
1540    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    /// Returns whether this domain relies on snapshot-covered commit history.
1559    pub fn has_snapshot_coverage(&self, domain_id: DomainId) -> Result<bool, StoreError> {
1560        Ok(!self.snapshot_coverage(domain_id)?.is_empty())
1561    }
1562
1563    /// Lists the exact opaque commit IDs represented by the installed snapshot.
1564    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    /// Loads the immutable snapshot recovery baseline for a compacted domain.
1587    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    /// Idempotently records local ownership of an immutable blob by one security domain.
1612    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    /// Returns whether an accepted commit or claimed immutable blob belongs to a domain.
1623    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    /// Lists commit IDs whose immutable envelope blobs are stored locally.
1648    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    /// Lists locally stored accepted commits visible through an epoch.
1674    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    /// Exports durable per-author sequence baselines for a snapshot.
1711    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    /// Exports complete canonical derived state for snapshots and equivalence checks.
1733    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    /// Atomically replaces all derived state for one domain.
1794    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    /// Atomically installs snapshot-derived state and its causal acceptance baseline.
1806    #[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    /// Loads and strictly decodes an accepted immutable commit blob.
1886    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    /// Finds canonical commit blobs left behind before their metadata transaction committed.
1913    ///
1914    /// Arbitrary application blobs are ignored. The returned IDs are deterministic and can be
1915    /// offered to a repair workflow; scanning never mutates durable state.
1916    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    /// Atomically records a domain genesis descriptor and its owner capability.
1938    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    /// Atomically imports a descriptor, chain, active root/leaf, and trust set.
1989    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    /// Atomically imports a signed current authority baseline from a targeted invitation.
2083    #[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    /// Installs a delegated capability after validating its durable issuer link.
2181    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    /// Atomically installs a validated root-to-leaf capability chain.
2206    ///
2207    /// Only the leaf becomes active. Intermediate certificates are retained as
2208    /// immutable delegation evidence and cannot replace another subject's
2209    /// active capability as a side effect of authenticating the leaf.
2210    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    /// Loads one verified capability certificate by content ID.
2231    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    /// Resolves a bounded root-to-leaf capability chain.
2247    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    /// Loads the immutable descriptor when this is a capability-controlled domain.
2275    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    /// Returns the active capability selected for an endpoint in a domain.
2291    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    /// Lists exact active capability IDs for a domain in canonical order.
2314    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    /// Checks active, non-revoked, endpoint-bound permission at an epoch.
2339    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    /// Loads and verifies the current single-successor control head.
2368    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    /// Returns the accepted control chain in ascending sequence order.
2385    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    /// Loads and verifies the optional history-compacting authority baseline.
2443    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    /// Checks persisted checkpoint evidence against immutable authority objects.
2470    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    /// Loads one retained canonical control object by content ID.
2495    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    /// Loads one retained accepted transition by its signed sequence.
2504    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    /// Returns the epoch authorized by the durable control chain.
2516    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    /// Returns the accepted domain-control sequence, or zero at genesis.
2524    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    /// Checks the rebuildable subject-revocation index against signed control.
2532    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    /// Returns the canonical revocation accumulator derived from signed authority.
2554    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    /// Atomically rebuilds subject revocations from the accepted signed chain.
2565    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    /// Returns the endpoint exclusively authorized to sign the next control transition.
2589    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    /// Returns the currently designated single writer.
2600    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    /// Verifies and retains a control object before separately delivered key material is activated.
2611    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    /// Atomically advances authority, records the causal cut, and revokes endpoint access.
2671    pub fn apply_control(&self, control: &ControlTransition) -> Result<(), StoreError> {
2672        self.apply_control_bundle(control, &[])
2673    }
2674
2675    /// Atomically installs successor capability evidence and advances authority.
2676    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    /// Adds a device to the local domain authorization set.
3169    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    /// Returns whether a device is locally authorized for this domain.
3180    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    /// Lists locally authorized devices in canonical identity order.
3195    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    /// Replaces bounded application-facing group metadata without changing protocol identity.
3214    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    /// Lists device groups for one domain in label order.
3243    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    /// Atomically marks an offline invitation ticket consumed, returning false on replay.
3278    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// Error conversion consumes std's typed error while only retaining safe display text.
4629#[allow(clippy::needless_pass_by_value)]
4630fn io_error(error: std::io::Error) -> StoreError {
4631    StoreError::Io(error.to_string())
4632}
4633
4634// Error conversion consumes redb's typed errors while only retaining safe display text.
4635#[allow(clippy::needless_pass_by_value)]
4636fn database_error(error: impl std::fmt::Display) -> StoreError {
4637    StoreError::Database(error.to_string())
4638}