1use crypto::thread_operation::SignedOperation;
4use objects::{
5 object::{Blob, State, StateId, Tree},
6 store::{ExternalObjectSource, FsStore},
7};
8use repo::thread_replication::{
9 delegated_import::AcceptedAuthority,
10 hosted_trust::{Clock, HostedTrust},
11};
12
13use super::*;
14use crate::hybrid::authority::{PublicEvidence, PublicProof};
15
16pub struct HostedReplica<C, A> {
17 local: LocalReplica<FsStore>,
18 directory: PathBuf,
19 trust: Arc<HostedTrust<C>>,
20 authority: Arc<A>,
21 export_bundle: Option<PublicProof>,
22}
23impl<C, A> Clone for HostedReplica<C, A> {
24 fn clone(&self) -> Self {
25 Self {
26 local: self.local.clone(),
27 directory: self.directory.clone(),
28 trust: self.trust.clone(),
29 authority: self.authority.clone(),
30 export_bundle: self.export_bundle.clone(),
31 }
32 }
33}
34impl LocalReplica<FsStore> {
35 pub fn with_hosted_authority<C: Clock, A: AcceptedAuthority>(
38 self,
39 directory: PathBuf,
40 trust: Arc<HostedTrust<C>>,
41 authority: Arc<A>,
42 ) -> HostedReplica<C, A> {
43 HostedReplica {
44 local: self,
45 directory,
46 trust,
47 authority,
48 export_bundle: None,
49 }
50 }
51}
52struct ExistingObjects(Arc<FsStore>);
53impl ExternalObjectSource for ExistingObjects {
54 fn get_blob(&self, hash: &ContentHash) -> objects::store::Result<Option<Blob>> {
55 self.0.get_blob(hash)
56 }
57 fn get_tree(&self, hash: &ContentHash) -> objects::store::Result<Option<Tree>> {
58 self.0.get_tree(hash)
59 }
60 fn get_state(&self, id: &StateId) -> objects::store::Result<Option<State>> {
61 self.0.get_state(id)
62 }
63 fn list_states(&self) -> objects::store::Result<Vec<StateId>> {
64 self.0.list_states()
65 }
66}
67impl<C: Clock + 'static, A: AcceptedAuthority + Send + Sync + 'static> HostedReplica<C, A> {
68 pub fn with_export_bundle(
71 mut self,
72 bundle: crate::contract::ImportPublicProofBundleV1,
73 ) -> Self {
74 self.export_bundle = Some(PublicProof::from(bundle));
75 self
76 }
77 pub fn with_native_export_bundle(
78 mut self,
79 bundle: crate::contract::NativePublicProofBundleV1,
80 ) -> Self {
81 self.export_bundle = Some(PublicProof::from(bundle));
82 self
83 }
84 pub fn with_export_proof(mut self, bundle: PublicProof) -> Self {
85 self.export_bundle = Some(bundle);
86 self
87 }
88 fn export_bundle(&self, mut bundle: PublicProof) -> Result<PublicProof, Error> {
89 if let Some(prepared) = &self.export_bundle {
90 bundle
91 .replace_receiver_metadata(prepared.clone())
92 .map_err(|e| Error::Store(e.into()))?;
93 }
94 Ok(bundle)
95 }
96 pub fn publish_source(
97 &self,
98 source: crate::fetch::StagedSource,
99 repository: &repo::Repository,
100 now_seconds: i64,
101 publication: crate::fetch::hosted::HostedPublication<'_>,
102 response: impl FnOnce(
103 &repo::thread_replication::hosted_trust::TrustTransaction<'_>,
104 ) -> repo::thread_replication::Result<Vec<u8>>,
105 ) -> Result<Vec<u8>, crate::fetch::Error> {
106 if repository.heddle_dir().canonicalize()? != self.directory.canonicalize()? {
107 return Err(api::hybrid_codec::Reject::Root.into());
108 }
109 source.publish_hosted(
110 repository,
111 &self.trust,
112 self.authority.as_ref(),
113 now_seconds,
114 publication,
115 response,
116 )
117 }
118
119 pub fn authenticate_import_carriers(
123 &self,
124 bundle: &crate::contract::ImportPublicProofBundleV1,
125 now_millis: i64,
126 ) -> repo::thread_replication::Result<crypto::import_authority::VerifiedImportCarriers> {
127 let snapshot = self.trust.snapshot()?;
128 let keys = snapshot
129 .known_job_associations
130 .iter()
131 .map(|(key, _)| key.clone())
132 .collect::<Vec<_>>();
133 let set = api::witness_trust::verify_set(
134 bundle
135 .witness_set
136 .as_ref()
137 .ok_or(api::hybrid_codec::Reject::Canonical)?,
138 &api::witness_trust::SetExpectation {
139 authority: &snapshot.root.authority,
140 root_id: &snapshot.root.root_id,
141 root_public_key: &snapshot.root.public_key,
142 root_epoch: snapshot.root_epoch,
143 now_unix_millis: now_millis,
144 clock_floor_unix_millis: snapshot.clock_floor_millis,
145 known_job_keys: &keys,
146 },
147 snapshot.previous.as_ref(),
148 )?;
149 let mut forbidden = vec![snapshot.root.public_key.to_vec()];
150 forbidden.extend(
151 set.body()
152 .entries
153 .iter()
154 .map(|entry| entry.public_key.clone()),
155 );
156 repo::thread_replication::delegated_import::authenticate_import_carriers(
157 bundle,
158 self.authority.as_ref(),
159 &api::import_authority::ImportWitnessRootPin {
160 authority: snapshot.root.authority,
161 root_id: snapshot.root.root_id,
162 public_key: snapshot.root.public_key.to_vec(),
163 epoch: snapshot.root_epoch,
164 },
165 now_millis,
166 &snapshot.known_job_associations,
167 &forbidden,
168 |_| Ok(()),
169 )
170 }
171
172 pub fn recheck_selected(
175 &self,
176 bundle: &crate::contract::ImportPublicProofBundleV1,
177 records: &[crate::contract::SignedRecord],
178 ) -> repo::thread_replication::Result<()> {
179 self.recheck_proof(&PublicProof::from(bundle.clone()), records)
180 }
181 pub fn recheck_native_selected(
182 &self,
183 bundle: &crate::contract::NativePublicProofBundleV1,
184 records: &[crate::contract::SignedRecord],
185 ) -> repo::thread_replication::Result<()> {
186 self.recheck_proof(&PublicProof::from(bundle.clone()), records)
187 }
188 pub fn recheck_proof(
189 &self,
190 bundle: &PublicProof,
191 records: &[crate::contract::SignedRecord],
192 ) -> repo::thread_replication::Result<()> {
193 if self.local.objects.root().canonicalize()? != self.directory.canonicalize()? {
194 return Err(api::hybrid_codec::Reject::Root.into());
195 }
196 let scratch = tempfile::tempdir()?;
197 let mut staged = FsStore::new(scratch.path());
198 staged.set_external_source(Arc::new(ExistingObjects(self.local.objects.clone())));
199 staged.init()?;
200 if matches!(bundle, PublicProof::Native(_)) {
201 ThreadReplica::install_hybrid_native(
202 &self.directory,
203 &self.trust,
204 &bundle.encode_to_vec(),
205 records,
206 self.authority.as_ref(),
207 &staged,
208 |artifacts| crate::fetch::hosted::publish_store(staged.root(), artifacts),
209 )
210 } else {
211 ThreadReplica::install_hybrid_import(
212 &self.directory,
213 &self.trust,
214 &bundle.encode_to_vec(),
215 records,
216 self.authority.as_ref(),
217 &staged,
218 |artifacts| crate::fetch::hosted::publish_store(staged.root(), artifacts),
219 )
220 }?;
221 self.local.objects.reload_packs()?;
222 Ok(())
223 }
224 async fn admit(
225 &self,
226 original: SignedOperation,
227 bundle: PublicProof,
228 ) -> Result<Admission, Error> {
229 let receiver = self.clone();
230 Ok(tokio::task::spawn_blocking(move || {
231 let operation = original.verify()?;
232 if operation.thread != receiver.local.replica.thread_id() {
233 return Err(repo::thread_replication::Error::Hybrid(
234 api::hybrid_codec::Reject::Scope,
235 ));
236 }
237 if receiver.local.objects.root().canonicalize()? != receiver.directory.canonicalize()? {
238 return Err(repo::thread_replication::Error::Hybrid(
239 api::hybrid_codec::Reject::Root,
240 ));
241 }
242 let id = operation.id()?;
243 let mut pending = vec![original];
244 let mut selected = BTreeSet::new();
245 let mut records = Vec::new();
246 while let Some(original) = pending.pop() {
247 let operation = original.verify()?;
248 if !selected.insert(operation.id()?) {
249 continue;
250 }
251 if selected.len() > 10_000 {
252 return Err(repo::thread_replication::Error::Hybrid(
253 api::hybrid_codec::Reject::Bounds,
254 ));
255 }
256 for parent in &operation.parents {
257 let (original, status) = receiver
258 .local
259 .replica
260 .operation(parent)?
261 .ok_or(api::hybrid_codec::Reject::Scope)?;
262 if !matches!(status, Admission::Accepted) {
263 return Err(repo::thread_replication::Error::Hybrid(
264 api::hybrid_codec::Reject::Scope,
265 ));
266 }
267 pending.push(original);
268 }
269 records.push(crate::contract::SignedRecord {
270 format: heddle_object_model::object::thread_replication::OPERATION_FORMAT
271 .into(),
272 canonical_record: original.canonical,
273 signatures: vec![crate::contract::RecordSignature {
274 public_key: operation.publisher.to_vec(),
275 signature: original.signature,
276 }],
277 });
278 }
279 receiver.recheck_proof(&bundle, &records)?;
280 receiver
281 .local
282 .replica
283 .operation(&id)?
284 .map(|(_, status)| status)
285 .ok_or_else(|| api::hybrid_codec::Reject::Scope.into())
286 })
287 .await??)
288 }
289}
290impl<C: Clock + 'static, A: AcceptedAuthority + Send + Sync + 'static> ReplicaStore
291 for HostedReplica<C, A>
292{
293 type Error = Error;
294 fn thread_id(&self) -> ContentHash {
295 self.local.thread_id()
296 }
297 async fn generation(&self) -> Result<i64, Error> {
298 self.local.generation().await
299 }
300 async fn sharing(&self, peer: [u8; 32]) -> Result<BTreeSet<ThreadFacet>, Error> {
301 self.local.sharing(peer).await
302 }
303 async fn frontier_page(
304 &self,
305 facet: ThreadFacet,
306 after: Option<ContentHash>,
307 limit: usize,
308 ) -> Result<Vec<ContentHash>, Error> {
309 self.local.frontier_page(facet, after, limit).await
310 }
311 async fn operation(
312 &self,
313 id: ContentHash,
314 ) -> Result<Option<(ReceivedOperation, Admission)>, Error> {
315 let stored = self
316 .local
317 .execute(move |replica, _| {
318 Ok((
319 replica.operation_with_authority_admission(&id)?,
320 replica.hybrid_import_bundle()?,
321 replica.hybrid_native_bundle()?,
322 ))
323 })
324 .await?;
325 let (Some(stored), imported, native) = stored else {
326 return Ok(None);
327 };
328 let bundle = match (imported, native) {
329 (Some(b), None) => PublicProof::from(b),
330 (None, Some(b)) => PublicProof::from(b),
331 (None, None) => return Err(Error::HostedTrustRequired),
332 _ => return Err(Error::HostedTrustRequired),
333 };
334 if stored.authority_admission.is_some() {
335 return Err(Error::HostedTrustRequired);
336 }
337 let bundle = self.export_bundle(bundle)?;
338 let status = self.admit(stored.original.clone(), bundle.clone()).await?;
339 Ok(Some((
340 ReceivedOperation {
341 native_authority: bundle.native().cloned().map(Arc::new),
342 original: stored.original,
343 authority_admission: None,
344 import_authority: bundle.imported().cloned().map(Arc::new),
345 },
346 status,
347 )))
348 }
349 async fn receive(&self, received: ReceivedOperation) -> Result<Admission, Error> {
350 if received.authority_admission.is_some() {
351 return Err(Error::HostedTrustRequired);
352 }
353 let bundle = match (received.import_authority, received.native_authority) {
354 (Some(b), None) => PublicProof::from(b.as_ref().clone()),
355 (None, Some(b)) => PublicProof::from(b.as_ref().clone()),
356 _ => return Err(Error::HostedTrustRequired),
357 };
358 let bundle = self.export_bundle(bundle)?;
359 self.admit(received.original, bundle).await
360 }
361
362 async fn remember_peer_heads(
363 &self,
364 peer: [u8; 32],
365 heads: Vec<(ThreadFacet, ContentHash)>,
366 ) -> Result<(), Error> {
367 self.local.remember_peer_heads(peer, heads).await
368 }
369 async fn record_peer_receipt(
370 &self,
371 peer: [u8; 32],
372 id: ContentHash,
373 admission: Admission,
374 ) -> Result<(), Error> {
375 self.local.record_peer_receipt(peer, id, admission).await
376 }
377 async fn settled_peer_heads(
378 &self,
379 peer: [u8; 32],
380 facets: BTreeSet<ThreadFacet>,
381 limit: usize,
382 ) -> Result<Vec<(ContentHash, Admission)>, Error> {
383 self.local.settled_peer_heads(peer, facets, limit).await
384 }
385 async fn needed_from_peer(
386 &self,
387 peer: [u8; 32],
388 facets: BTreeSet<ThreadFacet>,
389 limit: usize,
390 ) -> Result<Vec<ContentHash>, Error> {
391 self.local.needed_from_peer(peer, facets, limit).await
392 }
393}