Skip to main content

heddle_thread_api/replication/
native.rs

1//! Local SQLite adapter. Blocking work stays off the stream runtime.
2use std::{collections::BTreeSet, path::PathBuf, sync::Arc};
3
4use heddle_object_model::object::{
5    ContentHash,
6    thread_replication::{Admission, ThreadFacet},
7};
8use objects::store::ObjectStore;
9use repo::thread_replication::ThreadReplica;
10
11use super::store::{ReceivedOperation, ReplicaStore};
12
13mod hosted;
14pub use hosted::HostedReplica;
15
16#[derive(Debug, thiserror::Error)]
17pub enum Error {
18    #[error(transparent)]
19    Store(#[from] repo::thread_replication::Error),
20    #[error("local replica worker: {0}")]
21    Worker(#[from] tokio::task::JoinError),
22    #[error("HYBRID original requires independently selected hosted trust")]
23    HostedTrustRequired,
24}
25pub struct LocalReplica<S> {
26    replica: ThreadReplica,
27    objects: Arc<S>,
28    authority_home: Option<PathBuf>,
29}
30impl<S> Clone for LocalReplica<S> {
31    fn clone(&self) -> Self {
32        Self {
33            replica: self.replica.clone(),
34            objects: self.objects.clone(),
35            authority_home: self.authority_home.clone(),
36        }
37    }
38}
39impl<S: ObjectStore + Send + Sync + 'static> LocalReplica<S> {
40    pub fn new(replica: ThreadReplica, objects: Arc<S>) -> Self {
41        Self {
42            replica,
43            objects,
44            authority_home: None,
45        }
46    }
47    /// Use independently enrolled local account authority for original metadata
48    /// authors. Without this binding, new metadata admission fails closed.
49    pub fn with_device_authority(mut self, home: PathBuf) -> Self {
50        self.authority_home = Some(home);
51        self
52    }
53    async fn execute<T: Send + 'static>(
54        &self,
55        operation: impl FnOnce(&ThreadReplica, &S) -> repo::thread_replication::Result<T>
56        + Send
57        + 'static,
58    ) -> Result<T, Error> {
59        let local = self.clone();
60        Ok(
61            tokio::task::spawn_blocking(move || operation(&local.replica, &local.objects))
62                .await??,
63        )
64    }
65}
66impl<S: ObjectStore + Send + Sync + 'static> ReplicaStore for LocalReplica<S> {
67    type Error = Error;
68    fn thread_id(&self) -> ContentHash {
69        self.replica.thread_id()
70    }
71    async fn generation(&self) -> Result<i64, Error> {
72        self.execute(|replica, _| replica.generation()).await
73    }
74    async fn sharing(&self, destination: [u8; 32]) -> Result<BTreeSet<ThreadFacet>, Error> {
75        self.execute(move |replica, _| replica.sharing(&destination).map(|(facets, _)| facets))
76            .await
77    }
78    async fn frontier_page(
79        &self,
80        facet: ThreadFacet,
81        after: Option<ContentHash>,
82        limit: usize,
83    ) -> Result<Vec<ContentHash>, Error> {
84        self.execute(move |replica, _| replica.frontier_page(facet, after, limit))
85            .await
86    }
87    async fn operation(
88        &self,
89        id: ContentHash,
90    ) -> Result<Option<(ReceivedOperation, Admission)>, Error> {
91        let stored = self
92            .execute(move |replica, _| replica.operation_with_authority_admission(&id))
93            .await?;
94        let Some(stored) = stored else {
95            return Ok(None);
96        };
97        if stored.authority_admission.is_some()
98            || self
99                .execute(|replica, _| replica.hybrid_import_bundle())
100                .await?
101                .is_some()
102        {
103            return Err(Error::HostedTrustRequired);
104        }
105        Ok(Some((
106            ReceivedOperation {
107                native_authority: None,
108                original: stored.original,
109                authority_admission: None,
110                import_authority: None,
111            },
112            stored.status,
113        )))
114    }
115    async fn receive(&self, received: ReceivedOperation) -> Result<Admission, Error> {
116        if received.import_authority.is_some() || received.authority_admission.is_some() {
117            return Err(Error::HostedTrustRequired);
118        }
119        let authority_home = self.authority_home.clone();
120        self.execute(move |replica, objects| {
121            let operation = received.original;
122            replica.receive(&operation, objects, |native| {
123                use heddle_object_model::object::thread_replication::ThreadOperationBody;
124                if !matches!(native.body, ThreadOperationBody::Metadata(_)) && native.source_author()?.is_none() {
125                    return Ok(());
126                }
127                // Durable original-author receipt remains valid while causal
128                // parents arrive later. Neither the envelope nor claimed time
129                // can synthesize the atomically retained admission marker.
130                if replica.original_authority_admitted(&operation)? {
131                    return Ok(());
132                }
133                let genesis = replica.genesis()?;
134                if matches!(native.source_author()?, Some(heddle_object_model::object::thread_replication::SourceAuthor::LocalKey)) {
135                    return replica.verify_local_source_owner(native);
136                }
137                let home = authority_home.as_ref().ok_or_else(|| {
138                    repo::thread_replication::Error::Invalid(
139                        "original operation requires independently enrolled account authority"
140                            .into(),
141                    )
142                })?;
143                let now = chrono::Utc::now().timestamp();
144                let authority = repo::device_authority::load(home, now).map_err(authority_error)?;
145                if let ThreadOperationBody::Metadata(bytes) = &native.body {
146                    heddle_object_model::object::thread_replication::metadata::ThreadControl::decode(bytes)?.validate_parents(&genesis, &[])?;
147                }
148                let spool = genesis.spool.parse().map_err(authority_error)?;
149                let registered =
150                    repo::device_catalog::load(home, spool).map_err(authority_error)?;
151                if native.source_author()?.is_some() {
152                    replica.verify_source_authority(native, &authority, &registered.capability_path, now)
153                } else {
154                    repo::thread_replication::metadata::verify_control_authority(native, &authority, &registered.capability_path, now)
155                }
156
157            })
158        })
159        .await
160    }
161    async fn remember_peer_heads(
162        &self,
163        peer: [u8; 32],
164        heads: Vec<(ThreadFacet, ContentHash)>,
165    ) -> Result<(), Error> {
166        self.execute(move |replica, _| replica.remember_peer_heads(peer, &heads))
167            .await
168    }
169    async fn record_peer_receipt(
170        &self,
171        peer: [u8; 32],
172        id: ContentHash,
173        admission: Admission,
174    ) -> Result<(), Error> {
175        self.execute(move |replica, _| replica.record_peer_receipt(peer, id, &admission))
176            .await
177    }
178    async fn settled_peer_heads(
179        &self,
180        peer: [u8; 32],
181        facets: BTreeSet<ThreadFacet>,
182        limit: usize,
183    ) -> Result<Vec<(ContentHash, Admission)>, Error> {
184        self.execute(move |replica, _| replica.settled_peer_heads(peer, &facets, limit))
185            .await
186    }
187    async fn needed_from_peer(
188        &self,
189        peer: [u8; 32],
190        facets: BTreeSet<ThreadFacet>,
191        limit: usize,
192    ) -> Result<Vec<ContentHash>, Error> {
193        self.execute(move |replica, _| replica.needed_from_peer(peer, &facets, limit))
194            .await
195    }
196}
197
198fn authority_error(error: impl std::fmt::Display) -> repo::thread_replication::Error {
199    repo::thread_replication::Error::Invalid(error.to_string())
200}