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::SelectedStoreSnapshot;
12use coven_database::StoreDatabase;
13use coven_protocol::objects::StoreObjectError;
14use coven_protocol::objects::{ProtocolObjectContext, ProtocolObjectDomain};
15use coven_protocol::store_commit::{ack_slot_prefix, CommitFrontier, StoreAck, SuccessorLink};
16use coven_storage::CloudSyncObjectStorage;
17use std::sync::Arc;
18use tracing::debug;
19
20#[derive(Debug, thiserror::Error)]
21pub enum StoreAckError {
22 #[error("Store acknowledgement membership: {0}")]
23 Membership(#[from] coven_protocol::membership::MembershipError),
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)]
69pub enum ReplayBaselineAdvance {
70 Advanced(coven_database::AdvancedReplayBaseline),
71 Declined(ReplayBaselineDecline),
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub enum ReplayBaselineDecline {
76 NoAcceptedSnapshot,
78 NonPrefixCut {
81 snapshot: coven_protocol::store_commit::StoreSnapshotRef,
82 },
83 BaselineAtCoverage {
86 snapshot: coven_protocol::store_commit::StoreSnapshotRef,
87 },
88}
89
90impl ReplayBaselineDecline {
91 pub fn as_str(&self) -> &'static str {
92 match self {
93 Self::NoAcceptedSnapshot => "accepted Store history contains no snapshot",
94 Self::NonPrefixCut { .. } => {
95 "the accepted snapshot is not a prefix of accepted replay order"
96 }
97 Self::BaselineAtCoverage { .. } => "the baseline already covers it",
98 }
99 }
100
101 pub fn snapshot(&self) -> Option<&coven_protocol::store_commit::StoreSnapshotRef> {
102 match self {
103 Self::NoAcceptedSnapshot => None,
104 Self::NonPrefixCut { snapshot } | Self::BaselineAtCoverage { snapshot } => {
105 Some(snapshot)
106 }
107 }
108 }
109}
110
111pub(crate) struct AuthorizedAcknowledgements<'operation, 'storage> {
112 writer: &'operation mut AuthorizedWriterOperation<'storage>,
113 database: StoreDatabase,
114 storage: Arc<dyn CloudSyncObjectStorage>,
115 local_writer: Arc<LocalStoreWriter>,
116}
117
118impl<'operation, 'storage> AuthorizedAcknowledgements<'operation, 'storage> {
119 pub(crate) fn new(
120 writer: &'operation mut AuthorizedWriterOperation<'storage>,
121 database: StoreDatabase,
122 storage: Arc<dyn CloudSyncObjectStorage>,
123 local_writer: Arc<LocalStoreWriter>,
124 ) -> Self {
125 Self {
126 writer,
127 database,
128 storage,
129 local_writer,
130 }
131 }
132
133 pub(crate) async fn stage_and_publish(
137 &mut self,
138 sync_time: &str,
139 ) -> Result<(), SyncCycleFailure> {
140 self.drain_acknowledgements().await.map_err(|error| {
141 SyncCycleFailure::operation("publish queued Store acknowledgement", error)
142 })?;
143 let frontier =
144 CommitFrontier::from_refs(self.database.materialized_frontier().await.map_err(
145 |error| SyncCycleFailure::operation("read Store acknowledgement frontier", error),
146 )?)
147 .map_err(|error| {
148 SyncCycleFailure::operation("shape Store acknowledgement frontier", error)
149 })?;
150 Box::pin(
154 self.writer
155 .circles()
156 .stage_acknowledgements(&frontier, sync_time),
157 )
158 .await
159 .map_err(|error| SyncCycleFailure::operation("stage Circle acknowledgements", error))?;
160 let StagedStoreAcknowledgement { acknowledgement } =
161 Box::pin(self.stage_acknowledgement(frontier.clone(), sync_time.to_owned()))
162 .await
163 .map_err(|error| {
164 SyncCycleFailure::operation("stage Store acknowledgement", error)
165 })?;
166 if let Some(acknowledgement) = &acknowledgement {
167 debug!(
168 sequence = acknowledgement.sequence,
169 "Staged a Store acknowledgement"
170 );
171 }
172 self.drain_acknowledgements()
173 .await
174 .map_err(|error| SyncCycleFailure::operation("publish Store acknowledgement", error))?;
175 Ok(())
176 }
177
178 pub(crate) async fn stand_on_accepted_snapshot(
185 &mut self,
186 routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
187 ) -> Result<ReplayBaselineAdvance, StoreAckError> {
188 let resolved = self.writer.resolve_accepted_snapshot().await?;
189 let selected = match resolved {
190 Ok(selected) => selected,
191 Err(decline) => return Ok(ReplayBaselineAdvance::Declined(decline)),
192 };
193 let snapshot = selected.snapshot.reference.clone();
194 let advanced = match self.advance_over(selected, routing_encryption).await {
195 Ok(advanced) => advanced,
196 Err(StoreAckError::Database(coven_database::DbError::ReplayRetirementCutNotPrefix)) => {
197 return Ok(ReplayBaselineAdvance::Declined(
198 ReplayBaselineDecline::NonPrefixCut { snapshot },
199 ));
200 }
201 Err(error) => return Err(error),
202 };
203 match advanced {
204 Some(advanced) => Ok(ReplayBaselineAdvance::Advanced(advanced)),
205 None => Ok(ReplayBaselineAdvance::Declined(
209 ReplayBaselineDecline::BaselineAtCoverage { snapshot },
210 )),
211 }
212 }
213
214 async fn advance_over(
215 &mut self,
216 snapshot: SelectedStoreSnapshot,
217 routing_encryption: Option<&coven_keys::encryption::EncryptionService>,
218 ) -> Result<Option<coven_database::AdvancedReplayBaseline>, StoreAckError> {
219 Ok(self
220 .database
221 .advance_snapshot_replay_baseline(
222 self.writer.store_root().clone(),
223 snapshot.verified,
224 routing_encryption.cloned(),
225 self.local_writer.local_membership(self.writer.membership()),
226 )
227 .await?)
228 }
229
230 pub(crate) async fn stage_acknowledgement(
250 &mut self,
251 frontier: CommitFrontier,
252 sync_time: String,
253 ) -> Result<StagedStoreAcknowledgement, StoreAckError> {
254 let history_cut =
255 coven_protocol::store_commit::StoreHistoryCut::from_commits(frontier.commits().clone());
256 let (device_state, _) = self
257 .database
258 .store_device_state_for_history_cut(&history_cut)
259 .await?;
260 let previous = self.database.latest_local_store_ack().await?;
261 let acknowledgement = self
262 .say_acknowledgement(history_cut, device_state, previous, sync_time)
263 .await?;
264 Ok(StagedStoreAcknowledgement { acknowledgement })
265 }
266
267 async fn say_acknowledgement(
270 &mut self,
271 history_cut: coven_protocol::store_commit::StoreHistoryCut,
272 device_state: coven_protocol::store_commit::StoreDeviceStateRef,
273 previous: Option<coven_database::PublishedStoreAck>,
274 sync_time: String,
275 ) -> Result<Option<StoreAck>, StoreAckError> {
276 let device_id = self.writer.local_device_id().to_string();
277 let root = self.writer.store_root().clone();
278 if self.database.oldest_outbound_store_ack().await?.is_some() {
279 return Err(StoreAckError::InvalidOutbound(
280 "a prior acknowledgement remains queued".to_string(),
281 ));
282 }
283 let assertion = self
284 .local_writer
285 .device_acknowledgement_assertion(history_cut, device_state);
286 let carries_circle_acknowledgements = self.database.outbound_circle_acks_pending().await?;
290 let standing_still_holds = match previous
291 .as_ref()
292 .and_then(|previous| previous.standing.as_ref())
293 {
294 Some(standing) if standing.still_holds(&assertion) => true,
295 Some(standing) if standing.assertion.same_state_as(&assertion) => self
296 .writer
297 .history_has_only_acknowledgements(
298 &standing.assertion.store_cut,
299 &assertion.store_cut,
300 )
301 .await
302 .map_err(StoreError::from)?,
303 Some(_) | None => false,
304 };
305 if !carries_circle_acknowledgements && standing_still_holds {
306 debug!("skip Store acknowledgement: the standing one still holds");
307 return Ok(None);
308 }
309 let (sequence, predecessor, current_slot) = match previous {
310 Some(previous) => (
311 previous.reference.sequence.checked_add(1).ok_or_else(|| {
312 StoreAckError::InvalidOutbound(
313 "Store acknowledgement sequence overflow".to_string(),
314 )
315 })?,
316 Some(previous.reference.object),
317 previous.successor_slot,
318 ),
319 None => (1, None, self.local_writer.first_acknowledgement_slot()),
320 };
321 let context = ProtocolObjectContext::signed_plaintext(
322 root.store_root_hash,
323 ProtocolObjectDomain::StoreAck,
324 );
325 let semantic_prefix = ack_slot_prefix(&device_id, sequence);
326 let next_slot = self
327 .storage
328 .allocate_protocol_slot(
329 &context,
330 &ack_slot_prefix(
331 &device_id,
332 sequence.checked_add(1).ok_or_else(|| {
333 StoreAckError::InvalidOutbound(
334 "Store acknowledgement sequence overflow".to_string(),
335 )
336 })?,
337 ),
338 ".json",
339 )
340 .await
341 .map_err(StoreObjectError::from)?;
342 let activation = self
343 .local_writer
344 .acknowledgement_activation_id()
345 .map_err(StoreAckError::from)?;
346 let acknowledgement = self
347 .local_writer
348 .sign_device_acknowledgement(
349 sequence,
350 assertion,
351 sync_time,
352 SuccessorLink {
353 activation,
354 predecessor,
355 next_slot,
356 },
357 )
358 .map_err(StoreAckError::from)?;
359 let prepared = self
360 .storage
361 .prepare_protocol_object(
362 &context,
363 current_slot,
364 &semantic_prefix,
365 acknowledgement.to_bytes(),
366 )
367 .map_err(StoreObjectError::from)?;
368 self.database
369 .stage_store_ack(acknowledgement.clone(), prepared)
370 .await?;
371 Ok(Some(acknowledgement))
372 }
373
374 pub(crate) async fn drain_acknowledgements(&mut self) -> Result<u64, StoreAckError> {
375 let mut authorship = self.database.author_own_stream().await;
376 let device_id = self.writer.local_device_id().to_string();
377 let mut published = 0_u64;
378 while let Some(outbound) = self.database.oldest_outbound_store_ack().await? {
379 if let Some(active) = self.database.active_store_publication().await? {
380 if active.owner()
381 == &coven_database::ActiveStorePublicationOwner::StoreAcknowledgement
382 {
383 crate::sync::store::authorization::retire_store_write_candidates(
384 &self.database,
385 self.storage.as_ref(),
386 active,
387 )
388 .await?;
389 }
390 }
391 if let Some(activated) = self
392 .database
393 .activated_store_ack(&outbound.reference.registration)
394 .await?
395 {
396 if activated.reference == outbound.reference {
397 self.database
398 .complete_outbound_store_ack(
399 outbound.reference,
400 activated.activating_commit,
401 )
402 .await?;
403 published = published
404 .checked_add(1)
405 .ok_or(StoreAckError::PublishCountExhausted)?;
406 continue;
407 }
408 if activated.reference.sequence > outbound.reference.sequence {
409 return Err(StoreAckError::InvalidOutbound(
410 "queued Store acknowledgement differs from the activated exact ref"
411 .to_string(),
412 ));
413 }
414 }
415 let candidate = match outbound.activation.clone() {
416 coven_database::OutboundStoreAckActivation::AwaitingCandidate
417 | coven_database::OutboundStoreAckActivation::Created => {
418 let plan = self.writer.prepare_plan_with_authorship(authorship).await?;
419 let created = matches!(
420 outbound.activation,
421 coven_database::OutboundStoreAckActivation::Created
422 );
423 if created
424 && self
425 .database
426 .activated_store_ack(&outbound.reference.registration)
427 .await?
428 .is_some_and(|activated| activated.reference == outbound.reference)
429 {
430 authorship = plan.into_authorship();
431 continue;
432 }
433 let (reference, acknowledgement) = if created {
434 if outbound.ack.value.store_cut != plan.predecessor_cut()?
435 || &outbound.ack.value.device_state != plan.device_state()
436 {
437 self.prepare_acknowledgement_successor(&outbound, &plan)
438 .await?
439 } else {
440 (outbound.reference.clone(), outbound.ack.clone())
441 }
442 } else {
443 self.prepare_acknowledgement_object(
444 &outbound,
445 &plan,
446 outbound.reference.sequence,
447 outbound.reference.object.slot().clone(),
448 outbound.ack.value.successor.clone(),
449 )?
450 };
451 plan.validate_acknowledgement(&acknowledgement.value)?;
452 let mut candidate = Box::pin(self.writer.prepare_candidate(
453 &plan,
454 crate::sync::store::commit_publication::operation::commit_plan::StoreOperationBatch::Acknowledgement {
455 reference: reference.clone(),
456 value: acknowledgement.value.clone(),
457 circle_acknowledgements: outbound.circle_acknowledgements.clone(),
458 },
459 ))
460 .await?;
461 if created && reference != outbound.reference {
462 candidate
463 .history_evidence
464 .acknowledgement
465 .as_mut()
466 .ok_or_else(|| {
467 StoreAckError::InvalidOutbound(
468 "prepared acknowledgement omits its retained proof".into(),
469 )
470 })?
471 .predecessors =
472 vec![(outbound.reference.clone(), outbound.ack.value.clone())];
473 candidate.validate_closed_shape()?;
474 }
475 let claimed = self
476 .database
477 .prepare_acknowledgement_activation(
478 outbound.reference.clone(),
479 acknowledgement,
480 candidate,
481 )
482 .await?;
483 if !claimed {
484 return Ok(published);
485 }
486 authorship = plan.into_authorship();
487 continue;
488 }
489 coven_database::OutboundStoreAckActivation::Prepared(candidate) => candidate,
490 };
491 let context = ProtocolObjectContext::signed_plaintext(
492 outbound.ack.value.store_root_hash,
493 ProtocolObjectDomain::StoreAck,
494 );
495 let semantic_prefix = ack_slot_prefix(&device_id, outbound.reference.sequence);
496 if let Err(error) = self
497 .storage
498 .create_verified_protocol_object(
499 &context,
500 &outbound.ack.prepared,
501 &semantic_prefix,
502 &outbound.ack.bytes,
503 )
504 .await
505 {
506 if !matches!(
507 error,
508 coven_protocol::objects::StorageError::SlotCollision(_)
509 ) {
510 return Err(StoreObjectError::from(error).into());
511 }
512 let (winner_bytes, winner_prepared) = self
513 .storage
514 .read_prepared_protocol_slot(
515 &context,
516 outbound.reference.object.slot(),
517 &semantic_prefix,
518 )
519 .await
520 .map_err(StoreObjectError::from)?;
521 self.database
522 .adopt_outbound_store_ack_slot_winner(
523 outbound.reference.clone(),
524 winner_bytes,
525 winner_prepared,
526 )
527 .await?;
528 continue;
529 }
530 let acknowledgement_remote = candidate
531 .acknowledgement_remote_objects(&outbound.ack)?
532 .into_iter()
533 .find(|remote| remote.object() == &outbound.reference.object)
534 .ok_or_else(|| {
535 StoreAckError::InvalidOutbound(
536 "prepared activation does not own its acknowledgement object".to_string(),
537 )
538 })?;
539 self.database
540 .mark_remote_object_uploaded(acknowledgement_remote.into_record())
541 .await?;
542 self.writer
543 .circles()
544 .publish_acknowledgement_objects(&outbound, &candidate)
545 .await?;
546 let outcome = Box::pin(self.writer.publish_prepared_attempt(
547 Box::new(candidate.clone()),
548 None,
549 None,
550 ))
551 .await?;
552 let accepted = match outcome {
553 crate::sync::store::commit_publication::operation::StoreOperationPublicationOutcome::Accepted(accepted) => accepted,
554 crate::sync::store::commit_publication::operation::StoreOperationPublicationOutcome::SnapshotRetired(snapshot) => {
555 authorship = self.replace_snapshot_acknowledgement(&outbound, &candidate, snapshot, authorship).await?;
556 continue;
557 }
558 };
559 self.database
560 .complete_outbound_store_ack(outbound.reference, accepted.commit_ref().clone())
561 .await?;
562 published = published
563 .checked_add(1)
564 .ok_or(StoreAckError::PublishCountExhausted)?;
565 }
566 Ok(published)
567 }
568
569 async fn prepare_acknowledgement_successor(
570 &self,
571 outbound: &coven_database::OutboundStoreAck,
572 plan: &crate::sync::store::commit_publication::operation::commit_plan::StoreOperationCommitPlan,
573 ) -> Result<
574 (
575 coven_protocol::store_commit::StoreAckRef,
576 coven_protocol::objects::ExactProtocolObject<StoreAck>,
577 ),
578 StoreAckError,
579 > {
580 let sequence = outbound.reference.sequence.checked_add(1).ok_or_else(|| {
581 StoreAckError::InvalidOutbound("Store acknowledgement sequence overflow".into())
582 })?;
583 let following = sequence.checked_add(1).ok_or_else(|| {
584 StoreAckError::InvalidOutbound("Store acknowledgement sequence overflow".into())
585 })?;
586 let context = ProtocolObjectContext::signed_plaintext(
587 plan.root().store_root_hash,
588 ProtocolObjectDomain::StoreAck,
589 );
590 let next_slot = self
591 .storage
592 .allocate_protocol_slot(
593 &context,
594 &ack_slot_prefix(&self.writer.local_device_id().to_string(), following),
595 ".json",
596 )
597 .await
598 .map_err(StoreObjectError::from)?;
599 self.prepare_acknowledgement_object(
600 outbound,
601 plan,
602 sequence,
603 outbound.ack.value.successor.next_slot.clone(),
604 SuccessorLink {
605 activation: outbound.ack.value.successor.activation,
606 predecessor: Some(outbound.reference.object.clone()),
607 next_slot,
608 },
609 )
610 }
611
612 fn prepare_acknowledgement_object(
613 &self,
614 outbound: &coven_database::OutboundStoreAck,
615 plan: &crate::sync::store::commit_publication::operation::commit_plan::StoreOperationCommitPlan,
616 sequence: u64,
617 slot: coven_protocol::objects::ObjectSlot,
618 successor: SuccessorLink,
619 ) -> Result<
620 (
621 coven_protocol::store_commit::StoreAckRef,
622 coven_protocol::objects::ExactProtocolObject<StoreAck>,
623 ),
624 StoreAckError,
625 > {
626 let value = self.local_writer.sign_device_acknowledgement(
627 sequence,
628 self.local_writer.device_acknowledgement_assertion(
629 plan.predecessor_cut()?,
630 plan.device_state().clone(),
631 ),
632 outbound.ack.value.last_sync.clone(),
633 successor,
634 )?;
635 let bytes = value.to_bytes();
636 let prepared = self
637 .storage
638 .prepare_protocol_object(
639 &ProtocolObjectContext::signed_plaintext(
640 plan.root().store_root_hash,
641 ProtocolObjectDomain::StoreAck,
642 ),
643 slot,
644 &ack_slot_prefix(&self.writer.local_device_id().to_string(), sequence),
645 bytes.clone(),
646 )
647 .map_err(StoreObjectError::from)?;
648 let reference = coven_protocol::store_commit::StoreAckRef {
649 registration: value.registration.clone(),
650 sequence,
651 ack_hash: value.ack_hash(),
652 object: prepared.reference().clone(),
653 };
654 Ok((
655 reference,
656 coven_protocol::objects::ExactProtocolObject {
657 value,
658 bytes,
659 prepared,
660 },
661 ))
662 }
663
664 async fn replace_snapshot_acknowledgement(
665 &mut self,
666 outbound: &coven_database::OutboundStoreAck,
667 previous: &coven_protocol::prepared_commit::PreparedStoreOperationCommit,
668 snapshot: coven_protocol::store_commit::AcceptedStoreSnapshotRef,
669 authorship: coven_database::OwnStreamAuthorship,
670 ) -> Result<coven_database::OwnStreamAuthorship, StoreAckError> {
671 let plan = self.writer.prepare_plan_with_authorship(authorship).await?;
672 let (reference, acknowledgement) = self
673 .prepare_acknowledgement_successor(outbound, &plan)
674 .await?;
675 let coven_protocol::objects::ExactProtocolObject {
676 value,
677 bytes,
678 prepared,
679 } = acknowledgement;
680 plan.validate_acknowledgement(&value)?;
681 let mut candidate = self.writer.prepare_replacement_candidate(
682 &plan,
683 crate::sync::store::commit_publication::operation::commit_plan::StoreOperationBatch::Acknowledgement {
684 reference, value: value.clone(),
685 circle_acknowledgements: outbound.circle_acknowledgements.clone(),
686 },
687 previous,
688 ).await?;
689 let previous_proof = previous
690 .history_evidence
691 .acknowledgement
692 .as_ref()
693 .ok_or_else(|| {
694 StoreAckError::InvalidOutbound(
695 "queued acknowledgement omits its retained proof".into(),
696 )
697 })?;
698 let retained = candidate
699 .history_evidence
700 .acknowledgement
701 .as_mut()
702 .ok_or_else(|| {
703 StoreAckError::InvalidOutbound(
704 "replacement acknowledgement omits its retained proof".into(),
705 )
706 })?;
707 retained.predecessors = previous_proof.proof_objects().cloned().collect();
708 candidate.validate_closed_shape()?;
709 self.database
710 .replace_acknowledgement_activation(
711 outbound.reference.clone(),
712 snapshot,
713 coven_protocol::objects::ExactProtocolObject {
714 value,
715 bytes,
716 prepared,
717 },
718 candidate,
719 )
720 .await?;
721 Ok(plan.into_authorship())
722 }
723}
724
725#[cfg(test)]
726mod tests;
727
728#[cfg(test)]
729mod idle_tests;