Skip to main content

coven_replication/sync/store/commit_verification/merge_history/
membership_control.rs

1use super::membership;
2use super::*;
3
4pub(super) struct CurrentMergeAuthority {
5    pub(super) cut: StoreHistoryCut,
6    pub(super) state: ResolvedStoreDeviceState,
7    pub(super) registrations: BTreeMap<StoreDeviceId, ReferencedStoreDeviceRegistration>,
8}
9
10pub(crate) fn verify_merge_membership_state_ref(
11    state: &StoreMembershipStateRef,
12    membership: &MembershipChain,
13    device_state: &ResolvedStoreDeviceState,
14) -> Result<(), StorePullError> {
15    let expected = merge_membership_state_ref(membership, device_state)?;
16    if &expected != state {
17        return Err(StorePullError::InvalidState(
18            "Store history membership reference differs from its exact resolved state".to_string(),
19        ));
20    }
21    Ok(())
22}
23
24pub(crate) fn merge_membership_state_ref(
25    membership: &MembershipChain,
26    device_state: &ResolvedStoreDeviceState,
27) -> Result<StoreMembershipStateRef, StorePullError> {
28    StoreMembershipStateRef::from_membership(membership, device_state.recovery.clone())
29        .map_err(StorePullError::Protocol)
30}
31
32#[derive(Clone, PartialEq, Eq)]
33pub(crate) struct VerifiedMergeMembershipHeadActivation {
34    pub(super) commit: StoreBatchCommitRef,
35    pub(super) transition: protocol_membership::MergeMembershipHeadTransition,
36}
37
38impl VerifiedMergeMembershipHeadActivation {
39    pub(crate) fn verifies(
40        &self,
41        reference: &protocol_membership::MembershipHeadRef,
42        head: &protocol_membership::AuthorHead,
43        commit: &StoreBatchCommitRef,
44    ) -> bool {
45        &self.commit == commit && self.transition.matches_head(head, reference)
46    }
47}
48
49pub(crate) struct VerifiedMergeMembershipControl {
50    pub(crate) activations: VerifiedCircleActivations,
51    pub(crate) head_activation: VerifiedMergeMembershipHeadActivation,
52    pub(crate) conflict_resolution: Option<VerifiedMergeConflictResolutionActivation>,
53}
54
55impl VerifiedMergeMembershipControl {
56    pub(crate) fn verifies_head_activation(
57        &self,
58        reference: &protocol_membership::MembershipHeadRef,
59        head: &protocol_membership::AuthorHead,
60        commit: &StoreBatchCommitRef,
61    ) -> bool {
62        self.head_activation.verifies(reference, head, commit)
63    }
64}
65
66#[derive(Clone, Default)]
67pub struct VerifiedMergeMembershipPrefix {
68    commits: BTreeSet<StoreBatchCommitRef>,
69    predecessor_memberships: Vec<MembershipChain>,
70    head_activations: BTreeMap<StoreBatchCommitRef, VerifiedMergeMembershipHeadActivation>,
71    conflict_resolutions: BTreeMap<
72        protocol_membership::StoreMembershipConflictResolutionRef,
73        VerifiedMergeConflictResolutionActivation,
74    >,
75}
76
77#[derive(Clone, Copy, Debug, PartialEq, Eq)]
78pub(crate) enum VerifiedMergePrefixHeadStatus {
79    Included,
80    OutsidePrefix,
81}
82
83impl VerifiedMergeMembershipPrefix {
84    pub(super) fn extends(&self, verified: &Self) -> bool {
85        verified.commits.is_subset(&self.commits)
86    }
87
88    pub(crate) fn from_retained(
89        checkpoints: &[coven_database::RetainedMergeHistoryCheckpoint],
90    ) -> Result<Self, StorePullError> {
91        let mut prefix = Self::default();
92        for checkpoint in checkpoints {
93            match checkpoint {
94                coven_database::RetainedMergeHistoryCheckpoint::Snapshot(checkpoint) => {
95                    prefix
96                        .commits
97                        .extend(checkpoint.summary.causal_cut.values().cloned());
98                    for proof in checkpoint.summary.membership_proofs.values() {
99                        prefix.insert_retained_proof(proof)?;
100                    }
101                }
102                coven_database::RetainedMergeHistoryCheckpoint::Commit(materialization) => {
103                    prefix.commits.insert(materialization.commit_ref().clone());
104                    if let Some(proof) = &materialization.history_evidence().membership_proof {
105                        prefix.insert_retained_proof(proof)?;
106                    }
107                }
108            }
109        }
110        Ok(prefix)
111    }
112
113    fn insert_retained_proof(
114        &mut self,
115        proof: &store_commit::RetainedMergeMembershipProof,
116    ) -> Result<(), StorePullError> {
117        let Some(store_commit::StoreControl { transition }) = proof.commit_value.control() else {
118            return Err(StorePullError::InvalidState(
119                "retained Merge membership proof has no membership control".to_string(),
120            ));
121        };
122        let activation = VerifiedMergeMembershipHeadActivation {
123            commit: proof.commit.clone(),
124            transition: transition.clone(),
125        };
126        match self.head_activations.entry(proof.commit.clone()) {
127            std::collections::btree_map::Entry::Vacant(entry) => {
128                entry.insert(activation);
129            }
130            std::collections::btree_map::Entry::Occupied(entry) if entry.get() == &activation => {}
131            std::collections::btree_map::Entry::Occupied(_) => {
132                return Err(StorePullError::InvalidState(
133                    "retained checkpoints disagree on a membership activation".to_string(),
134                ));
135            }
136        }
137        if let Some(reference) = &proof.resolution {
138            let activation = VerifiedMergeConflictResolutionActivation {
139                reference: reference.clone(),
140            };
141            match self.conflict_resolutions.entry(reference.clone()) {
142                std::collections::btree_map::Entry::Vacant(entry) => {
143                    entry.insert(activation);
144                }
145                std::collections::btree_map::Entry::Occupied(entry)
146                    if entry.get() == &activation => {}
147                std::collections::btree_map::Entry::Occupied(_) => {
148                    return Err(StorePullError::InvalidState(
149                        "retained checkpoints disagree on a conflict resolution".to_string(),
150                    ));
151                }
152            }
153        }
154        Ok(())
155    }
156
157    pub(crate) fn head_activation(
158        &self,
159        commit: &StoreBatchCommitRef,
160    ) -> Option<&VerifiedMergeMembershipHeadActivation> {
161        self.head_activations.get(commit)
162    }
163
164    pub(crate) fn verifies_conflict_resolution(
165        &self,
166        reference: &protocol_membership::StoreMembershipConflictResolutionRef,
167    ) -> bool {
168        self.conflict_resolutions
169            .get(reference)
170            .is_some_and(|proof| proof.verifies(reference))
171    }
172
173    pub(crate) fn classify_head(
174        &self,
175        reference: &protocol_membership::MembershipHeadRef,
176        head: &protocol_membership::AuthorHead,
177        commit: &StoreBatchCommitRef,
178    ) -> Result<VerifiedMergePrefixHeadStatus, StorePullError> {
179        if !self.commits.contains(commit) {
180            return Ok(VerifiedMergePrefixHeadStatus::OutsidePrefix);
181        }
182        let proof = self.head_activations.get(commit).ok_or_else(|| {
183            StorePullError::InvalidState(
184                "in-prefix membership activation is absent from its verified Store control"
185                    .to_string(),
186            )
187        })?;
188        if !proof.verifies(reference, head, commit) {
189            return Err(StorePullError::InvalidState(
190                "membership head differs from its in-prefix verified Store control".to_string(),
191            ));
192        }
193        Ok(VerifiedMergePrefixHeadStatus::Included)
194    }
195
196    pub(crate) fn validate_complete_membership(
197        &self,
198        membership: &MembershipChain,
199    ) -> Result<(), StorePullError> {
200        if self
201            .predecessor_memberships
202            .iter()
203            .any(|predecessor| !membership.causally_includes(predecessor))
204        {
205            return Err(StorePullError::InvalidState(
206                "membership state regresses below an exact Store predecessor membership"
207                    .to_string(),
208            ));
209        }
210        if self
211            .head_activations
212            .values()
213            .any(|proof| !membership.contains_coord(&proof.transition.body.entry.coord))
214        {
215            return Err(StorePullError::InvalidState(
216                "membership state omits an accepted Store membership control".to_string(),
217            ));
218        }
219        if self.conflict_resolutions.keys().any(|reference| {
220            membership
221                .resolution_refs()
222                .binary_search(reference)
223                .is_err()
224        }) {
225            return Err(StorePullError::InvalidState(
226                "membership state omits an accepted Store conflict resolution".to_string(),
227            ));
228        }
229        Ok(())
230    }
231}
232
233/// The membership authority a commit's predecessors establish, down to the
234/// installed baseline.
235///
236/// A covered predecessor contributes its position and nothing else: the
237/// baseline's own membership floor is what stands behind it, and the control
238/// activations under it were validated when the image that restates them was
239/// verified.
240pub(crate) fn verified_merge_membership_prefix(
241    history: &VerifiedMergeHistory,
242    tips: impl IntoIterator<Item = StoreBatchCommitRef>,
243) -> Result<VerifiedMergeMembershipPrefix, StorePullError> {
244    let closure = verified_merge_commit_closure(history, tips)?;
245    let mut prefix = VerifiedMergeMembershipPrefix {
246        commits: closure.clone(),
247        ..VerifiedMergeMembershipPrefix::default()
248    };
249    for reference in closure {
250        let Some(verified) = history.commits.get(&reference) else {
251            continue;
252        };
253        prefix
254            .predecessor_memberships
255            .push(verified.predecessor_membership.clone());
256        if let Some(control) = &verified.membership_control {
257            prefix
258                .head_activations
259                .insert(reference, control.head_activation.clone());
260            if let Some(resolution) = &control.conflict_resolution {
261                prefix
262                    .conflict_resolutions
263                    .insert(resolution.reference.clone(), resolution.clone());
264            }
265        }
266    }
267    Ok(prefix)
268}
269
270impl<'a> MergeHistoryVerifier<'a> {
271    pub(crate) async fn verify_membership_control_with_retained_history(
272        &mut self,
273        commit_ref: &StoreBatchCommitRef,
274        commit: &StoreBatchCommit,
275        predecessor_membership: &MembershipChain,
276        predecessor_state: &ResolvedStoreDeviceState,
277        pending_resolution: Option<&VerifiedMergeConflictResolutionActivation>,
278    ) -> Result<
279        (
280            VerifiedCircleActivations,
281            Option<VerifiedMergeConflictResolutionActivation>,
282        ),
283        StorePullError,
284    > {
285        let Some(store_commit::StoreControl { transition }) = commit.control() else {
286            return Err(StorePullError::InvalidState(
287                "Merge membership verifier received another Store control".to_string(),
288            ));
289        };
290        let root = self.root.reference().clone();
291        let state = &commit.membership_state;
292        let commit_author = self
293            .commit_verifier
294            .load_registration(&commit.author_registration)
295            .await?;
296        if transition.body.author_registration != commit.author_registration
297            || transition.body.entry.coord.author_pubkey != commit_author.value.author_pubkey
298            || transition.body.resolutions != state.resolutions
299            || transition.body.successor.predecessor
300                != transition
301                    .body
302                    .predecessor
303                    .as_ref()
304                    .map(|reference| reference.object.clone())
305        {
306            return Err(StorePullError::InvalidState(
307                "Merge membership transition differs from its Store authority".to_string(),
308            ));
309        }
310        match &transition.body.predecessor {
311            Some(predecessor) if state.heads.binary_search(predecessor).is_err() => {
312                return Err(StorePullError::InvalidState(
313                    "Merge membership transition predecessor is absent from its signed state"
314                        .to_string(),
315                ));
316            }
317            None if state.heads.iter().any(|head| {
318                head.coord.stream_key() == transition.body.entry.coord.stream_key()
319            }) =>
320            {
321                return Err(StorePullError::InvalidState(
322                    "first Merge membership transition has an existing signed predecessor"
323                        .to_string(),
324                ));
325            }
326            _ => {}
327        }
328        let opened_entry = self
329            .commit_verifier
330            .membership_objects()
331            .load_entry(&transition.body.entry)
332            .await?;
333        if opened_entry.value.coord() != transition.body.entry.coord
334            || opened_entry.value.dependencies != predecessor_membership.effective_frontier()
335            || opened_entry.value.resolution_dependencies != transition.body.resolutions
336        {
337            return Err(StorePullError::InvalidState(
338                "Merge membership transition differs from its exact entry".to_string(),
339            ));
340        }
341        if let protocol_membership::MembershipChange::RemoveMember {
342            user_pubkey,
343            removes,
344            retirement_device_state,
345            ..
346        } = &opened_entry.value.change
347        {
348            let removes_exact_member =
349                removes == &predecessor_membership.active_grant_ids(user_pubkey);
350            let retires_owner = removes.iter().any(|grant| {
351                predecessor_membership
352                    .active_grant(grant)
353                    .is_some_and(|record| {
354                        matches!(
355                            record.role,
356                            protocol_membership::StoreMembershipRoleGrant::Owner { .. }
357                        )
358                    })
359            });
360            if !removes_exact_member
361                || !retires_owner
362                || retirement_device_state.as_ref() != Some(&commit.device_state)
363                || !commit.stream_activations().is_empty()
364            {
365                return Err(StorePullError::InvalidState(
366                    "Merge Owner-removal control differs from its exact membership entry"
367                        .to_string(),
368                ));
369            }
370            let mut successor_membership = predecessor_membership.clone();
371            successor_membership.add_entry(opened_entry.value)?;
372            return VerifiedCircleActivations::membership_control(commit, commit_ref)
373                .map(|activations| (activations, None))
374                .map_err(StorePullError::from);
375        }
376        if let protocol_membership::MembershipChange::ResolutionActivation { resolution } =
377            &opened_entry.value.change
378        {
379            let resolution = resolution.clone();
380            let resolution_proof = pending_resolution
381                .filter(|proof| proof.verifies(&resolution))
382                .ok_or_else(|| {
383                    StorePullError::InvalidState(
384                        "Merge conflict resolution lacks its verified Store activation".to_string(),
385                    )
386                })?
387                .clone();
388            let opened_resolution = self
389                .commit_verifier
390                .membership_objects()
391                .load_resolution(&resolution)
392                .await?;
393            let acceptance = &opened_resolution.value.replacement_acceptance;
394            let mut expected = vec![
395                store_commit::StreamActivation::grant_authorized(
396                    root.store_root_hash,
397                    acceptance.owner_registration.clone(),
398                    opened_resolution.value.replacement_grant.clone(),
399                    acceptance.membership.clone(),
400                ),
401                store_commit::StreamActivation::grant_authorized(
402                    root.store_root_hash,
403                    acceptance.owner_registration.clone(),
404                    opened_resolution.value.replacement_grant.clone(),
405                    acceptance.recovery.clone(),
406                ),
407            ];
408            expected.sort();
409            if transition.body.predecessor.is_some()
410                || transition
411                    .body
412                    .resolutions
413                    .binary_search(&resolution)
414                    .is_err()
415                || commit.stream_activations() != expected
416            {
417                return Err(StorePullError::InvalidState(
418                    "Merge conflict-resolution control differs from its exact membership entry"
419                        .to_string(),
420                ));
421            }
422            let mut successor_membership = predecessor_membership.clone();
423            successor_membership.add_entry(opened_entry.value)?;
424            return VerifiedCircleActivations::membership_control(commit, commit_ref)
425                .map(|activations| (activations, Some(resolution_proof)))
426                .map_err(StorePullError::from);
427        }
428        let protocol_membership::MembershipChange::SetMember {
429            user_pubkey,
430            role:
431                protocol_membership::StoreMembershipRoleGrant::Owner {
432                    recovery: protocol_membership::OwnerRecoveryAnchorRef::Promotion { acceptance },
433                },
434            grant_id,
435            membership: Some(membership_anchor),
436            replaces,
437            retirement_device_state,
438            ..
439        } = &opened_entry.value.change
440        else {
441            return Err(StorePullError::InvalidState(
442                "Merge membership control does not activate one Owner promotion".to_string(),
443            ));
444        };
445        if retirement_device_state.is_some()
446            || user_pubkey != &acceptance.request.member_pubkey
447            || grant_id != &acceptance.request.intended_owner_grant
448            || replaces != &BTreeSet::from([acceptance.request.member_grant.clone()])
449            || acceptance.request.promoter_registration != commit.author_registration
450        {
451            return Err(StorePullError::InvalidState(
452                "Merge Owner-promotion control differs from its exact membership entry".to_string(),
453            ));
454        }
455        self.verify_owner_promotion_acceptance_in_loaded_history(acceptance)
456            .await?;
457        let request_activation = acceptance.activation.commit();
458        let request_commit = self
459            .history
460            .commits
461            .get(request_activation)
462            .ok_or_else(|| {
463                StorePullError::InvalidState(
464                    "Merge Owner-promotion request activation is absent from its verified history"
465                        .to_string(),
466                )
467            })?;
468        let verified_membership_activations = verified_merge_membership_prefix(
469            &self.history,
470            commit_predecessor_references(request_commit.verified.value()),
471        )?;
472        let request_membership = self
473            .load_membership_at_verified_prefix(
474                &acceptance.request.predecessor_membership.heads,
475                &acceptance.request.predecessor_membership.resolutions,
476                &verified_membership_activations,
477                None,
478            )
479            .await?;
480        let predecessor_cut = commit.order.predecessor_cut()?;
481        let predecessor_frontier = predecessor_cut.commits();
482        let request_stream = request_activation.coord.stream_id;
483        let activation_is_covered = predecessor_frontier
484            .get(&request_stream)
485            .is_some_and(|head| head.coord.sequence() >= request_activation.coord.sequence());
486        let promoter_is_active = device_state_has_active_registration(
487            predecessor_state,
488            &acceptance.request.promoter_registration,
489        );
490        let candidate_is_active = device_state_has_active_registration(
491            predecessor_state,
492            &acceptance.request.member_registration,
493        );
494        let promoter_grant_is_active = predecessor_membership
495            .active_owner_grant(&commit_author.value.author_pubkey)
496            .as_ref()
497            == Some(&acceptance.request.promoter_owner_grant);
498        let candidate_grant_is_active = predecessor_membership
499            .active_grant(&acceptance.request.member_grant)
500            .is_some_and(|record| {
501                record.member_pubkey == acceptance.request.member_pubkey
502                    && record.role == protocol_membership::StoreMembershipRoleGrant::Member
503            });
504        if !predecessor_membership.causally_includes(&request_membership)
505            || !activation_is_covered
506            || !promoter_is_active
507            || !candidate_is_active
508            || !promoter_grant_is_active
509            || !candidate_grant_is_active
510        {
511            return Err(StorePullError::InvalidState(
512                "Merge Owner-promotion transition does not include its accepted authority"
513                    .to_string(),
514            ));
515        }
516        let store_commit::OwnerPromotionAnchors {
517            membership,
518            recovery,
519        } = &acceptance.anchors;
520        if membership != membership_anchor {
521            return Err(StorePullError::InvalidState(
522                "Merge Owner-promotion entry carries another membership anchor".to_string(),
523            ));
524        }
525        let mut expected = vec![
526            store_commit::StreamActivation::grant_authorized(
527                root.store_root_hash,
528                acceptance.request.member_registration.clone(),
529                acceptance.request.intended_owner_grant.clone(),
530                membership.clone(),
531            ),
532            store_commit::StreamActivation::grant_authorized(
533                root.store_root_hash,
534                acceptance.request.member_registration.clone(),
535                acceptance.request.intended_owner_grant.clone(),
536                recovery.clone(),
537            ),
538        ];
539        expected.sort();
540        if commit.stream_activations() != expected {
541            return Err(StorePullError::InvalidState(
542                "Merge Owner-promotion control carries different stream activations".to_string(),
543            ));
544        }
545        VerifiedCircleActivations::membership_control(commit, commit_ref)
546            .map(|activations| (activations, None))
547            .map_err(StorePullError::from)
548    }
549
550    pub(crate) async fn verified_membership_objects(
551        &self,
552        commit_ref: &StoreBatchCommitRef,
553        commit: &StoreBatchCommit,
554    ) -> Result<Option<VerifiedMergeMembershipClosure>, StorePullError> {
555        self.commit_verifier
556            .verified_merge_membership_objects(commit_ref, commit)
557            .await
558    }
559
560    pub(crate) async fn verify_accepted_provider_access_activation(
561        &mut self,
562        access: &coven_protocol::provider::ActivatedStoreMemberProviderAccessGrant,
563        provider_admin: &coven_protocol::provider::ProviderAdminGrantRecord,
564        administrator: &StoreDeviceRegistration,
565    ) -> Result<(), StorePullError> {
566        let grant = self
567            .load_provider_access_grant(&access.grant_ref, administrator)
568            .await?;
569        if grant.value != access.grant {
570            return Err(StorePullError::InvalidState(
571                "device provider approval embeds a different access grant than its exact reference"
572                    .to_string(),
573            ));
574        }
575        let activation = self.load_ref(&access.activation).await?;
576        if activation.value().provider_access_grants() != std::slice::from_ref(&access.grant_ref)
577            || activation.value().author_registration != access.grant.administrator
578            || activation.author() != administrator
579        {
580            return Err(StorePullError::InvalidState(
581                "device provider approval activation is not the administrator's exact sole access grant"
582                    .to_string(),
583            ));
584        }
585        let membership = self
586            .load_predecessor_membership(&activation.value().membership_state)
587            .await
588            .map_err(StorePullError::from)?;
589        if !predecessor_verifies_provider_administrator(
590            &membership,
591            &access.grant.administrator_grant,
592            &activation.value().author_registration,
593            provider_admin,
594        ) {
595            return Err(StorePullError::InvalidState(
596                "device provider approval activation lacks exact predecessor provider-administrator authority"
597                    .to_string(),
598            ));
599        }
600        if !self
601            .current_history_contains(&membership, &access.activation)
602            .await?
603        {
604            return Err(StorePullError::InvalidState(
605                "device provider approval activation is absent from current accepted Store history"
606                    .to_string(),
607            ));
608        }
609        Ok(())
610    }
611
612    async fn current_history_contains(
613        &mut self,
614        membership: &MembershipChain,
615        expected: &StoreBatchCommitRef,
616    ) -> Result<bool, StorePullError> {
617        self.verify_refs([expected.clone()]).await?;
618        let state = self
619            .history
620            .commits
621            .get(expected)
622            .ok_or_else(|| {
623                StorePullError::InvalidState(
624                    "provider-access activation is absent from its verified Merge graph"
625                        .to_string(),
626                )
627            })?
628            .state_after
629            .clone();
630        let authority = self.current_merge_authority_from(membership, state).await?;
631        let accepted_closure = verified_merge_commit_closure(
632            &self.history,
633            authority.cut.commits().values().cloned(),
634        )?;
635        Ok(accepted_closure.contains(expected))
636    }
637
638    pub(super) async fn current_merge_authority(
639        &mut self,
640        membership: &MembershipChain,
641    ) -> Result<CurrentMergeAuthority, StorePullError> {
642        self.current_merge_authority_from(membership, self.history.genesis.clone())
643            .await
644    }
645
646    pub(crate) async fn current_merge_authority_cut(
647        &mut self,
648        membership: &MembershipChain,
649    ) -> Result<StoreHistoryCut, StorePullError> {
650        self.current_merge_authority(membership)
651            .await
652            .map(|authority| authority.cut)
653    }
654
655    async fn current_merge_authority_from(
656        &mut self,
657        membership: &MembershipChain,
658        mut state: ResolvedStoreDeviceState,
659    ) -> Result<CurrentMergeAuthority, StorePullError> {
660        let mut registrations = BTreeMap::new();
661        let founder = self.commit_verifier.load_founder_registration().await?;
662        let founder_ref =
663            StoreDeviceRegistrationRef::from_registration(&founder.value, founder.object);
664        registrations.insert(
665            founder_ref.device_id,
666            ReferencedStoreDeviceRegistration::verified(founder_ref, founder.value)
667                .map_err(StorePullError::Protocol)?,
668        );
669        for recovered in self.discover_owner_recoveries(membership).await? {
670            registrations.insert(recovered.reference().device_id, recovered);
671        }
672        self.load_state_registrations(&state, &mut registrations)
673            .await?;
674
675        let mut observed_states = BTreeSet::new();
676        loop {
677            let mut next = BTreeMap::new();
678            for registration in registrations.values() {
679                let registration_ref = registration.reference();
680                let inactive_cut = match state.devices.get(&registration_ref.device_id) {
681                    Some(record) if record.registration != *registration_ref => {
682                        return Err(StorePullError::InvalidState(
683                            "current Merge device state names another registration revision"
684                                .to_string(),
685                        ));
686                    }
687                    Some(record) => match &record.status {
688                        StoreDeviceStatus::Active => None,
689                        StoreDeviceStatus::Inactive { accepted_cut, .. } => Some(accepted_cut),
690                    },
691                    None => None,
692                };
693                let discovered = self
694                    .discover_merge_stream(registration_ref, registration.value(), inactive_cut)
695                    .await?;
696                if matches!(discovered.block, Some(MergeStreamBlock::Authenticated(_))) {
697                    return Err(StorePullError::InvalidState(
698                        "an authenticated Merge stream position cannot be verified".to_string(),
699                    ));
700                }
701                if let Some(reference) = discovered
702                    .commits
703                    .last()
704                    .map(|(_, _, reference, _)| reference)
705                    .or_else(|| {
706                        self.commit_verifier
707                            .covered_announcement_commit(registration_ref)
708                    })
709                {
710                    let stream_id = reference.coord.stream_id;
711                    next.insert(stream_id, reference.clone());
712                }
713            }
714            self.verify_refs(next.values().cloned()).await?;
715            let next_state = if next.is_empty() {
716                self.history.genesis.clone()
717            } else {
718                ResolvedStoreDeviceState::merge(
719                    next.values()
720                        .map(|reference| {
721                            self.history.state_after(reference).cloned().ok_or_else(|| {
722                                StorePullError::InvalidState(
723                                    "current Merge frontier is absent from its verified graph"
724                                        .to_string(),
725                                )
726                            })
727                        })
728                        .collect::<Result<Vec<_>, _>>()?,
729                )
730                .map_err(StorePullError::Protocol)?
731            };
732            let registration_count = registrations.len();
733            self.load_state_registrations(&next_state, &mut registrations)
734                .await?;
735            let stable = next_state == state && registrations.len() == registration_count;
736            if stable {
737                return Ok(CurrentMergeAuthority {
738                    cut: StoreHistoryCut(next),
739                    state: next_state,
740                    registrations,
741                });
742            }
743            let state_fingerprint = ObjectHash::digest(
744                &serde_json::to_vec(&(&next, &next_state))
745                    .map_err(StorePullError::Serialization)?,
746            );
747            if !observed_states.insert(state_fingerprint) {
748                return Err(StorePullError::InvalidState(
749                    "current Merge authority discovery does not reach one stable frontier"
750                        .to_string(),
751                ));
752            }
753            state = next_state;
754        }
755    }
756
757    pub async fn load_exact_anchored_membership(
758        &mut self,
759        heads: &[protocol_membership::MembershipHeadRef],
760        owner: Option<&str>,
761    ) -> Result<MembershipChain, crate::sync::store::membership::AnchoredChainError> {
762        self.load_exact_anchored_membership_traversal(heads, owner)
763            .await
764            .map(|(membership, _)| membership)
765    }
766
767    /// The anchored walk, and every membership object it read on the way.
768    ///
769    /// A snapshot publisher needs the second half: what it publishes as the
770    /// membership rollup is exactly the set of objects a reader of this same
771    /// frontier would otherwise fetch one at a time.
772    pub(crate) async fn load_exact_anchored_membership_traversal(
773        &mut self,
774        heads: &[protocol_membership::MembershipHeadRef],
775        owner: Option<&str>,
776    ) -> Result<
777        (MembershipChain, membership::TraversedMembership),
778        crate::sync::store::membership::AnchoredChainError,
779    > {
780        let (membership, traversed) = membership::HistoryMembershipActivation::new(self)
781            .load_exact_anchored_chain(heads, owner)
782            .await?;
783        if self.history.commits.is_empty() {
784            let authority = VerifiedMergeMembershipPrefix::default();
785            authority
786                .validate_complete_membership(&membership)
787                .map_err(crate::sync::store::membership::AnchoredChainError::from)?;
788            self.remember_verified_membership(authority, membership.clone());
789        }
790        Ok((membership, traversed))
791    }
792
793    pub(crate) async fn load_membership_at_exact_heads(
794        &mut self,
795        heads: &[protocol_membership::MembershipHeadRef],
796        resolutions: &[protocol_membership::StoreMembershipConflictResolutionRef],
797    ) -> Result<MembershipChain, crate::sync::store::membership::AnchoredChainError> {
798        membership::HistoryMembershipActivation::new(self)
799            .load_at_exact_heads(heads, resolutions)
800            .await
801    }
802
803    pub(crate) async fn load_membership_at_verified_prefix(
804        &self,
805        heads: &[protocol_membership::MembershipHeadRef],
806        resolutions: &[protocol_membership::StoreMembershipConflictResolutionRef],
807        verified_activations: &VerifiedMergeMembershipPrefix,
808        pending_resolution: Option<&VerifiedMergeConflictResolutionActivation>,
809    ) -> Result<MembershipChain, crate::sync::store::membership::AnchoredChainError> {
810        VerifiedPrefixMembershipActivation::new(
811            &self.root,
812            &self.commit_verifier,
813            verified_activations,
814        )
815        .load_at_exact_heads(heads, resolutions, pending_resolution)
816        .await
817    }
818
819    pub(crate) async fn load_predecessor_membership(
820        &mut self,
821        state: &StoreMembershipStateRef,
822    ) -> Result<MembershipChain, RegistrationLoadError> {
823        self.load_membership_at_exact_heads(&state.heads, &state.resolutions)
824            .await
825            .map_err(RegistrationLoadError::from)
826    }
827
828    pub(crate) async fn load_predecessor_membership_at_verified_prefix(
829        &self,
830        state: &StoreMembershipStateRef,
831        verified_activations: &VerifiedMergeMembershipPrefix,
832        pending_resolution: Option<&VerifiedMergeConflictResolutionActivation>,
833    ) -> Result<MembershipChain, RegistrationLoadError> {
834        self.load_membership_at_verified_prefix(
835            &state.heads,
836            &state.resolutions,
837            verified_activations,
838            pending_resolution,
839        )
840        .await
841        .map_err(RegistrationLoadError::from)
842    }
843
844    pub(crate) async fn load_exact_membership_head(
845        &mut self,
846        reference: &protocol_membership::MembershipHeadRef,
847    ) -> Result<protocol_membership::AuthorHead, crate::sync::store::membership::AnchoredChainError>
848    {
849        self.commit_verifier
850            .membership_objects()
851            .load_head(reference)
852            .await
853            .map(|loaded| loaded.value)
854            .map_err(membership::map_membership_object_error)
855    }
856
857    pub(crate) async fn project_membership_to_verified_prefix(
858        &self,
859        candidate_heads: &[protocol_membership::MembershipHeadRef],
860        prefix: &VerifiedMergeMembershipPrefix,
861    ) -> Result<MembershipChain, crate::sync::store::membership::AnchoredChainError> {
862        VerifiedPrefixMembershipActivation::new(&self.root, &self.commit_verifier, prefix)
863            .project(candidate_heads)
864            .await
865    }
866
867    pub(crate) async fn verify_membership_grant_revocation_nonactivation(
868        &mut self,
869        grant_id: &protocol_membership::MembershipGrantId,
870        membership: &StoreMembershipStateRef,
871        activation_commit: &StoreBatchCommitRef,
872        activation_head: &store_commit::StoreDeviceHeadRef,
873        candidate: &VerifiedStoreBatchCommit,
874        candidate_head: &StoreDeviceHead,
875        candidate_head_object: &ExactObjectRef,
876    ) -> Result<remote_object::VerifiedCandidateNonactivation, StorePullError> {
877        let root = self.root.reference().clone();
878        let head_prefix =
879            store_commit::semantic_prefix_from_exact_object(&activation_head.object, ".json")
880                .map_err(StorePullError::Protocol)?;
881        let context = ProtocolObjectContext::signed_plaintext(
882            root.store_root_hash,
883            ProtocolObjectDomain::StoreHead,
884        );
885        let head_bytes = self
886            .commit_verifier
887            .read_protocol_object(&context, &activation_head.object, &head_prefix)
888            .await?;
889        activation_head.object.verify(&head_bytes)?;
890        let witness_head: StoreDeviceHead =
891            serde_json::from_slice(&head_bytes).map_err(|error| {
892                StorePullError::context("membership revocation witness head", error)
893            })?;
894        if witness_head.head_hash() != activation_head.head_hash
895            || &witness_head.commit != activation_commit
896        {
897            return Err(StorePullError::InvalidState(
898                "membership revocation witness head differs from its exact activation".to_string(),
899            ));
900        }
901        let witness_author = self
902            .commit_verifier
903            .load_registration(&witness_head.author_registration)
904            .await?;
905        let opened = self
906            .commit_verifier
907            .load_head(activation_head, &witness_author.value, &witness_head.commit)
908            .await?;
909        self.verify_refs([witness_head.commit.clone()]).await?;
910        let witness_commit = self
911            .history
912            .commits
913            .get(&witness_head.commit)
914            .ok_or_else(|| {
915                StorePullError::InvalidState(
916                    "membership revocation witness is absent from its verified history".to_string(),
917                )
918            })?
919            .verified
920            .clone();
921        if witness_commit.author() != &witness_author.value {
922            return Err(StorePullError::InvalidState(
923                "membership revocation witness commit belongs to another author".to_string(),
924            ));
925        }
926        let (_, exact_head) = self
927            .commit_verifier
928            .exact_next_announcement_slot(
929                &witness_head.author_registration,
930                &witness_author.value,
931                Some(&witness_commit),
932            )
933            .await
934            .map_err(|error| StorePullError::Store(Box::new(error)))?;
935        if exact_head.as_ref() != Some(activation_head) || opened.value != witness_head {
936            return Err(StorePullError::InvalidState(
937                "membership revocation witness is not an accepted exact head".to_string(),
938            ));
939        }
940        if witness_commit.value().membership_state != *membership {
941            return Err(StorePullError::InvalidState(
942                "membership revocation witness commit names another membership state".to_string(),
943            ));
944        }
945        let current_membership = self
946            .load_predecessor_membership(&witness_commit.value().membership_state)
947            .await
948            .map_err(StorePullError::from)?;
949        let MembershipStatus::Resolved(current) = current_membership.status() else {
950            return Err(StorePullError::InvalidState(
951                "membership revocation witness state is conflicted".to_string(),
952            ));
953        };
954        let Some(causal_grants::GrantState::Tombstoned {
955            record: current_record,
956            ..
957        }) = current.grants.get(grant_id)
958        else {
959            return Err(StorePullError::InvalidState(
960                "membership revocation witness grant is not tombstoned".to_string(),
961            ));
962        };
963        let candidate_ref = candidate.reference();
964        let candidate_commit = candidate.value();
965        let candidate_author = candidate.author();
966        let predecessor_membership = self
967            .load_predecessor_membership(&candidate_commit.membership_state)
968            .await
969            .map_err(StorePullError::from)?;
970        let MembershipStatus::Resolved(predecessor) = predecessor_membership.status() else {
971            return Err(StorePullError::InvalidState(
972                "membership revocation candidate predecessor is conflicted".to_string(),
973            ));
974        };
975        let Some(predecessor_record) = predecessor.active_grant(grant_id) else {
976            return Err(StorePullError::InvalidState(
977                "membership revocation grant was not active at the candidate predecessor"
978                    .to_string(),
979            ));
980        };
981        if predecessor_record != current_record
982            || predecessor_record.member_pubkey != candidate_author.author_pubkey
983            || candidate_commit.membership_authority.as_ref()
984                != Some(&predecessor_record.creation_authority)
985        {
986            return Err(StorePullError::InvalidState(
987                "membership revocation grant differs from the candidate's signed authority"
988                    .to_string(),
989            ));
990        }
991        let cap = witness_commit
992            .value()
993            .order
994            .predecessor_cut()
995            .map_err(StorePullError::Protocol)?;
996        let expected_stream = store_commit::StreamActivation::device_authorized_stream_id(
997            root.store_root_hash,
998            &candidate_commit.author_registration,
999            store_commit::StreamAnchorDomain::StoreAnnouncements,
1000        );
1001        let StoreCommitCoord {
1002            stream_id,
1003            sequence,
1004        } = candidate_ref.coord;
1005        if stream_id != expected_stream
1006            || cap
1007                .commits()
1008                .get(&expected_stream)
1009                .is_some_and(|covered| sequence <= covered.coord.sequence())
1010        {
1011            return Err(StorePullError::InvalidState(
1012                "membership revocation candidate is not beyond the accepted witness cut"
1013                    .to_string(),
1014            ));
1015        }
1016        let verified_candidate_head = self
1017            .commit_verifier
1018            .verify_terminal_candidate_head(candidate, candidate_head, candidate_head_object)
1019            .await?;
1020        let durable = remote_object::CandidateNonactivation::from_durable_parts(
1021            candidate_ref,
1022            candidate_commit,
1023            remote_object::CandidateNonactivationProof::MergeMembershipGrantRevocation {
1024                grant_id: grant_id.clone(),
1025                membership: membership.clone(),
1026                activation_commit: witness_head.commit.clone(),
1027                activation_head: activation_head.clone(),
1028            },
1029        )
1030        .map_err(StorePullError::RemoteObject)?;
1031        remote_object::VerifiedCandidateNonactivation::from_verified_membership_grant_revocation(
1032            durable,
1033            candidate_ref.clone(),
1034            verified_candidate_head,
1035        )
1036        .map_err(StorePullError::RemoteObject)
1037    }
1038
1039    #[cfg(any(test, feature = "test-utils"))]
1040    pub(crate) async fn assert_deep_membership_projection(
1041        &mut self,
1042        heads: &[protocol_membership::MembershipHeadRef],
1043    ) {
1044        membership::HistoryMembershipActivation::new(self)
1045            .assert_deep_valid_predecessor_path_is_iterative(heads)
1046            .await;
1047    }
1048}