Skip to main content

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

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