1mod 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#[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 NoAcknowledgedSnapshot,
80 SnapshotAuthorInactive { generation: u64 },
83 SnapshotUnavailable { generation: u64 },
85 SnapshotRejected { generation: u64 },
87 MissingWriterAcknowledgement {
90 generation: u64,
91 member: String,
92 device_id: String,
93 },
94 MembershipNotAccepted { generation: u64 },
97 PendingOwnerRecovery {
100 generation: u64,
101 member: String,
102 device_id: String,
103 },
104 NonPrefixCut { generation: u64 },
107 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 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 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 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(®istration, &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 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 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 #[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 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 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;