Skip to main content

coven_replication/sync/store/commit_publication/operation/
facades.rs

1use super::*;
2
3impl<'storage> AuthorizedWriterOperation<'storage> {
4    pub(crate) fn accepted_commit_membership_state(
5        &self,
6        reference: &coven_protocol::store_commit::StoreBatchCommitRef,
7    ) -> Option<&coven_protocol::circle_control::StoreMembershipStateRef> {
8        self.history.accepted_commit_membership_state(reference)
9    }
10
11    pub(super) fn membership_objects(&self) -> StoreMembershipObjectVerifier<'_, 'storage> {
12        self.history.membership_objects()
13    }
14
15    pub(crate) fn store_root(&self) -> &coven_protocol::store_commit::StoreRootRef {
16        self.history.root()
17    }
18
19    pub(crate) async fn snapshot_publication(
20        &self,
21    ) -> crate::sync::store::snapshots::AuthorizedSnapshotPublication<'_> {
22        crate::sync::store::snapshots::AuthorizedSnapshotPublication::begin(
23            &self.database,
24            self.storage.as_ref(),
25            self.store_dir,
26        )
27        .await
28    }
29
30    pub(crate) async fn resume_snapshot_publication(
31        &self,
32    ) -> Result<
33        Option<coven_protocol::store_commit::SnapshotMeta>,
34        crate::sync::store::snapshots::SnapshotError,
35    > {
36        self.snapshot_publication().await.resume_store().await
37    }
38
39    pub(super) fn protocol_root(&self) -> &coven_protocol::store_commit::StoreProtocolRoot {
40        &self.history.verified_root_object().value
41    }
42
43    pub(super) fn resolved_membership(
44        &self,
45    ) -> Result<
46        &coven_protocol::membership::MembershipChain,
47        crate::sync::store::membership::MembershipOpsError,
48    > {
49        match self.membership.conflict() {
50            Some(conflict) => Err(
51                crate::sync::store::membership::MembershipOpsError::SemanticConflict(Box::new(
52                    conflict.clone(),
53                )),
54            ),
55            None => Ok(&self.membership),
56        }
57    }
58
59    pub(super) async fn open_keyring(
60        &self,
61    ) -> Result<
62        coven_keys::encryption::EncryptionService,
63        crate::sync::store::commit_publication::membership::MembershipMutationError,
64    > {
65        self.keyrings.open(&self.membership).await
66    }
67
68    pub(super) async fn open_keyring_for_membership(
69        &self,
70        membership: &coven_protocol::membership::MembershipChain,
71    ) -> Result<
72        coven_keys::encryption::EncryptionService,
73        crate::sync::store::commit_publication::membership::MembershipMutationError,
74    > {
75        self.keyrings.open(membership).await
76    }
77
78    pub(crate) async fn open_keyring_or_for_membership(
79        &self,
80        membership: &coven_protocol::membership::MembershipChain,
81        initial: &coven_keys::encryption::EncryptionService,
82    ) -> Result<
83        coven_keys::encryption::EncryptionService,
84        crate::sync::store::commit_publication::membership::MembershipMutationError,
85    > {
86        self.keyrings.open_or(membership, initial).await
87    }
88
89    pub(crate) async fn prepare_wrapped_key(
90        &self,
91        recipient: &str,
92        value: coven_protocol::wrapped_store_key::WrappedStoreKey,
93    ) -> Result<
94        coven_protocol::wrapped_store_key::PreparedWrappedStoreKey,
95        coven_protocol::objects::StorageError,
96    > {
97        self.keyrings.prepare(recipient, value).await
98    }
99
100    /// Select the exact author stream without overwriting its committed prefix.
101    /// Streams are persisted per database, so independently restored devices use
102    /// different streams; copied state that reuses one exposes an immutable fork.
103    pub(super) async fn select_membership_author_stream(
104        &self,
105        chain: &coven_protocol::membership::MembershipChain,
106    ) -> Result<
107        coven_protocol::membership::AuthorStreamId,
108        crate::sync::store::commit_publication::membership::MembershipMutationError,
109    > {
110        let author = self.writer.author_pubkey();
111        let grant = chain.active_owner_grant(&author).ok_or_else(|| {
112            coven_protocol::membership::MembershipError::SignerIsNotOwner(author.clone())
113        })?;
114        let mut reusable = chain.reusable_author_streams(&author, &grant);
115        if let Some(anchored) = chain.membership_stream_id(&grant) {
116            reusable.insert(anchored);
117        }
118        Ok(self
119            .database
120            .select_membership_author_stream(&author, &grant, reusable)
121            .await?)
122    }
123
124    pub(super) async fn verify_membership_publication_author(
125        &self,
126        publication: &PreparedMembershipPublication,
127    ) -> Result<
128        coven_protocol::store_commit::StoreDeviceRegistration,
129        crate::sync::store::membership::MembershipMutationError,
130    > {
131        let author = self
132            .history
133            .load_registration(&publication.head.body.author_registration)
134            .await
135            .map_err(crate::sync::store::membership::MembershipMutationError::from)?
136            .value;
137        if !publication.head.verify(&author) {
138            return Err(
139                crate::sync::store::membership::MembershipMutationError::InvalidDurableMutation(
140                    "prepared membership head has an invalid certified-device signature"
141                        .to_string(),
142                ),
143            );
144        }
145        Ok(author)
146    }
147
148    pub(super) async fn blocked_candidate_nonactivation(
149        &mut self,
150        candidate: &coven_database::BlockedMergeCandidate,
151    ) -> Result<Option<coven_protocol::remote_object::VerifiedCandidateNonactivation>, StoreError>
152    {
153        let verified = self
154            .history
155            .authenticate_blocked_candidate(candidate)
156            .await?;
157        self.history
158            .merge_conflict()
159            .excluded_candidate_nonactivation(&verified, &candidate.head, &candidate.head_object)
160            .await
161    }
162
163    pub(super) async fn cleanup_merge_candidate_history(
164        &mut self,
165        write_id: coven_protocol::write::WriteId,
166    ) -> Result<(), crate::sync::store::pull::StorePullError> {
167        self.history.cleanup_merge_candidate(write_id).await
168    }
169
170    pub(crate) async fn resolve_acknowledged_snapshot(
171        &mut self,
172        registration: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
173        members: &coven_protocol::membership::MembershipChain,
174    ) -> Result<
175        Result<
176            crate::sync::store::commit_verification::merge_history::SelectedReplayBaselineRetirement,
177            crate::sync::store::ReplayBaselineDecline,
178        >,
179        crate::sync::store::acknowledgements::StoreAckError,
180    >{
181        self.history
182            .resolve_acknowledged_snapshot(registration, members)
183            .await
184    }
185
186    pub(crate) async fn select_acknowledgement_snapshot(
187        &mut self,
188        frontier: &coven_protocol::store_commit::CommitFrontier,
189        device_state: &coven_protocol::store_commit::StoreDeviceStateRef,
190    ) -> Result<
191        Option<
192            crate::sync::store::commit_verification::merge_history::SelectedInstallableStoreSnapshot,
193        >,
194        crate::sync::store::acknowledgements::StoreAckError,
195    >{
196        self.history
197            .select_acknowledgement_snapshot(frontier, device_state)
198            .await
199    }
200
201    pub(super) async fn stage_verified_blob_plaintext(
202        &self,
203        authority: &coven_protocol::blob::RowBlobAuthority,
204        stored: &coven_protocol::blob::locator::StoredBlobRef,
205        destination: &std::path::Path,
206    ) -> Result<coven_foundation::local_file::AtomicStagedFile, crate::sync::BlobCacheError> {
207        let stage = self
208            .store_dir
209            .stage_atomic_file(destination)
210            .await
211            .map_err(crate::sync::BlobCacheError::File)?;
212        self.history
213            .stage_verified_blob_plaintext(
214                authority,
215                stored,
216                stage,
217                coven_storage::cloud::no_download_progress(),
218            )
219            .await
220    }
221
222    pub(super) async fn authorize_retained_outbound(
223        &self,
224        order: &coven_protocol::store_commit::StoreCommitOrder,
225        membership_heads: &[coven_protocol::membership::MembershipHeadRef],
226    ) -> Result<
227        crate::sync::store::commit_verification::merge_history::MergeOutboundAuthorization,
228        crate::sync::store::pull::StorePullError,
229    > {
230        self.writer
231            .authorize_retained_outbound(&self.history, order, membership_heads)
232            .await
233    }
234
235    /// Seed the verifier from retained history before a walk over it, so the
236    /// walk reads nothing from the provider.
237    /// Retire the owner's journal for every join whose device has arrived.
238    ///
239    /// The owner's half of a join ends at a published activation commit, and
240    /// until now it ended there permanently: the row sat at
241    /// `ActivationPrepared` for the life of the store, still offering to hand
242    /// the activation over, because the owner has no artifact by which it could
243    /// learn the joining device took it — the same asymmetry that makes the
244    /// joiner, not the owner, delete the attempt's transport slots.
245    ///
246    /// The arrival it can see is the joined device's own first commit. Every
247    /// other trace of the join is something the owner wrote: the registration
248    /// goes Active from the owner's own activation commit, so it says nothing
249    /// about whether the device ever ran. A stream in the materialized frontier
250    /// under that device's announcement stream id is a commit the device signed
251    /// and this device verified, which it can only have published after
252    /// installing the Store.
253    ///
254    /// Reads nothing from the provider: the stream id is derived from the
255    /// registration the journal already holds, and the frontier is the row the
256    /// cycle reads anyway. Each retirement is one row delete, so a cycle that
257    /// fails partway leaves the rest for the next one to find.
258    pub(crate) async fn retire_arrived_device_joins(
259        &self,
260    ) -> Result<usize, crate::sync::store::DeviceJoinError> {
261        let awaiting = self.database.owner_device_joins_awaiting_arrival().await?;
262        if awaiting.is_empty() {
263            return Ok(0);
264        }
265        let frontier = self.database.materialized_frontier().await?;
266        let store_root_hash = self.store_root().store_root_hash;
267        let mut retired = 0;
268        for (attempt_id, registration) in awaiting {
269            let stream =
270                coven_protocol::store_commit::StreamActivation::device_authorized_stream_id(
271                    store_root_hash,
272                    &registration,
273                    coven_protocol::store_commit::StreamAnchorDomain::StoreAnnouncements,
274                );
275            if !frontier.contains_key(&stream.to_string()) {
276                continue;
277            }
278            self.database
279                .retire_device_join(attempt_id, crate::sync::store::DeviceJoinRole::Owner)
280                .await?;
281            retired += 1;
282        }
283        Ok(retired)
284    }
285
286    pub(crate) async fn seed_retained_history(
287        &mut self,
288    ) -> Result<(), crate::sync::store::pull::StorePullError> {
289        self.history.seed_retained_history().await
290    }
291
292    pub(super) async fn prepare_merge_history_successor(
293        &self,
294        commit: &coven_protocol::store_commit::VerifiedStoreBatchCommit,
295        membership: &coven_protocol::membership::MembershipChain,
296        recovery_author: Option<&coven_protocol::store_commit::StoreDeviceRegistrationRef>,
297        state_after: coven_protocol::store_commit::ResolvedStoreDeviceState,
298        evidence: crate::sync::store::commit_verification::merge_history::MergeHistorySuccessorEvidence,
299    ) -> Result<
300        crate::sync::store::commit_verification::merge_history::PreparedMergeHistorySuccessor,
301        crate::sync::store::pull::StorePullError,
302    > {
303        self.history
304            .prepare_merge_history_successor(
305                commit,
306                membership,
307                recovery_author,
308                state_after,
309                evidence,
310            )
311            .await
312    }
313
314    pub(super) async fn observe_occupied_merge_head(
315        &mut self,
316        expected: &coven_protocol::store_commit::StoreDeviceHead,
317        expected_commit: &coven_protocol::store_commit::VerifiedStoreBatchCommit,
318        slot: &coven_protocol::objects::ObjectSlot,
319        semantic_prefix: &str,
320    ) -> Result<crate::sync::store::merge_conflict::VerifiedMergeWinner, StoreError> {
321        self.history
322            .merge_conflict()
323            .observe_occupied_merge_head(expected, expected_commit, slot, semantic_prefix)
324            .await
325    }
326
327    pub(super) async fn upload_commit(
328        &self,
329        candidate: &commit_plan::PreparedStoreOperationCommit,
330    ) -> Result<(), StoreError> {
331        let stream_id = candidate.reference.coord.stream_id;
332        let context = coven_protocol::objects::ProtocolObjectContext::signed_plaintext(
333            candidate.commit.store_root_hash,
334            coven_protocol::objects::ProtocolObjectDomain::StoreCommit,
335        );
336        let prefix = coven_protocol::store_commit::commit_semantic_prefix(
337            candidate.commit.candidate_family(),
338            &stream_id.to_string(),
339            candidate.commit.seq(),
340            candidate.commit.commit_hash(),
341        );
342        self.storage
343            .as_ref()
344            .create_verified_protocol_object(
345                &context,
346                &candidate.prepared_commit()?,
347                &prefix,
348                &candidate.commit.to_bytes(),
349            )
350            .await
351            .map_err(StoreError::prepared_object)
352    }
353
354    pub async fn pull(
355        &mut self,
356        routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
357    ) -> Result<crate::sync::store::StorePullResult, SyncCycleFailure> {
358        let membership = self.membership.clone();
359        let execution = self
360            .writer
361            .pull(&mut self.history, &membership, routing_encryption)
362            .await
363            .map_err(|error| SyncCycleFailure::operation("pull Store commits", error))?;
364        self.membership = execution.membership;
365        Ok(execution.result)
366    }
367
368    pub(crate) fn require_current_owner(
369        &self,
370        author_pubkey: &str,
371    ) -> Result<(), coven_protocol::membership::MembershipError> {
372        if self.membership.is_owner_now(author_pubkey) {
373            Ok(())
374        } else {
375            Err(
376                coven_protocol::membership::MembershipError::SignerIsNotOwner(
377                    author_pubkey.to_string(),
378                ),
379            )
380        }
381    }
382
383    pub(crate) async fn prepare_merge_snapshot_history_summary(
384        &self,
385        coverage: &coven_protocol::store_commit::CommitFrontier,
386        membership: &coven_protocol::membership::MembershipChain,
387        state: &coven_protocol::store_commit::ResolvedStoreDeviceState,
388    ) -> Result<
389        coven_protocol::store_commit::RetainedVerifiedMergeHistorySummary,
390        crate::sync::store::pull::StorePullError,
391    > {
392        self.writer
393            .prepare_merge_snapshot_history_summary(&self.history, coverage, membership, state)
394            .await
395    }
396
397    /// The membership objects a reader of this Store's current frontier would
398    /// otherwise fetch one at a time — what a snapshot publishes as its
399    /// membership rollup.
400    ///
401    /// This runs the anchored walk again rather than keeping what an earlier
402    /// one read, because the rollup has to describe the whole chain and not
403    /// whatever part of it this operation happened to touch. The walk is
404    /// served from this verifier's own slot and object memos, so on a device
405    /// that has already resolved its membership this costs the terminating
406    /// probe per stream and nothing else.
407    pub(crate) async fn membership_rollup_parts(
408        &mut self,
409        membership: &coven_protocol::membership::MembershipChain,
410    ) -> Result<
411        (
412            Vec<coven_protocol::store_commit::MembershipRollupStream>,
413            Vec<coven_protocol::store_commit::MembershipRollupResolution>,
414        ),
415        crate::sync::store::membership::AnchoredChainError,
416    > {
417        self.history.membership_rollup_parts(membership).await
418    }
419
420    pub(crate) fn snapshots(
421        &mut self,
422    ) -> crate::sync::store::snapshots::AuthorizedSnapshots<'_, 'storage> {
423        let database = self.database.clone();
424        let storage = Arc::clone(self.storage);
425        let store_dir = self.store_dir;
426        let membership = self.membership.clone();
427        let local_writer = Arc::clone(&self.writer);
428        crate::sync::store::snapshots::AuthorizedSnapshots::new(
429            self,
430            database,
431            storage,
432            store_dir,
433            membership,
434            local_writer,
435        )
436    }
437
438    pub(crate) fn acknowledgements(
439        &mut self,
440    ) -> crate::sync::store::acknowledgements::AuthorizedAcknowledgements<'_, 'storage> {
441        let database = self.database.clone();
442        let storage = Arc::clone(self.storage);
443        let local_writer = Arc::clone(&self.writer);
444        crate::sync::store::acknowledgements::AuthorizedAcknowledgements::new(
445            self,
446            database,
447            storage,
448            local_writer,
449        )
450    }
451
452    pub(crate) fn reclaim_history(
453        &mut self,
454    ) -> crate::sync::store::reclaim::ReclaimHistory<'_, 'storage> {
455        self.history.reclaim()
456    }
457
458    pub(crate) fn owner_promotion(
459        &mut self,
460    ) -> crate::sync::store::owner_role_promotion::AuthorizedOwnerPromotion<'_, 'storage> {
461        let database = self.database.clone();
462        let storage = self.storage.clone();
463        let root = self.store_root().clone();
464        let membership = self.membership.clone();
465        crate::sync::store::owner_role_promotion::AuthorizedOwnerPromotion::new(
466            self, database, storage, root, membership,
467        )
468    }
469
470    pub(crate) fn owner_promotion_history(
471        &mut self,
472    ) -> crate::sync::store::owner_role_promotion::OwnerPromotionHistory<'_, 'storage> {
473        self.history.owner_promotion()
474    }
475
476    pub(crate) async fn refresh_authorization_state(
477        &self,
478        cipher: &dyn coven_storage::CloudSyncCipherStateAccess,
479        pending_rotation: &dyn coven_storage::CloudSyncRotationStateAccess,
480        master_keys: Option<&dyn coven_keys::keys::MasterKeyCustody>,
481    ) -> Result<(), SyncCycleFailure> {
482        let result = async {
483            if cipher.is_plaintext() {
484                tracing::debug!("refresh: plaintext home, nothing to refresh");
485                return Ok(());
486            }
487
488            let recipient = self.writer.author_pubkey();
489            let wrapped_keys = self
490                .membership
491                .wrapped_key_authority_for(&recipient)
492                .map_err(AuthorizationRefreshError::Membership)?;
493            if wrapped_keys.is_empty() {
494                tracing::debug!(
495                    "refresh: no activated wrapped key for this device; keeping the live key"
496                );
497                return Ok(());
498            }
499
500            match self.open_keyring().await {
501                Ok(new_encryption) => {
502                    let merged = cipher
503                        .merged_keyring(&new_encryption)
504                        .map_err(AuthorizationRefreshError::InvalidKeyring)?;
505                    if merged.merged_key_count() == merged.live_key_count() {
506                        if pending_rotation.gate().is_some() {
507                            let gate = self
508                                .database
509                                .complete_peer_rotation_adoption(merged.merged_generation())
510                                .await
511                                .map_err(AuthorizationRefreshError::Database)?;
512                            pending_rotation.install_durable_gate(gate);
513                        }
514                        tracing::debug!(
515                            "refresh: wrapped store key is already held by the live keyring"
516                        );
517                    } else {
518                        let gate = self
519                            .database
520                            .record_peer_rotation(merged.merged_generation())
521                            .await
522                            .map_err(AuthorizationRefreshError::Database)?;
523                        pending_rotation.install_durable_gate(Some(gate));
524                        match master_keys {
525                            None => {
526                                tracing::info!(
527                                    committed_generation = merged.merged_generation(),
528                                    "refresh: found a rotated store key but this cycle has no \
529                                     master-key custody to adopt it; sealing is paused until a \
530                                     cycle with custody adopts it"
531                                );
532                            }
533                            Some(master_keys) => {
534                                let adopted = cipher
535                                    .adopt_key_rotation(&new_encryption, master_keys)
536                                    .map_err(AuthorizationRefreshError::KeyAdoption)?;
537                                let gate = self
538                                    .database
539                                    .complete_peer_rotation_adoption(adopted.generation())
540                                    .await
541                                    .map_err(AuthorizationRefreshError::Database)?;
542                                pending_rotation.install_durable_gate(gate);
543                                tracing::info!(
544                                    fingerprint = adopted.fingerprint(),
545                                    "Adopted rotated store key"
546                                );
547                            }
548                        }
549                    }
550                }
551                Err(error) => return Err(AuthorizationRefreshError::WrappedKey(error)),
552            }
553
554            Ok(())
555        }
556        .await;
557
558        result.map_err(|error| SyncCycleFailure::operation("refresh authorization state", error))
559    }
560
561    pub(crate) fn circles(
562        &mut self,
563    ) -> crate::sync::store::circles::AuthorizedCircleWriter<'_, 'storage> {
564        let database = self.database.clone();
565        let storage = Arc::clone(self.storage);
566        let store_dir = self.store_dir;
567        let root = self.store_root().clone();
568        let membership = self.membership.clone();
569        let local_writer = Arc::clone(&self.writer);
570        crate::sync::store::circles::AuthorizedCircleWriter::from_parts(
571            self,
572            database,
573            storage,
574            store_dir,
575            root,
576            membership,
577            local_writer,
578        )
579    }
580
581    pub(crate) fn circle_history(
582        &mut self,
583    ) -> crate::sync::store::commit_publication::circles::VerifiedCircleHistory<'_, 'storage> {
584        self.history.circles()
585    }
586
587    pub(crate) fn join_history(
588        &mut self,
589    ) -> crate::sync::store::device_join::history::DeviceJoinHistory<'_, 'storage> {
590        self.history.device_join()
591    }
592
593    pub(crate) fn device_exclusion_history(
594        &mut self,
595    ) -> crate::sync::store::device_exclusion::DeviceExclusionHistory<'_, 'storage> {
596        self.history.device_exclusion()
597    }
598
599    pub(crate) fn device_exclusion(
600        &mut self,
601    ) -> crate::sync::store::device_exclusion::AuthorizedDeviceExclusion<'_, 'storage> {
602        let database = self.database.clone();
603        let storage = Arc::clone(self.storage);
604        crate::sync::store::device_exclusion::AuthorizedDeviceExclusion::new(
605            self, database, storage,
606        )
607    }
608
609    pub(crate) fn join_operation(
610        &mut self,
611    ) -> crate::sync::store::commit_publication::device_join::AuthorizedJoin<'_, 'storage> {
612        let database = self.database.clone();
613        let storage = Arc::clone(self.storage);
614        let root = self.store_root().clone();
615        let verified_root = self.history.verified_root_object().clone();
616        let membership = self.membership.clone();
617        let local_writer = Arc::clone(&self.writer);
618        crate::sync::store::commit_publication::device_join::AuthorizedJoin::from_parts(
619            self,
620            database,
621            storage,
622            root,
623            verified_root,
624            membership,
625            local_writer,
626        )
627    }
628
629    pub(super) async fn membership_mutation_permit(
630        &self,
631    ) -> coven_database::store::MembershipMutationPermit {
632        self.database.membership_mutation_permit().await
633    }
634
635    pub(super) fn writer_pubkey(&self) -> String {
636        self.writer.author_pubkey()
637    }
638
639    pub(crate) fn local_author_pubkey(&self) -> String {
640        self.writer.author_pubkey()
641    }
642
643    pub(crate) fn is_local_registration(
644        &self,
645        registration: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
646    ) -> bool {
647        self.writer.is_authored_by_registration(registration)
648    }
649
650    pub(crate) fn is_current_owner(
651        &self,
652        membership: &coven_protocol::membership::MembershipChain,
653    ) -> bool {
654        self.writer.is_current_owner(membership)
655    }
656
657    pub(crate) fn matches_local_author(
658        &self,
659        registration: &coven_protocol::store_commit::StoreDeviceRegistrationRef,
660        author_pubkey: &str,
661    ) -> bool {
662        self.writer.matches_author(registration, author_pubkey)
663    }
664
665    pub(crate) fn grant_authorized_stream_id(
666        &self,
667        grant: &coven_protocol::membership::MembershipGrantId,
668        domain: coven_protocol::store_commit::StreamAnchorDomain,
669    ) -> coven_protocol::membership::AuthorStreamId {
670        self.writer
671            .grant_authorized_stream_id(self.store_root().store_root_hash, grant, domain)
672    }
673}