Skip to main content

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

1use super::*;
2
3pub struct PreparedMergeHistorySuccessor {
4    pub(crate) history_evidence: store_commit::RetainedMergeCommitEvidence,
5    pub(crate) head_slot: coven_protocol::objects::ObjectSlot,
6    pub(crate) predecessor_head: Option<store_commit::StoreDeviceHeadRef>,
7}
8
9pub struct MergeHistorySuccessorEvidence {
10    pub(crate) registrations: Vec<ReferencedStoreDeviceRegistration>,
11    pub(crate) acknowledgement: Option<store_commit::RetainedVerifiedActivatedAck>,
12    pub(crate) membership_proof: Option<store_commit::RetainedMergeMembershipProof>,
13}
14
15impl MergeHistorySuccessorEvidence {
16    pub(crate) fn none() -> Self {
17        Self {
18            registrations: Vec::new(),
19            acknowledgement: None,
20            membership_proof: None,
21        }
22    }
23}
24
25fn insert_exact<K, V>(
26    target: &mut BTreeMap<K, V>,
27    key: K,
28    value: V,
29    conflict: &str,
30) -> Result<(), StorePullError>
31where
32    K: Ord,
33    V: PartialEq,
34{
35    match target.entry(key) {
36        std::collections::btree_map::Entry::Vacant(entry) => {
37            entry.insert(value);
38            Ok(())
39        }
40        std::collections::btree_map::Entry::Occupied(entry) if entry.get() == &value => Ok(()),
41        std::collections::btree_map::Entry::Occupied(_) => {
42            Err(StorePullError::InvalidState(conflict.to_string()))
43        }
44    }
45}
46
47/// Merge one predecessor summary's acknowledgement chain for a device into the
48/// chain being composed. Two chains for one device must agree — one extending
49/// the other is a longer view of the same history; anything else is a fork.
50pub(crate) fn insert_latest_acknowledgement(
51    target: &mut BTreeMap<store_commit::StoreDeviceId, store_commit::RetainedAcknowledgementChain>,
52    device_id: store_commit::StoreDeviceId,
53    value: store_commit::RetainedAcknowledgementChain,
54) -> Result<(), StorePullError> {
55    match target.entry(device_id) {
56        std::collections::btree_map::Entry::Vacant(entry) => {
57            entry.insert(value);
58            Ok(())
59        }
60        std::collections::btree_map::Entry::Occupied(entry) if entry.get() == &value => Ok(()),
61        std::collections::btree_map::Entry::Occupied(mut entry)
62            if value.exactly_extends(entry.get()) =>
63        {
64            entry.insert(value);
65            Ok(())
66        }
67        std::collections::btree_map::Entry::Occupied(entry)
68            if entry.get().exactly_extends(&value) =>
69        {
70            Ok(())
71        }
72        std::collections::btree_map::Entry::Occupied(_) => Err(StorePullError::InvalidState(
73            "Merge predecessor checkpoints contain forked acknowledgement proof chains".to_string(),
74        )),
75    }
76}
77
78/// Fold the one acknowledgement a retained commit activated into the chain being
79/// composed for its device.
80///
81/// The rows in a cut carry the acknowledgements made within it, which is enough
82/// to identify each device's latest — but not enough to reach sequence one when
83/// the cut starts above it. A summary states contiguity, so the caller completes
84/// each chain by walking it from the latest entry this fold found.
85pub(crate) fn extend_acknowledgement_chain(
86    target: &mut BTreeMap<store_commit::StoreDeviceId, store_commit::RetainedAcknowledgementChain>,
87    device_id: store_commit::StoreDeviceId,
88    activated: &store_commit::RetainedVerifiedActivatedAck,
89    activating_commit_value: &store_commit::StoreBatchCommit,
90) -> Result<(), StorePullError> {
91    let extended = match target.entry(device_id) {
92        std::collections::btree_map::Entry::Vacant(entry) => {
93            entry.insert(store_commit::RetainedAcknowledgementChain::activated(
94                activated,
95                activating_commit_value,
96            ));
97            true
98        }
99        std::collections::btree_map::Entry::Occupied(mut entry) => {
100            entry.get_mut().extend(activated, activating_commit_value)
101        }
102    };
103    if extended {
104        Ok(())
105    } else {
106        Err(StorePullError::InvalidState(
107            "retained acknowledgements fork at one sequence".to_string(),
108        ))
109    }
110}
111
112fn insert_latest_announcement(
113    target: &mut BTreeMap<
114        protocol_membership::AuthorStreamId,
115        store_commit::RetainedAcceptedStoreAnnouncement,
116    >,
117    stream_id: protocol_membership::AuthorStreamId,
118    value: store_commit::RetainedAcceptedStoreAnnouncement,
119) -> Result<(), StorePullError> {
120    match target.entry(stream_id) {
121        std::collections::btree_map::Entry::Vacant(entry) => {
122            entry.insert(value);
123            Ok(())
124        }
125        std::collections::btree_map::Entry::Occupied(entry) if entry.get() == &value => Ok(()),
126        std::collections::btree_map::Entry::Occupied(mut entry)
127            if entry.get().value.commit.coord.sequence() < value.value.commit.coord.sequence() =>
128        {
129            entry.insert(value);
130            Ok(())
131        }
132        std::collections::btree_map::Entry::Occupied(entry)
133            if entry.get().value.commit.coord.sequence() > value.value.commit.coord.sequence() =>
134        {
135            Ok(())
136        }
137        std::collections::btree_map::Entry::Occupied(_) => Err(StorePullError::InvalidState(
138            "Merge predecessor checkpoints contain conflicting announcement heads at one sequence"
139                .to_string(),
140        )),
141    }
142}
143
144pub(crate) struct MergedRetainedMergeHistory {
145    causal_cut: BTreeMap<StoreCommitCoord, StoreBatchCommitRef>,
146    registrations: BTreeMap<store_commit::StoreDeviceId, ReferencedStoreDeviceRegistration>,
147    acknowledgements:
148        BTreeMap<store_commit::StoreDeviceId, store_commit::RetainedAcknowledgementChain>,
149    membership_proofs: BTreeMap<StoreBatchCommitRef, store_commit::RetainedMergeMembershipProof>,
150    announcement_frontier: BTreeMap<
151        protocol_membership::AuthorStreamId,
152        store_commit::RetainedAcceptedStoreAnnouncement,
153    >,
154}
155
156impl MergedRetainedMergeHistory {
157    fn insert_membership_proof(
158        &mut self,
159        reference: StoreBatchCommitRef,
160        value: store_commit::RetainedMergeMembershipProof,
161    ) -> Result<(), StorePullError> {
162        if self
163            .membership_proofs
164            .keys()
165            .any(|existing| existing.coord == reference.coord && existing != &reference)
166        {
167            return Err(StorePullError::InvalidState(
168                "Merge predecessor checkpoints contain conflicting membership proofs at one Store coordinate"
169                    .to_string(),
170            ));
171        }
172        insert_exact(
173            &mut self.membership_proofs,
174            reference,
175            value,
176            "Merge predecessor checkpoints disagree on a membership proof",
177        )
178    }
179}
180
181pub(crate) fn merge_retained_merge_history(
182    root: &StoreRootRef,
183    membership: &MembershipChain,
184    predecessors: Vec<OpenedRetainedMergeHistorySummary>,
185) -> Result<MergedRetainedMergeHistory, StorePullError> {
186    let mut merged = MergedRetainedMergeHistory {
187        causal_cut: BTreeMap::new(),
188        registrations: BTreeMap::new(),
189        acknowledgements: BTreeMap::new(),
190        membership_proofs: BTreeMap::new(),
191        announcement_frontier: BTreeMap::new(),
192    };
193    for predecessor in predecessors {
194        let predecessor_cut = predecessor.summary.causal_cut.clone();
195        if predecessor.summary.store_root_hash != root.store_root_hash {
196            return Err(StorePullError::InvalidState(
197                "Merge predecessor checkpoint belongs to another Store".to_string(),
198            ));
199        }
200        if predecessor
201            .summary
202            .membership_floor
203            .effective_coordinates
204            .iter()
205            .any(|coordinate| !membership.effectively_contains_coord(coordinate))
206            || predecessor
207                .summary
208                .membership_floor
209                .resolutions
210                .iter()
211                .any(|reference| {
212                    membership
213                        .resolution_refs()
214                        .binary_search(reference)
215                        .is_err()
216                })
217        {
218            return Err(StorePullError::InvalidState(
219                "Merge successor membership omits its retained causal floor".to_string(),
220            ));
221        }
222        for (key, value) in predecessor.summary.causal_cut {
223            insert_exact(
224                &mut merged.causal_cut,
225                key,
226                value,
227                "Merge predecessor checkpoints disagree on a Store coordinate",
228            )?;
229        }
230        for (key, value) in predecessor.summary.registrations {
231            insert_exact(
232                &mut merged.registrations,
233                key,
234                value,
235                "Merge predecessor checkpoints disagree on a device registration",
236            )?;
237        }
238        for (key, value) in predecessor.summary.acknowledgements {
239            insert_latest_acknowledgement(&mut merged.acknowledgements, key, value)?;
240        }
241        for (key, mut value) in predecessor.summary.membership_proofs {
242            if predecessor_cut.get(&value.commit.coord) == Some(&value.commit)
243                && value.announcement.is_none()
244            {
245                let stream_id = value.commit.coord.stream_id;
246                value.announcement = predecessor
247                    .announcement_frontier
248                    .get(&stream_id)
249                    .filter(|announcement| announcement.value.commit == value.commit)
250                    .cloned();
251            }
252            merged.insert_membership_proof(key, value)?;
253        }
254        for (key, value) in predecessor.announcement_frontier {
255            insert_latest_announcement(&mut merged.announcement_frontier, key, value)?;
256        }
257    }
258    Ok(merged)
259}
260
261pub(crate) fn compose_merge_snapshot_history_summary(
262    root: &StoreRootRef,
263    coverage: &CommitFrontier,
264    membership: &MembershipChain,
265    state: &ResolvedStoreDeviceState,
266    author_ref: &StoreDeviceRegistrationRef,
267    author: &StoreDeviceRegistration,
268    predecessors: Vec<coven_database::RetainedMergeHistoryCheckpoint>,
269) -> Result<RetainedVerifiedMergeHistorySummary, StorePullError> {
270    let frontier = &coverage.0;
271    let snapshot_predecessors = predecessors
272        .iter()
273        .filter_map(|checkpoint| match checkpoint {
274            coven_database::RetainedMergeHistoryCheckpoint::Snapshot(checkpoint) => {
275                Some(checkpoint.clone())
276            }
277            coven_database::RetainedMergeHistoryCheckpoint::Commit(_) => None,
278        })
279        .collect();
280    let mut merged = merge_retained_merge_history(root, membership, snapshot_predecessors)?;
281    for checkpoint in predecessors {
282        let coven_database::RetainedMergeHistoryCheckpoint::Commit(materialization) = checkpoint
283        else {
284            continue;
285        };
286        insert_snapshot_commit(
287            &mut merged,
288            root,
289            materialization.commit_ref(),
290            materialization.commit(),
291            materialization.verified_commit().author(),
292            materialization.registrations(),
293            materialization.history_evidence(),
294            materialization.activation_head(),
295            materialization.activation_head_object(),
296        )?;
297    }
298    let MergedRetainedMergeHistory {
299        causal_cut,
300        mut registrations,
301        acknowledgements,
302        membership_proofs,
303        announcement_frontier,
304    } = merged;
305    author_ref
306        .verify_registration(author)
307        .map_err(StorePullError::Protocol)?;
308    insert_exact(
309        &mut registrations,
310        author_ref.device_id,
311        ReferencedStoreDeviceRegistration::verified(author_ref.clone(), author.clone())
312            .map_err(StorePullError::Protocol)?,
313        "Merge snapshot author registration conflicts with retained authority",
314    )?;
315    let summary = RetainedVerifiedMergeHistorySummary {
316        version: store_commit::STORE_PROTOCOL_VERSION,
317        store_root_hash: root.store_root_hash,
318        causal_cut,
319        post_state: StoreDeviceStateRef::from_resolved(coverage.clone(), state)
320            .map_err(StorePullError::Protocol)?,
321        membership_floor: store_commit::MembershipCausalFloor::from_membership(membership),
322        registrations,
323        acknowledgements,
324        membership_proofs,
325        announcement_frontier,
326    };
327    // Assembled, not yet valid — see
328    // `validate_composed_snapshot_history_summary`, which the caller runs once
329    // each device's acknowledgement chain has been walked back to sequence one.
330    let _ = frontier;
331    Ok(summary)
332}
333
334#[allow(clippy::too_many_arguments)]
335fn insert_snapshot_commit(
336    merged: &mut MergedRetainedMergeHistory,
337    root: &StoreRootRef,
338    commit_ref: &StoreBatchCommitRef,
339    commit: &StoreBatchCommit,
340    author: &StoreDeviceRegistration,
341    registrations: &[ActivatedStoreDeviceRegistration],
342    evidence: &store_commit::RetainedMergeCommitEvidence,
343    activation_head: &StoreDeviceHead,
344    activation_head_object: &ExactObjectRef,
345) -> Result<(), StorePullError> {
346    if commit.store_root_hash != root.store_root_hash {
347        return Err(StorePullError::InvalidState(
348            "retained Merge commit belongs to another Store".to_string(),
349        ));
350    }
351    insert_exact(
352        &mut merged.causal_cut,
353        commit_ref.coord.clone(),
354        commit_ref.clone(),
355        "retained Merge commits disagree on a Store coordinate",
356    )?;
357    for registration in registrations {
358        insert_exact(
359            &mut merged.registrations,
360            registration.reference().device_id,
361            registration.registration().clone(),
362            "retained Merge commits disagree on a device registration",
363        )?;
364    }
365    let author = ReferencedStoreDeviceRegistration::verified(
366        commit.author_registration.clone(),
367        author.clone(),
368    )
369    .map_err(StorePullError::Protocol)?;
370    insert_exact(
371        &mut merged.registrations,
372        author.reference().device_id,
373        author,
374        "retained Merge commit author conflicts with retained authority",
375    )?;
376    if let Some(acknowledgement) = &evidence.acknowledgement {
377        let device_id = acknowledgement.acknowledgement().0.registration.device_id;
378        extend_acknowledgement_chain(
379            &mut merged.acknowledgements,
380            device_id,
381            acknowledgement,
382            commit,
383        )?;
384    }
385    let announcement = store_commit::RetainedAcceptedStoreAnnouncement {
386        reference: store_commit::StoreDeviceHeadRef {
387            head_hash: activation_head.head_hash(),
388            object: activation_head_object.clone(),
389        },
390        value: activation_head.clone(),
391    };
392    if let Some(proof) = &evidence.membership_proof {
393        let mut proof = proof.clone();
394        proof.announcement = Some(announcement.clone());
395        merged.insert_membership_proof(commit_ref.clone(), *proof)?;
396    }
397    insert_latest_announcement(
398        &mut merged.announcement_frontier,
399        commit_ref.coord.stream_id,
400        announcement,
401    )
402}
403
404/// Recompose a snapshot's history summary from the commits it covers, resuming
405/// at whatever `baseline` restates.
406///
407/// `baseline` is the summary this device's own replay baseline rests on, when
408/// its coverage lies inside the snapshot's. Below it the commits are retired,
409/// so the composition starts from the summary that stands for them instead of
410/// from history the device dropped — the same resume the publisher composes
411/// from. A device on a genesis baseline passes `None` and the walk runs to the
412/// bottom.
413pub(crate) fn compose_verified_merge_snapshot_history_summary<'a>(
414    root: &StoreRootRef,
415    coverage: &CommitFrontier,
416    membership: &MembershipChain,
417    state: &ResolvedStoreDeviceState,
418    author_ref: &StoreDeviceRegistrationRef,
419    author: &StoreDeviceRegistration,
420    baseline: Option<OpenedRetainedMergeHistorySummary>,
421    commits: impl IntoIterator<Item = &'a VerifiedMergeHistoryCommit>,
422) -> Result<RetainedVerifiedMergeHistorySummary, StorePullError> {
423    let mut merged =
424        merge_retained_merge_history(root, membership, baseline.into_iter().collect())?;
425    for verified in commits {
426        insert_snapshot_commit(
427            &mut merged,
428            root,
429            verified.verified.reference(),
430            verified.verified.value(),
431            verified.verified.author(),
432            &verified.registrations,
433            &verified.history_evidence,
434            &verified.activation_head,
435            &verified.activation_head_object,
436        )?;
437    }
438    let MergedRetainedMergeHistory {
439        causal_cut,
440        mut registrations,
441        acknowledgements,
442        membership_proofs,
443        announcement_frontier,
444    } = merged;
445    author_ref
446        .verify_registration(author)
447        .map_err(StorePullError::Protocol)?;
448    insert_exact(
449        &mut registrations,
450        author_ref.device_id,
451        ReferencedStoreDeviceRegistration::verified(author_ref.clone(), author.clone())
452            .map_err(StorePullError::Protocol)?,
453        "Merge snapshot author registration conflicts with retained authority",
454    )?;
455    let summary = RetainedVerifiedMergeHistorySummary {
456        version: store_commit::STORE_PROTOCOL_VERSION,
457        store_root_hash: root.store_root_hash,
458        causal_cut,
459        post_state: StoreDeviceStateRef::from_resolved(coverage.clone(), state)
460            .map_err(StorePullError::Protocol)?,
461        membership_floor: store_commit::MembershipCausalFloor::from_membership(membership),
462        registrations,
463        acknowledgements,
464        membership_proofs,
465        announcement_frontier,
466    };
467    // Assembled, not yet valid: each device's acknowledgement chain still has to
468    // be completed back to sequence one, which needs a walker this function does
469    // not have. `validate_composed_snapshot_history_summary` is the other half
470    // and runs once the caller has completed them.
471    Ok(summary)
472}
473
474/// Check a composed snapshot summary once its acknowledgement chains are whole.
475pub(crate) fn validate_composed_snapshot_history_summary(
476    summary: &RetainedVerifiedMergeHistorySummary,
477    coverage: &CommitFrontier,
478) -> Result<(), StorePullError> {
479    summary
480        .validate_snapshot_baseline()
481        .map_err(StorePullError::Protocol)?;
482    if summary.frontier().map_err(StorePullError::Protocol)? != coverage.0 {
483        return Err(StorePullError::InvalidState(
484            "Merge snapshot history does not exactly cover its signed frontier".to_string(),
485        ));
486    }
487    Ok(())
488}