Skip to main content

coven_replication/sync/store/acknowledgements/
mod.rs

1//! Store and Circle acknowledgement publication.
2
3mod circle;
4
5pub(crate) use circle::CircleAcknowledgementReader;
6
7use super::snapshots as snapshot;
8use super::{AuthorizedWriterOperation, StoreError};
9use crate::sync::cycle::SyncCycleFailure;
10use crate::sync::store::commit_publication::LocalStoreWriter;
11use crate::sync::store::commit_verification::merge_history::SelectedReplayBaselineRetirement;
12use coven_database::StoreDatabase;
13use coven_protocol::objects::StoreObjectError;
14use coven_protocol::objects::{ProtocolObjectContext, ProtocolObjectDomain};
15use coven_protocol::store_commit::{
16    ack_slot_prefix, CommitFrontier, StoreAck, StoreSnapshotLocator, SuccessorLink,
17};
18use coven_storage::CloudSyncObjectStorage;
19use std::sync::Arc;
20use tracing::debug;
21
22#[derive(Debug, thiserror::Error)]
23pub enum StoreAckError {
24    #[error("database: {0}")]
25    Database(#[from] coven_database::DbError),
26    #[error("Store protocol: {0}")]
27    Protocol(#[from] coven_protocol::store_commit::StoreProtocolError),
28    #[error("published Store acknowledgement count has no representable successor")]
29    PublishCountExhausted,
30    #[error("{0}")]
31    Object(#[from] StoreObjectError),
32    #[error("outbound Store acknowledgement is invalid: {0}")]
33    InvalidOutbound(String),
34    #[error("outbound Store acknowledgement prepared commit: {0}")]
35    PreparedCommit(#[from] coven_protocol::prepared_commit::PreparedCommitError),
36    #[error("Store acknowledgement activation: {0}")]
37    Outbound(#[from] StoreError),
38    #[error("Store acknowledgement sync cycle: {0}")]
39    SyncCycle(#[source] Box<crate::sync::cycle::SyncCycleFailure>),
40    #[error("Store acknowledgement writer authorization: {0}")]
41    WriterAuthorization(#[source] Box<crate::sync::store::StoreWriterAuthorizationError>),
42    #[error("Store acknowledgement snapshot: {0}")]
43    Snapshot(#[from] snapshot::SnapshotError),
44}
45
46impl From<crate::sync::cycle::SyncCycleFailure> for StoreAckError {
47    fn from(error: crate::sync::cycle::SyncCycleFailure) -> Self {
48        Self::SyncCycle(Box::new(error))
49    }
50}
51
52impl From<crate::sync::store::StoreWriterAuthorizationError> for StoreAckError {
53    fn from(error: crate::sync::store::StoreWriterAuthorizationError) -> Self {
54        Self::WriterAuthorization(Box::new(error))
55    }
56}
57
58pub struct StagedStoreAcknowledgement {
59    pub acknowledgement: Option<StoreAck>,
60}
61
62/// What standing on this device's acknowledged snapshot did, or why it did
63/// nothing.
64///
65/// A decline is a value rather than a swallowed nothing for the same reason the
66/// reclaim report's is: a stage that speaks only when it acts is
67/// indistinguishable from one that is not running, and this one spent weeks
68/// looking exactly like that on a live store.
69#[derive(Debug, Clone, PartialEq, Eq)]
70pub enum ReplayBaselineAdvance {
71    Advanced(coven_database::AdvancedReplayBaseline),
72    Declined(ReplayBaselineDecline),
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub enum ReplayBaselineDecline {
77    /// This device has published no acknowledgement naming a snapshot, so it
78    /// has never said it holds one.
79    NoAcknowledgedSnapshot,
80    /// The device that authored the acknowledged snapshot is no longer an
81    /// activated registration, so its stream is not this device's to read.
82    SnapshotAuthorInactive { generation: u64 },
83    /// The acknowledged snapshot is gone from its author's stream.
84    SnapshotUnavailable { generation: u64 },
85    /// The acknowledged snapshot did not verify as installable now.
86    SnapshotRejected { generation: u64 },
87    /// A current writer has not published an acknowledgement whose Store cut
88    /// covers this snapshot.
89    MissingWriterAcknowledgement {
90        generation: u64,
91        member: String,
92        device_id: String,
93    },
94    /// Current membership has not been named by accepted Store history, so its
95    /// writer set cannot yet license retirement.
96    MembershipNotAccepted { generation: u64 },
97    /// A published Owner recovery registration has not yet activated its
98    /// replacement writer, so retiring its predecessor history would strand it.
99    PendingOwnerRecovery {
100        generation: u64,
101        member: String,
102        device_id: String,
103    },
104    /// Current accepted history applies a commit outside the snapshot cut
105    /// before a commit inside it, so the cut cannot become a replay baseline.
106    NonPrefixCut { generation: u64 },
107    /// The steady state: the baseline already restates everything the
108    /// acknowledged snapshot does.
109    BaselineAtCoverage { generation: u64 },
110}
111
112impl ReplayBaselineDecline {
113    pub fn as_str(&self) -> &'static str {
114        match self {
115            Self::NoAcknowledgedSnapshot => "this device has acknowledged no snapshot",
116            Self::SnapshotAuthorInactive { .. } => "the acknowledged snapshot's author is inactive",
117            Self::SnapshotUnavailable { .. } => "the acknowledged snapshot is gone from its stream",
118            Self::SnapshotRejected { .. } => "the acknowledged snapshot did not verify",
119            Self::MissingWriterAcknowledgement { .. } => {
120                "a current writer has not crossed the acknowledged snapshot"
121            }
122            Self::MembershipNotAccepted { .. } => {
123                "current membership is not yet in accepted Store history"
124            }
125            Self::PendingOwnerRecovery { .. } => {
126                "an Owner recovery is waiting to activate its replacement writer"
127            }
128            Self::NonPrefixCut { .. } => {
129                "the acknowledged snapshot is not a prefix of accepted replay order"
130            }
131            Self::BaselineAtCoverage { .. } => "the baseline already covers it",
132        }
133    }
134
135    pub fn generation(&self) -> Option<u64> {
136        match self {
137            Self::NoAcknowledgedSnapshot => None,
138            Self::SnapshotAuthorInactive { generation }
139            | Self::SnapshotUnavailable { generation }
140            | Self::SnapshotRejected { generation }
141            | Self::MissingWriterAcknowledgement { generation, .. }
142            | Self::MembershipNotAccepted { generation }
143            | Self::PendingOwnerRecovery { generation, .. }
144            | Self::NonPrefixCut { generation }
145            | Self::BaselineAtCoverage { generation } => Some(*generation),
146        }
147    }
148}
149
150pub(crate) struct AuthorizedAcknowledgements<'operation, 'storage> {
151    writer: &'operation mut AuthorizedWriterOperation<'storage>,
152    database: StoreDatabase,
153    storage: Arc<dyn CloudSyncObjectStorage>,
154    local_writer: Arc<LocalStoreWriter>,
155}
156
157impl<'operation, 'storage> AuthorizedAcknowledgements<'operation, 'storage> {
158    pub(crate) fn new(
159        writer: &'operation mut AuthorizedWriterOperation<'storage>,
160        database: StoreDatabase,
161        storage: Arc<dyn CloudSyncObjectStorage>,
162        local_writer: Arc<LocalStoreWriter>,
163    ) -> Self {
164        Self {
165            writer,
166            database,
167            storage,
168            local_writer,
169        }
170    }
171
172    /// Publish anything queued, then acknowledge where this device now stands.
173    /// Local replay retirement runs independently after acknowledgement
174    /// publication, once every current writer has crossed the selected cut.
175    pub(crate) async fn stage_and_publish(
176        &mut self,
177        sync_time: &str,
178        settled: &crate::sync::store::SettledCycle,
179    ) -> Result<(), SyncCycleFailure> {
180        Box::pin(self.drain_acknowledgements())
181            .await
182            .map_err(|error| {
183                SyncCycleFailure::operation("publish queued Store acknowledgement", error)
184            })?;
185        let frontier =
186            CommitFrontier::from_refs(self.database.materialized_frontier().await.map_err(
187                |error| SyncCycleFailure::operation("read Store acknowledgement frontier", error),
188            )?)
189            .map_err(|error| {
190                SyncCycleFailure::operation("shape Store acknowledgement frontier", error)
191            })?;
192        // Circle acknowledgements first: an outbound Store acknowledgement is what
193        // carries them to the cloud, so the Store one below has to know whether
194        // any are waiting before it decides it has nothing to say.
195        Box::pin(
196            self.writer
197                .circles()
198                .stage_acknowledgements(&frontier, sync_time),
199        )
200        .await
201        .map_err(|error| SyncCycleFailure::operation("stage Circle acknowledgements", error))?;
202        let StagedStoreAcknowledgement { acknowledgement } =
203            Box::pin(self.stage_acknowledgement_against(
204                frontier.clone(),
205                sync_time.to_owned(),
206                Some(settled),
207            ))
208            .await
209            .map_err(|error| SyncCycleFailure::operation("stage Store acknowledgement", error))?;
210        if let Some(acknowledgement) = &acknowledgement {
211            debug!(
212                sequence = acknowledgement.sequence,
213                snapshot = acknowledgement
214                    .snapshot
215                    .as_ref()
216                    .map(|locator| locator.snapshot.generation),
217                "Staged a Store acknowledgement"
218            );
219        }
220        Box::pin(self.drain_acknowledgements())
221            .await
222            .map_err(|error| SyncCycleFailure::operation("publish Store acknowledgement", error))?;
223        Ok(())
224    }
225
226    /// Stand on the snapshot this device has already acknowledged.
227    ///
228    /// Its own cycle stage, not a step of publishing an acknowledgement,
229    /// because the licence is the statement the device has already made and not
230    /// the act of making another. A device with nothing new to say never stages
231    /// one; a device whose store has moved past every published snapshot can no
232    /// longer name one to stage; and both of those describe a device with a full
233    /// retained history to retire. Reading the licence out of what a device is
234    /// about to say finds nothing in exactly those cases.
235    ///
236    /// Idempotent: adopting a cut the baseline already holds retires nothing,
237    /// and the ordinary answer once a device has caught up is
238    /// [`ReplayBaselineDecline::BaselineAtCoverage`], reached without reading
239    /// anything from the provider.
240    pub(crate) async fn stand_on_acknowledged_snapshot(
241        &mut self,
242        routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
243    ) -> Result<ReplayBaselineAdvance, StoreAckError> {
244        let registration = self.writer.local_registration_ref().clone();
245        let members = self.writer.membership().clone();
246        let resolved = self
247            .writer
248            .resolve_acknowledged_snapshot(&registration, &members)
249            .await?;
250        let selected = match resolved {
251            Ok(selected) => selected,
252            Err(decline) => return Ok(ReplayBaselineAdvance::Declined(decline)),
253        };
254        let generation = selected.snapshot.reference.generation;
255        let advanced = match self.advance_over(Some(selected), routing_encryption).await {
256            Ok(advanced) => advanced,
257            Err(StoreAckError::Database(coven_database::DbError::ReplayRetirementCutNotPrefix)) => {
258                return Ok(ReplayBaselineAdvance::Declined(
259                    ReplayBaselineDecline::NonPrefixCut { generation },
260                ));
261            }
262            Err(error) => return Err(error),
263        };
264        match advanced {
265            Some(advanced) => Ok(ReplayBaselineAdvance::Advanced(advanced)),
266            // The cut this snapshot covers does not move the baseline forward,
267            // which the coverage check above did not catch: the baseline is at
268            // or past it by a route the coverage comparison did not see.
269            None => Ok(ReplayBaselineAdvance::Declined(
270                ReplayBaselineDecline::BaselineAtCoverage { generation },
271            )),
272        }
273    }
274
275    async fn snapshot_it_will_name(
276        &mut self,
277        frontier: &CommitFrontier,
278        device_state: &coven_protocol::store_commit::StoreDeviceStateRef,
279        settled: Option<&crate::sync::store::SettledCycle>,
280    ) -> Result<Option<StoreSnapshotLocator>, StoreAckError> {
281        // Which snapshot a device could acknowledge next is settled by the same
282        // local facts reclaim's answer is, and reaching it means reading every
283        // activated device's snapshot stream. A device with nothing new to say
284        // asks that question every cycle and gets the same answer, so the
285        // answer is remembered against the facts it was reached from and the
286        // read is skipped until one of them moves.
287        let inputs = match settled {
288            Some(settled) => {
289                let inputs =
290                    crate::sync::store::CycleInputs::read(&self.database, self.writer.membership())
291                        .await?;
292                if let Some(remembered) = settled.acknowledgeable_snapshot(&inputs) {
293                    return Ok(remembered);
294                }
295                Some((settled, inputs))
296            }
297            None => None,
298        };
299        let selected = self
300            .writer
301            .select_acknowledgement_snapshot(frontier, device_state)
302            .await?;
303        if let Some((settled, inputs)) = &inputs {
304            settled.record_acknowledgeable_snapshot(
305                inputs.clone(),
306                selected.as_ref().map(|selected| StoreSnapshotLocator {
307                    author_registration: selected.snapshot.meta.author_registration.clone(),
308                    snapshot: selected.snapshot.reference.clone(),
309                }),
310            );
311        }
312        let Some(selected) = selected else {
313            return Ok(None);
314        };
315        let locator = StoreSnapshotLocator {
316            author_registration: selected.snapshot.meta.author_registration.clone(),
317            snapshot: selected.snapshot.reference.clone(),
318        };
319        Ok(Some(locator))
320    }
321
322    async fn advance_over(
323        &mut self,
324        snapshot: Option<SelectedReplayBaselineRetirement>,
325        routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
326    ) -> Result<Option<coven_database::AdvancedReplayBaseline>, StoreAckError> {
327        let Some(snapshot) = snapshot else {
328            return Ok(None);
329        };
330        Ok(self
331            .database
332            .advance_snapshot_replay_baseline(
333                self.writer.store_root().clone(),
334                snapshot.verified,
335                routing_encryption.cloned(),
336            )
337            .await?)
338    }
339
340    /// Stage this device's acknowledgement of `frontier`, unless the one it
341    /// already published still says the same thing.
342    ///
343    /// Publishing an acknowledgement appends a commit, so an acknowledgement that
344    /// asserts nothing new still lands in every device's history, every retained
345    /// materialization, and every snapshot taken afterwards. Without a guard the
346    /// device acknowledges its own acknowledgement and a Store where nothing is
347    /// happening grows one commit per device per sync cycle, without end.
348    ///
349    /// [`StoreAckAssertion`] is what an acknowledgement claims; the rest of it —
350    /// the sequence, the wall clock, the links to its neighbours — differs by
351    /// construction and says nothing. The one subtlety is the frontier: an
352    /// acknowledgement cannot cover the commit that carries it, so the standing
353    /// state records that commit and the comparison treats it as covered.
354    /// Anything else in the frontier having moved is new material to acknowledge.
355    ///
356    /// Returns the acknowledgement it staged, or `None` when the standing one
357    /// still holds, alongside what advancing the baseline retired.
358    #[cfg(any(test, feature = "test-utils"))]
359    pub(crate) async fn stage_acknowledgement(
360        &mut self,
361        frontier: CommitFrontier,
362        sync_time: String,
363    ) -> Result<StagedStoreAcknowledgement, StoreAckError> {
364        self.stage_acknowledgement_against(frontier, sync_time, None)
365            .await
366    }
367
368    async fn stage_acknowledgement_against(
369        &mut self,
370        frontier: CommitFrontier,
371        sync_time: String,
372        settled: Option<&crate::sync::store::SettledCycle>,
373    ) -> Result<StagedStoreAcknowledgement, StoreAckError> {
374        let history_cut =
375            coven_protocol::store_commit::StoreHistoryCut::from_commits(frontier.commits().clone());
376        let (device_state, _) = self
377            .database
378            .store_device_state_for_history_cut(&history_cut)
379            .await?;
380        let previous = self.database.latest_local_store_ack().await?;
381        let snapshot = self
382            .snapshot_it_will_name(&frontier, &device_state, settled)
383            .await?;
384        let acknowledgement = self
385            .say_acknowledgement(history_cut, device_state, previous, snapshot, sync_time)
386            .await?;
387        Ok(StagedStoreAcknowledgement { acknowledgement })
388    }
389
390    /// Stage the acknowledgement itself after selecting the snapshot it names.
391    /// Publishing the durable promise does not retire local replay inputs;
392    /// retirement later requires every current writer to have crossed the cut.
393    async fn say_acknowledgement(
394        &mut self,
395        history_cut: coven_protocol::store_commit::StoreHistoryCut,
396        device_state: coven_protocol::store_commit::StoreDeviceStateRef,
397        previous: Option<coven_database::PublishedStoreAck>,
398        snapshot: Option<StoreSnapshotLocator>,
399        sync_time: String,
400    ) -> Result<Option<StoreAck>, StoreAckError> {
401        let device_id = self.writer.local_device_id().to_string();
402        let root = self.writer.store_root().clone();
403        let exclusions = coven_protocol::store_commit::StoreAckExclusionState {
404            proposal_freezes: self.database.store_device_exclusion_freezes().await?,
405        };
406        if self.database.oldest_outbound_store_ack().await?.is_some() {
407            return Err(StoreAckError::InvalidOutbound(
408                "a prior acknowledgement remains queued".to_string(),
409            ));
410        }
411        let assertion = self.local_writer.device_acknowledgement_assertion(
412            history_cut,
413            device_state.clone(),
414            snapshot,
415            exclusions,
416        );
417        let membership_state =
418            coven_protocol::circle_control::StoreMembershipStateRef::from_membership(
419                self.writer.membership(),
420                device_state.recovery().to_vec(),
421            )?;
422        // A queued Circle acknowledgement travels to the cloud inside the Store
423        // acknowledgement's commit, so one waiting is reason enough to publish
424        // even when this device has nothing of its own left to say.
425        let carries_circle_acknowledgements = self.database.outbound_circle_acks_pending().await?;
426        let standing_still_holds = previous
427            .as_ref()
428            .and_then(|previous| previous.standing.as_ref())
429            .is_some_and(|standing| {
430                standing.still_holds(&assertion)
431                    && standing
432                        .activating_commit
433                        .as_ref()
434                        .and_then(|commit| self.writer.accepted_commit_membership_state(commit))
435                        == Some(&membership_state)
436            });
437        if !carries_circle_acknowledgements && standing_still_holds {
438            debug!("skip Store acknowledgement: the standing one still holds");
439            return Ok(None);
440        }
441        let (sequence, predecessor, current_slot) = match previous {
442            Some(previous) => (
443                previous.reference.sequence.checked_add(1).ok_or_else(|| {
444                    StoreAckError::InvalidOutbound(
445                        "Store acknowledgement sequence overflow".to_string(),
446                    )
447                })?,
448                Some(previous.reference.object),
449                previous.successor_slot,
450            ),
451            None => (1, None, self.local_writer.first_acknowledgement_slot()),
452        };
453        let context = ProtocolObjectContext::signed_plaintext(
454            root.store_root_hash,
455            ProtocolObjectDomain::StoreAck,
456        );
457        let semantic_prefix = ack_slot_prefix(&device_id, sequence);
458        let next_slot = self
459            .storage
460            .allocate_protocol_slot(
461                &context,
462                &ack_slot_prefix(
463                    &device_id,
464                    sequence.checked_add(1).ok_or_else(|| {
465                        StoreAckError::InvalidOutbound(
466                            "Store acknowledgement sequence overflow".to_string(),
467                        )
468                    })?,
469                ),
470                ".json",
471            )
472            .await
473            .map_err(StoreObjectError::from)?;
474        let activation = self
475            .local_writer
476            .acknowledgement_activation_id()
477            .map_err(StoreAckError::from)?;
478        let acknowledgement = self
479            .local_writer
480            .sign_device_acknowledgement(
481                root.store_root_hash,
482                sequence,
483                assertion,
484                sync_time,
485                SuccessorLink {
486                    activation,
487                    predecessor,
488                    next_slot,
489                },
490            )
491            .map_err(StoreAckError::from)?;
492        let prepared = self
493            .storage
494            .prepare_protocol_object(
495                &context,
496                current_slot,
497                &semantic_prefix,
498                acknowledgement.to_bytes(),
499            )
500            .map_err(StoreObjectError::from)?;
501        self.database
502            .stage_store_ack(acknowledgement.clone(), prepared)
503            .await?;
504        Ok(Some(acknowledgement))
505    }
506
507    pub(crate) async fn drain_acknowledgements(&mut self) -> Result<u64, StoreAckError> {
508        let device_id = self.writer.local_device_id().to_string();
509        let mut published = 0_u64;
510        while let Some(outbound) = self.database.oldest_outbound_store_ack().await? {
511            if let Some(activated) = self
512                .database
513                .activated_store_ack(&outbound.reference.registration)
514                .await?
515            {
516                if activated.reference == outbound.reference {
517                    self.database
518                        .complete_outbound_store_ack(
519                            outbound.reference,
520                            activated.activating_commit,
521                        )
522                        .await?;
523                    published = published
524                        .checked_add(1)
525                        .ok_or(StoreAckError::PublishCountExhausted)?;
526                    continue;
527                }
528                if activated.reference.sequence >= outbound.reference.sequence {
529                    return Err(StoreAckError::InvalidOutbound(
530                        "queued Store acknowledgement differs from the activated exact ref"
531                            .to_string(),
532                    ));
533                }
534            }
535            let candidate = match outbound.activation.clone() {
536                coven_database::OutboundStoreAckActivation::AwaitingCandidate => {
537                    let plan = self.writer.prepare_plan().await?;
538                    plan.validate_acknowledgement(&outbound.ack.value)?;
539                    let candidate = Box::pin(self.writer.prepare_candidate(
540                        plan,
541                        crate::sync::store::commit_publication::operation::commit_plan::StoreOperationBatch::Acknowledgement {
542                            reference: outbound.reference.clone(),
543                            value: outbound.ack.value.clone(),
544                            circle_acknowledgements: outbound.circle_acknowledgements.clone(),
545                        },
546                    ))
547                    .await?;
548                    self.database
549                        .prepare_acknowledgement_activation(outbound.reference.clone(), candidate)
550                        .await?;
551                    continue;
552                }
553                coven_database::OutboundStoreAckActivation::Prepared(candidate) => candidate,
554                coven_database::OutboundStoreAckActivation::Nonactivating(_) => {
555                    self.writer
556                        .finish_nonactivating_acknowledgement(outbound.reference)
557                        .await?;
558                    published = published
559                        .checked_add(1)
560                        .ok_or(StoreAckError::PublishCountExhausted)?;
561                    continue;
562                }
563            };
564            let context = ProtocolObjectContext::signed_plaintext(
565                outbound.ack.value.store_root_hash,
566                ProtocolObjectDomain::StoreAck,
567            );
568            let semantic_prefix = ack_slot_prefix(&device_id, outbound.reference.sequence);
569            if let Err(error) = self
570                .storage
571                .create_verified_protocol_object(
572                    &context,
573                    &outbound.ack.prepared,
574                    &semantic_prefix,
575                    &outbound.ack.bytes,
576                )
577                .await
578            {
579                if !matches!(
580                    error,
581                    coven_protocol::objects::StorageError::SlotCollision(_)
582                ) {
583                    return Err(StoreObjectError::from(error).into());
584                }
585                let (winner_bytes, winner_prepared) = self
586                    .storage
587                    .read_prepared_protocol_slot(
588                        &context,
589                        outbound.reference.object.slot(),
590                        &semantic_prefix,
591                    )
592                    .await
593                    .map_err(StoreObjectError::from)?;
594                self.database
595                    .adopt_outbound_store_ack_slot_winner(
596                        outbound.reference.clone(),
597                        winner_bytes,
598                        winner_prepared,
599                    )
600                    .await?;
601                continue;
602            }
603            let acknowledgement_remote = candidate
604                .acknowledgement_remote_objects(&outbound.ack)?
605                .into_iter()
606                .find(|remote| remote.object() == &outbound.reference.object)
607                .ok_or_else(|| {
608                    StoreAckError::InvalidOutbound(
609                        "prepared activation does not own its acknowledgement object".to_string(),
610                    )
611                })?;
612            self.database
613                .mark_remote_object_uploaded(acknowledgement_remote.into_record())
614                .await?;
615            self.writer
616                .circles()
617                .publish_acknowledgement_objects(&outbound, &candidate)
618                .await?;
619            let _authorship = self.database.author_own_stream().await;
620            let publication = Box::pin(self.writer.publish_prepared(
621                Box::new(candidate),
622                None,
623                None,
624            ))
625            .await?;
626            match publication
627            {
628                crate::sync::store::commit_publication::operation::commit_plan::StoreOperationPublicationOutcome::Activated(activating_commit) => {
629                    self.database
630                        .complete_outbound_store_ack(outbound.reference, activating_commit)
631                        .await?;
632                }
633                crate::sync::store::commit_publication::operation::commit_plan::StoreOperationPublicationOutcome::Nonactivated(_) => {}
634                crate::sync::store::commit_publication::operation::commit_plan::StoreOperationPublicationOutcome::Reprepared => {
635                    continue;
636                }
637                crate::sync::store::commit_publication::operation::commit_plan::StoreOperationPublicationOutcome::RepreparedCandidate(_)
638                | crate::sync::store::commit_publication::operation::commit_plan::StoreOperationPublicationOutcome::NonactivatedCandidate { .. } => {
639                    return Err(StoreAckError::InvalidOutbound(
640                        "acknowledgement publication returned non-acknowledgement conflict state"
641                            .to_string(),
642                    ));
643                }
644            }
645            published = published
646                .checked_add(1)
647                .ok_or(StoreAckError::PublishCountExhausted)?;
648        }
649        Ok(published)
650    }
651}
652
653#[cfg(test)]
654mod tests;