Skip to main content

heddle_thread_api/replication/
store.rs

1//! Async durable-store boundary shared by device and hosted replication.
2use std::{collections::BTreeSet, future::Future, sync::Arc};
3
4use crypto::{
5    thread_authority_admission::SignedAuthorityAdmission, thread_operation::SignedOperation,
6};
7use heddle_object_model::object::{
8    ContentHash,
9    thread_replication::{Admission, ThreadFacet},
10};
11
12/// Original immutable bytes plus optional independently signed first-authority
13/// admission. A receipt never replaces the original signature or current courier
14/// authorization, and remains attached across later peer relays.
15#[derive(Clone, Debug, PartialEq)]
16pub struct ReceivedOperation {
17    pub original: SignedOperation,
18    pub authority_admission: Option<SignedAuthorityAdmission>,
19    /// Complete public evidence stays attached through staging, admission and
20    /// relay. Its presence is not authority; the store verifies it at mutation.
21    pub native_authority: Option<Arc<crate::contract::NativePublicProofBundleV1>>,
22    pub import_authority: Option<Arc<crate::contract::ImportPublicProofBundleV1>>,
23}
24impl From<SignedOperation> for ReceivedOperation {
25    fn from(original: SignedOperation) -> Self {
26        Self {
27            original,
28            authority_admission: None,
29            import_authority: None,
30            native_authority: None,
31        }
32    }
33}
34
35/// Implementations commit before returning receipts. A store is bound to one
36/// authorized Thread; operations and peer metadata may never escape that scope.
37/// Notification delivery is separate from this durable contract.
38pub trait ReplicaStore: Clone + Send + Sync + 'static {
39    type Error: std::error::Error + Send + Sync + 'static;
40    fn thread_id(&self) -> ContentHash;
41    fn generation(&self) -> impl Future<Output = Result<i64, Self::Error>> + Send;
42    fn sharing(
43        &self,
44        destination: [u8; 32],
45    ) -> impl Future<Output = Result<BTreeSet<ThreadFacet>, Self::Error>> + Send;
46    fn frontier_page(
47        &self,
48        facet: ThreadFacet,
49        after: Option<ContentHash>,
50        limit: usize,
51    ) -> impl Future<Output = Result<Vec<ContentHash>, Self::Error>> + Send;
52    /// Recheck retained hosted authority under the receiver's mutation
53    /// serialization before exporting an accepted original. Cached acceptance
54    /// and structural sidecars grant no authority; keep the complete public
55    /// evidence attached for the next receiver.
56    fn operation(
57        &self,
58        id: ContentHash,
59    ) -> impl Future<Output = Result<Option<(ReceivedOperation, Admission)>, Self::Error>> + Send;
60    fn receive(
61        &self,
62        operation: ReceivedOperation,
63    ) -> impl Future<Output = Result<Admission, Self::Error>> + Send;
64    fn remember_peer_heads(
65        &self,
66        peer: [u8; 32],
67        heads: Vec<(ThreadFacet, ContentHash)>,
68    ) -> impl Future<Output = Result<(), Self::Error>> + Send;
69    fn record_peer_receipt(
70        &self,
71        peer: [u8; 32],
72        id: ContentHash,
73        admission: Admission,
74    ) -> impl Future<Output = Result<(), Self::Error>> + Send;
75    fn settled_peer_heads(
76        &self,
77        peer: [u8; 32],
78        facets: BTreeSet<ThreadFacet>,
79        limit: usize,
80    ) -> impl Future<Output = Result<Vec<(ContentHash, Admission)>, Self::Error>> + Send;
81    fn needed_from_peer(
82        &self,
83        peer: [u8; 32],
84        facets: BTreeSet<ThreadFacet>,
85        limit: usize,
86    ) -> impl Future<Output = Result<Vec<ContentHash>, Self::Error>> + Send;
87}