Skip to main content

heddle_thread_api/replication/native/
hosted.rs

1//! A native relay retains the complete public bundle and rechecks durable trust
2//! on receive and export. Current access remains the receiver's callback.
3use 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    /// The directory, selected root, owner history and current access callback
36    /// come from the receiver, independently of an incoming authority bundle.
37    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    /// Prepared receiver metadata remains bound to the selected authority.
69    /// Overlay only the public set/proofs; every original byte must agree.
70    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    /// Authenticate import parents for read-only source staging under the
120    /// receiver's selected root/history and trusted current time. Publication
121    /// rechecks the newest durable authority inside its commit transaction.
122    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    /// Recheck retained originals before a source relay. Run this blocking
173    /// operation on the caller's storage worker, with its current access gate.
174    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}