heddle_thread_api/replication/
native.rs1use 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 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 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, ®istered.capability_path, now)
153 } else {
154 repo::thread_replication::metadata::verify_control_authority(native, &authority, ®istered.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}