Skip to main content

heddle_thread_api/fetch/
hosted.rs

1//! Selected hosted installation keeps repository artifacts inside the same
2//! trust transaction as native admission. Staging stores supply no authority.
3use std::path::Path;
4
5use heddle_object_model::object::StateId;
6use objects::{lock::RepositoryLockExt, store::ObjectStore};
7use prost::Message;
8use repo::{
9    Repository,
10    thread_replication::{
11        ThreadReplica,
12        delegated_import::AcceptedAuthority,
13        hosted_trust::{Clock, HostedTrust, TrustTransaction},
14        install_artifacts::InstallArtifacts,
15    },
16};
17
18use super::{Error, StagedSource};
19
20#[cfg(test)]
21#[path = "import_ancestry_tests.rs"]
22mod import_ancestry_tests;
23
24pub struct HostedPublication<'a> {
25    pub replica: &'a ThreadReplica,
26    pub prepared: repo::thread_replication::source_publication::PreparedPublication<'a>,
27    pub command: repo::thread_replication::source_publication::Command<'a>,
28}
29
30impl StagedSource {
31    /// `authority` and `trust` are independently selected by the receiver.
32    /// Every selected original is admitted again under the durable trust lock;
33    /// the complete public history never selects unrelated native branches.
34    pub fn install_hosted(
35        self,
36        repository: &Repository,
37        trust: &HostedTrust<impl Clock>,
38        authority: &impl AcceptedAuthority,
39        now_seconds: i64,
40    ) -> Result<StateId, Error> {
41        self.install_hosted_commit(repository, trust, authority, now_seconds, None, |_| {
42            Ok(Vec::new())
43        })
44        .map(|(state, _)| state)
45    }
46    pub fn publish_hosted(
47        self,
48        repository: &Repository,
49        trust: &HostedTrust<impl Clock>,
50        authority: &impl AcceptedAuthority,
51        now_seconds: i64,
52        publication: HostedPublication<'_>,
53        response: impl FnOnce(
54            &repo::thread_replication::hosted_trust::TrustTransaction<'_>,
55        ) -> repo::thread_replication::Result<Vec<u8>>,
56    ) -> Result<Vec<u8>, Error> {
57        self.install_hosted_commit(
58            repository,
59            trust,
60            authority,
61            now_seconds,
62            Some(publication),
63            response,
64        )
65        .map(|(_, receipt)| receipt)
66    }
67    #[allow(clippy::too_many_arguments)]
68    fn install_hosted_commit(
69        self,
70        repository: &Repository,
71        trust: &HostedTrust<impl Clock>,
72        authority: &impl AcceptedAuthority,
73        now_seconds: i64,
74        publication: Option<HostedPublication<'_>>,
75        response: impl FnOnce(
76            &repo::thread_replication::hosted_trust::TrustTransaction<'_>,
77        ) -> repo::thread_replication::Result<Vec<u8>>,
78    ) -> Result<(StateId, Vec<u8>), Error> {
79        crate::hybrid::transfer_ready(&self.ready).map_err(Error::Invalid)?;
80        let imported = self.import_authority();
81        let native = self.native_authority();
82        let (bundle, bundle_owner) = match (imported, native) {
83            (Some(b), None) => (b.encode_to_vec(), b.owner_genesis.as_ref()),
84            (None, Some(b)) => (b.encode_to_vec(), b.owner_genesis.as_ref()),
85            _ => return Err(Error::HostedTrustRequired),
86        };
87        let spool = self
88            .ready
89            .thread
90            .as_ref()
91            .and_then(|t| t.spool.as_ref())
92            .ok_or(Error::Invalid("Spool absent"))?
93            .id
94            .parse::<uuid::Uuid>()
95            .map_err(preparation)?;
96        let owner_genesis = self
97            .ready
98            .owner_genesis
99            .as_ref()
100            .ok_or(Error::Invalid("owner genesis absent"))?;
101        let owner = self
102            .ready
103            .ownership
104            .as_ref()
105            .ok_or(Error::Invalid("owner history absent"))?;
106        let selected =
107            repo::verify_spool_owner_observation(owner_genesis, owner, spool, now_seconds)
108                .map_err(preparation)?;
109        if bundle_owner != Some(selected.owner_genesis().signed()) {
110            return Err(api::hybrid_codec::Reject::Root.into());
111        }
112        let _write_lock = repository.locker().write().map_err(preparation)?;
113        let next_pin = repository
114            .prepare_owner_observation_pin(
115                owner_genesis,
116                owner,
117                spool,
118                &selected.wire().canonical_spool_path_segments,
119                now_seconds,
120            )
121            .map_err(preparation)?;
122        let previous_spool = read_optional(&repository.heddle_dir().join("spool-id"))?;
123        if previous_spool.as_ref().is_some_and(|bytes| {
124            std::str::from_utf8(bytes)
125                .ok()
126                .and_then(|s| s.trim().parse::<uuid::Uuid>().ok())
127                != Some(spool)
128        }) {
129            return Err(api::hybrid_codec::Reject::Root.into());
130        }
131        let pin_path = repository.heddle_dir().join("owner-authorization.bin");
132        let previous_pin = read_optional(&pin_path)?;
133        let staging = tempfile::tempdir_in(self.directory.path())?;
134        let staged_repo = Repository::init(staging.path()).map_err(preparation)?;
135        self.install_source_objects(&staged_repo)?;
136        // Native checkout comparison needs the immutable empty base. Keep its
137        // tree/state in staging so every imported artifact shares the journal.
138        let seed = objects::object::thread_replication::initial_base::synthetic_initial_base()
139            .map_err(preparation)?;
140        staged_repo
141            .store()
142            .put_snapshot_objects_packed(Vec::new(), &objects::object::Tree::new(), &seed)
143            .map_err(preparation)?;
144        let main = self
145            .ready
146            .thread_genesis
147            .as_ref()
148            .ok_or(Error::Invalid("Thread genesis absent"))?;
149        let main_id = crate::replication::opening::verify_genesis(
150            main.genesis
151                .as_ref()
152                .ok_or(Error::Invalid("signed genesis absent"))?,
153            self.ready
154                .thread
155                .as_ref()
156                .ok_or(Error::Invalid("Thread absent"))?,
157        )?
158        .id()
159        .map_err(preparation)?;
160        let mut records = Vec::new();
161        for wrapper in std::iter::once(main).chain(&self.dependencies) {
162            records.push(
163                wrapper
164                    .genesis
165                    .clone()
166                    .ok_or(Error::Invalid("original genesis absent"))?,
167            );
168            records.extend(wrapper.ownership_claims.clone());
169            records.extend(wrapper.ownership_resolutions.clone());
170        }
171        for signed in &self.operations {
172            let operation = signed.verify().map_err(preparation)?;
173            records.push(crate::contract::SignedRecord {
174                format: heddle_object_model::object::thread_replication::OPERATION_FORMAT.into(),
175                canonical_record: signed.canonical.clone(),
176                signatures: vec![crate::contract::RecordSignature {
177                    public_key: operation.publisher.to_vec(),
178                    signature: signed.signature.clone(),
179                }],
180            });
181        }
182        records.extend(self.prefix_original.iter().cloned());
183        if let Some(bundle) = native {
184            for wrapper in std::iter::once(main).chain(&self.dependencies) {
185                let genesis = wrapper
186                    .genesis
187                    .as_ref()
188                    .ok_or(Error::Invalid("genesis absent"))?;
189                let (_, decoded) = crypto::import_authority::verify_native_genesis(genesis)
190                    .map_err(preparation)?;
191                let id = decoded.id().map_err(preparation)?;
192                if !bundle
193                    .foreign_dependencies
194                    .iter()
195                    .any(|r| r.thread_genesis_digest.as_slice() == id.as_bytes())
196                {
197                    require_native_genesis_match(bundle, wrapper)?;
198                }
199            }
200        }
201        let state = self.state.id();
202        let publish = |artifacts: &mut InstallArtifacts<'_>| {
203            // An owner/spool update while staging requires fresh preparation.
204            if read_optional(&pin_path).map_err(replica_error)? != previous_pin
205                || read_optional(&repository.heddle_dir().join("spool-id"))
206                    .map_err(replica_error)?
207                    != previous_spool
208            {
209                return Err(repo::thread_replication::Error::Hybrid(
210                    api::hybrid_codec::Reject::StaleContext,
211                ));
212            }
213            publish_store(staged_repo.heddle_dir(), artifacts)?;
214            artifacts.write_file(Path::new("owner-authorization.bin"), &next_pin)?;
215            artifacts.write_file(Path::new("spool-id"), spool.to_string().as_bytes())?;
216            Ok(())
217        };
218        let publish_with_tiers =
219            |context: &TrustTransaction<'_>, artifacts: &mut InstallArtifacts<'_>| {
220                for floor in &self.ancestry {
221                    context.record_import_floor_tiers(main_id, floor.tip, &floor.tiers)?;
222                }
223                publish(artifacts)
224            };
225        let (replicas, receipt) = if let Some(publication) = publication {
226            let receipt = if native.is_some() {
227                ThreadReplica::publish_native_source(
228                    publication.replica,
229                    trust,
230                    &bundle,
231                    &records,
232                    authority,
233                    staged_repo.store(),
234                    publication.prepared,
235                    publication.command,
236                    |context, artifacts| {
237                        publish_with_tiers(context, artifacts)?;
238                        response(context)
239                    },
240                )
241            } else {
242                ThreadReplica::publish_hybrid_source(
243                    publication.replica,
244                    trust,
245                    &bundle,
246                    &records,
247                    authority,
248                    staged_repo.store(),
249                    publication.prepared,
250                    publication.command,
251                    |context, artifacts| {
252                        publish_with_tiers(context, artifacts)?;
253                        response(context)
254                    },
255                )
256            }
257            .map_err(Error::from)?;
258            (Vec::new(), receipt)
259        } else {
260            let replicas = if native.is_some() {
261                ThreadReplica::install_hybrid_native_with(
262                    repository.heddle_dir(),
263                    trust,
264                    &bundle,
265                    &records,
266                    authority,
267                    staged_repo.store(),
268                    |_| Ok(()),
269                    publish_with_tiers,
270                )
271            } else {
272                ThreadReplica::install_hybrid_import_with(
273                    repository.heddle_dir(),
274                    trust,
275                    &bundle,
276                    &records,
277                    authority,
278                    staged_repo.store(),
279                    |_| Ok(()),
280                    publish_with_tiers,
281                )
282            }
283            .map_err(Error::from)?;
284            (replicas, Vec::new())
285        };
286        repository.store().reload_packs().map_err(preparation)?;
287        if !replicas.is_empty() {
288            let selected = replicas
289                .iter()
290                .find(|r| r.thread_id() == main_id)
291                .ok_or(Error::Invalid("selected replica absent"))?;
292            if self.is_complete() {
293                selected
294                    .record_source_possession(state)
295                    .map_err(preparation)?;
296            }
297        }
298        if !replicas.is_empty() {
299            // Account/device lookup uses the committed Spool identity. Registration
300            // is local discovery after authority commit, while the repository lock
301            // still prevents another writer from replacing the installed identity.
302            repo::device_catalog::register(&repo::identity::heddle_home_dir(), repository, spool)
303                .map_err(preparation)?;
304        }
305        Ok((state, receipt))
306    }
307}
308
309fn require_native_genesis_match(
310    bundle: &crate::contract::NativePublicProofBundleV1,
311    wrapper: &crate::contract::ThreadGenesisRecord,
312) -> Result<(), Error> {
313    if !bundle.genesis_witnesses.iter().any(|p| {
314        p.original_genesis == wrapper.genesis
315            && p.creator_authority_envelope == wrapper.creator_authority
316            && p.binding == wrapper.native_genesis_authority
317    }) {
318        return Err(api::hybrid_codec::Reject::GenesisBinding.into());
319    }
320    Ok(())
321}
322
323#[cfg(test)]
324mod binding_tests {
325    use super::*;
326    #[test]
327    fn native_ready_binding_must_match_the_exact_witnessed_genesis() {
328        let fixture: serde_json::Value = serde_json::from_str(include_str!(
329            "../../tests/fixtures/native-host-witness-v1.json"
330        ))
331        .expect("vectors");
332        let bundle: crate::contract::NativePublicProofBundleV1 = api::hybrid_codec::strict_decode(
333            &hex::decode(
334                fixture["wire_vectors"]["start_thread"]["wire_hex"]
335                    .as_str()
336                    .expect("wire"),
337            )
338            .expect("hex"),
339            api::import_authority::MAX_BUNDLE_BYTES,
340        )
341        .expect("bundle");
342        let p = &bundle.genesis_witnesses[0];
343        let control = crate::contract::ThreadGenesisRecord {
344            genesis: p.original_genesis.clone(),
345            creator_authority: p.creator_authority_envelope.clone(),
346            native_genesis_authority: p.binding.clone(),
347            ..Default::default()
348        };
349        require_native_genesis_match(&bundle, &control).expect("exact Ready control");
350        for field in 0..3 {
351            let mut changed = control.clone();
352            match field {
353                0 => changed.genesis = None,
354                1 => changed.creator_authority.push(0),
355                _ => changed.native_genesis_authority = None,
356            }
357            assert!(
358                matches!(
359                    require_native_genesis_match(&bundle, &changed),
360                    Err(Error::Hybrid(api::hybrid_codec::Reject::GenesisBinding))
361                ),
362                "Ready original, envelope and binding must match"
363            );
364        }
365        require_native_genesis_match(&bundle, &control).expect("unchanged Ready control");
366    }
367}
368
369fn read_optional(path: &Path) -> Result<Option<Vec<u8>>, Error> {
370    match std::fs::read(path) {
371        Ok(bytes) => Ok(Some(bytes)),
372        Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
373        Err(e) => Err(e.into()),
374    }
375}
376fn replica_error(error: impl std::fmt::Display) -> repo::thread_replication::Error {
377    repo::thread_replication::Error::Invalid(error.to_string())
378}
379
380pub(crate) fn publish_store(
381    staged: &Path,
382    artifacts: &mut InstallArtifacts<'_>,
383) -> repo::thread_replication::Result<()> {
384    // Only immutable source storage crosses the boundary. Staging's refs,
385    // identities, locks, configuration and local metadata never join the clone.
386    for name in ["packs", "objects"] {
387        publish_directory(&staged.join(name), Path::new(name), artifacts)?;
388    }
389    Ok(())
390}
391fn publish_directory(
392    staged: &Path,
393    destination: &Path,
394    artifacts: &mut InstallArtifacts<'_>,
395) -> repo::thread_replication::Result<()> {
396    for entry in std::fs::read_dir(staged)? {
397        let entry = entry?;
398        if entry.file_type()?.is_dir() {
399            // Pack install bookkeeping is local to the staging store.
400            if entry
401                .file_name()
402                .to_str()
403                .is_some_and(|name| name.starts_with('.'))
404            {
405                continue;
406            }
407            publish_directory(
408                &entry.path(),
409                &destination.join(entry.file_name()),
410                artifacts,
411            )?;
412        } else if entry.file_type()?.is_file() {
413            if entry
414                .file_name()
415                .to_str()
416                .is_some_and(|name| name.starts_with('.'))
417            {
418                continue;
419            }
420            artifacts.install_file(&entry.path(), &destination.join(entry.file_name()))?;
421        }
422    }
423    Ok(())
424}
425fn preparation(error: impl std::fmt::Display) -> Error {
426    Error::Preparation(error.to_string())
427}
428
429#[cfg(test)]
430pub(crate) mod tests {
431    use std::{
432        collections::BTreeMap,
433        sync::atomic::{AtomicUsize, Ordering},
434    };
435
436    use objects::{
437        object::{State, Tree},
438        store::{
439            ObjectStore,
440            pack::{ObjectType, PackBuilder, PackObjectId},
441        },
442    };
443    use repo::thread_replication::hosted_trust::*;
444
445    use super::*;
446    use crate::{
447        contract::*,
448        hybrid::authority::{
449            AcceptedHistory, SelectedAuthority,
450            tests::{bundle, selected},
451        },
452    };
453
454    pub(super) struct ReceiverClock;
455    impl Clock for ReceiverClock {
456        fn now_millis(&self) -> repo::thread_replication::Result<i64> {
457            Ok(1_350_000)
458        }
459        fn elapsed_millis(&self) -> repo::thread_replication::Result<u64> {
460            Ok(0)
461        }
462    }
463    pub(super) fn record<T: Message + Default>(fixture: &serde_json::Value, name: &str) -> T {
464        let vector = fixture["wire_vectors"]
465            .get(name)
466            .or_else(|| fixture["signed_vectors"].get(name))
467            .expect("published vector");
468        T::decode(
469            hex::decode(vector["wire_hex"].as_str().expect("wire bytes"))
470                .expect("hex")
471                .as_slice(),
472        )
473        .expect("original record")
474    }
475    pub(crate) fn source(
476        scratch: &Path,
477        seed_only: bool,
478    ) -> (
479        StagedSource,
480        RootSelection,
481        heddleco_capability_verifier::VerifiedCloneKeyring,
482    ) {
483        let fixture: serde_json::Value =
484            serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
485                .expect("release fixture");
486        let mut bundle = bundle();
487        bundle.history_proofs = [
488            "genesis_proof",
489            "genesis_dev_proof",
490            "publication_proof",
491            "dev_publication_proof",
492        ]
493        .map(|name| record(&fixture, name))
494        .to_vec();
495        let limits = heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
496            .expect("limits");
497        let pinned = selected(&bundle, limits);
498        let original: SignedRecord = record(&fixture, "converted_main");
499        let signed = crate::replication::decode_record(original).expect("converted original");
500        let operation = signed.verify().expect("original signature");
501        let genesis_record = bundle
502            .original_geneses
503            .iter()
504            .find(|record| {
505                crypto::thread_operation::SignedGenesis {
506                    canonical: record.canonical_record.clone(),
507                    signature: record.signatures[0].signature.clone(),
508                }
509                .verify()
510                .is_ok_and(|g| g.id().expect("genesis ID") == operation.thread)
511            })
512            .expect("selected genesis")
513            .clone();
514        let genesis = crypto::thread_operation::SignedGenesis {
515            canonical: genesis_record.canonical_record.clone(),
516            signature: genesis_record.signatures[0].signature.clone(),
517        }
518        .verify()
519        .expect("genesis signature");
520        let state: State = if seed_only {
521            objects::object::thread_replication::initial_base::synthetic_initial_base()
522                .expect("empty seed")
523        } else {
524            operation.source_state().expect("source").expect("State")
525        };
526        let tree = Tree::new();
527        assert_eq!(
528            tree.hash(),
529            state.tree,
530            "published native conversion has the empty source tree"
531        );
532        let mut builder = PackBuilder::for_repack(Default::default(), 0);
533        builder.add_id(
534            PackObjectId::StateId(state.id()),
535            ObjectType::State,
536            state.encode_current_msgpack().expect("State"),
537        );
538        builder.add_id(
539            PackObjectId::Hash(tree.hash()),
540            ObjectType::Tree,
541            tree.encode_canonical().expect("Tree"),
542        );
543        let (pack, index, _) = builder.build().expect("source pack");
544        let directory = tempfile::tempdir_in(scratch).expect("staging");
545        std::fs::write(directory.path().join("source.pack"), pack).expect("pack");
546        std::fs::write(directory.path().join("source.idx"), index).expect("index");
547        let spool = SpoolRef {
548            id: genesis.spool.clone(),
549        };
550        let owner = pinned.owner_state();
551        let history = AcceptedHistory::from_selected_spool(&bundle, &pinned, 1350, limits)
552            .expect("verified import owner history");
553        let authority = SelectedAuthority::new(
554            history,
555            bundle.clone(),
556            |_: &ImportPublicProofBundleV1, _: i64, _: &TrustTransaction<'_>| Ok(()),
557        );
558        let pin = api::import_authority::ImportWitnessRootPin {
559            authority: "https://weft.example.test".into(),
560            root_id: "descriptor-root-1".into(),
561            public_key: hex::decode(
562                fixture["keys"]["root"]["public_key_hex"]
563                    .as_str()
564                    .expect("root"),
565            )
566            .expect("hex"),
567            epoch: 1,
568        };
569        let carriers = repo::thread_replication::delegated_import::authenticate_import_carriers(
570            &bundle,
571            &authority,
572            &pin,
573            1_350_000,
574            &[],
575            &[],
576            |_| Ok(()),
577        )
578        .expect("independently authenticated import carriers");
579        let ready = TransferReady {
580            thread: Some(ThreadRef {
581                spool: Some(spool.clone()),
582                id: Some(ThreadId {
583                    value: operation.thread.as_bytes().to_vec(),
584                }),
585            }),
586            current: Some(RevisionRef {
587                spool: Some(spool),
588                revision: Some(revision_ref::Revision::State(
589                    api::heddle::api::common::StateId {
590                        value: state.id().as_bytes().to_vec(),
591                    },
592                )),
593            }),
594            thread_genesis: Some(ThreadGenesisRecord {
595                creator_authority: bundle
596                    .genesis_witnesses
597                    .iter()
598                    .find(|p| p.original_genesis.as_ref() == Some(&genesis_record))
599                    .expect("selected creator authority")
600                    .creator_authority_envelope
601                    .clone(),
602                genesis: Some(genesis_record),
603                ..Default::default()
604            }),
605            owner_genesis: bundle.owner_genesis.clone(),
606            ownership: Some(OwnerState {
607                owner: Some(PrincipalRef {
608                    id: uuid::Uuid::from_bytes(
609                        owner
610                            .signed_root()
611                            .root
612                            .as_ref()
613                            .expect("root")
614                            .account_uuid
615                            .as_slice()
616                            .try_into()
617                            .expect("UUID"),
618                    )
619                    .to_string(),
620                }),
621                root: Some(owner.signed_root().clone()),
622                accepted_transitions: pinned.wire().accepted_transitions.clone(),
623                version: owner.state_hash().to_vec(),
624                resource_keyring: Some(pinned.wire().clone()),
625                ..Default::default()
626            }),
627            full_closure_available: true,
628            import_authority: Some(bundle),
629            protocol: Some(crate::hybrid::protocol()),
630            ..Default::default()
631        };
632        let staged = super::super::staging::validate_with_receipts_and_carriers(
633            directory,
634            ready,
635            if seed_only { vec![] } else { vec![signed] },
636            vec![],
637            vec![],
638            Some(carriers),
639            Default::default(),
640        )
641        .expect("selected structural source");
642        let root = RootSelection {
643            authority: "https://weft.example.test".into(),
644            root_id: "descriptor-root-1".into(),
645            public_key: hex::decode(
646                fixture["keys"]["root"]["public_key_hex"]
647                    .as_str()
648                    .expect("root"),
649            )
650            .expect("hex")
651            .try_into()
652            .expect("key"),
653        };
654        (staged, root, pinned)
655    }
656    fn artifacts(path: &Path) -> BTreeMap<std::path::PathBuf, Vec<u8>> {
657        fn walk(root: &Path, path: &Path, values: &mut BTreeMap<std::path::PathBuf, Vec<u8>>) {
658            if !path.exists() {
659                return;
660            }
661            for entry in std::fs::read_dir(path).expect("directory") {
662                let entry = entry.expect("entry");
663                if entry.file_type().expect("type").is_dir() {
664                    walk(root, &entry.path(), values);
665                } else {
666                    values.insert(
667                        entry
668                            .path()
669                            .strip_prefix(root)
670                            .expect("relative")
671                            .to_path_buf(),
672                        std::fs::read(entry.path()).expect("bytes"),
673                    );
674                }
675            }
676        }
677        let mut values = BTreeMap::new();
678        for name in ["objects", "packs"] {
679            walk(path, &path.join(name), &mut values);
680        }
681        for name in ["owner-authorization.bin", "spool-id"] {
682            if let Ok(bytes) = std::fs::read(path.join(name)) {
683                values.insert(name.into(), bytes);
684            }
685        }
686        values
687    }
688    #[test]
689    fn selected_hosted_capture_and_genesis_only_install_without_sibling_branch() {
690        for seed_only in [false, true] {
691            let scratch = tempfile::tempdir().expect("scratch");
692            let directory = tempfile::tempdir().expect("receiver");
693            let repo = Repository::init(directory.path()).expect("repository");
694            let (staged, root, pinned) = source(scratch.path(), seed_only);
695            let state = staged.state().id();
696            let selected_thread = staged
697                .ready()
698                .thread
699                .as_ref()
700                .expect("Thread")
701                .id
702                .as_ref()
703                .expect("ID")
704                .value
705                .clone();
706            let bundle = staged.import_authority().expect("public history").clone();
707            let history = AcceptedHistory::from_selected_spool(
708                &bundle,
709                &pinned,
710                1350,
711                heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
712                    .expect("limits"),
713            )
714            .expect("selected history");
715            select_root(repo.heddle_dir(), &root).expect("root pin");
716            select_spool(
717                repo.heddle_dir(),
718                pinned.owner_genesis().spool_uuid(),
719                *history.genesis(),
720                *history.initial_owner(),
721            )
722            .expect("Spool selection");
723            let trust = HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock)
724                .expect("trust");
725            let authority = SelectedAuthority::new(
726                history,
727                bundle.clone(),
728                |_: &ImportPublicProofBundleV1,
729                 _: i64,
730                 _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
731                    Ok(())
732                },
733            );
734            assert_eq!(
735                staged
736                    .install_hosted(&repo, &trust, &authority, 1350)
737                    .expect("selected hosted install"),
738                state
739            );
740            assert!(repo.heddle_dir().join("spool-id").exists());
741            assert!(repo.heddle_dir().join("owner-authorization.bin").exists());
742            let (_, retained_owner) = repo
743                .pinned_owner_observation(1350)
744                .expect("independently retained owner observation");
745            assert_eq!(
746                retained_owner.owner_genesis().signed(),
747                pinned.owner_genesis().signed()
748            );
749            assert!(
750                repo.store()
751                    .get_state(&state)
752                    .expect("stored State")
753                    .is_some()
754            );
755            for record in &bundle.original_geneses {
756                let genesis = crypto::thread_operation::SignedGenesis {
757                    canonical: record.canonical_record.clone(),
758                    signature: record.signatures[0].signature.clone(),
759                }
760                .verify()
761                .expect("original genesis");
762                let id = genesis.id().expect("ID");
763                if id.as_bytes().as_slice() == selected_thread {
764                    assert_eq!(
765                        ThreadReplica::open(repo.heddle_dir(), id)
766                            .expect("selected replica")
767                            .hybrid_import_bundle()
768                            .expect("retained bundle"),
769                        Some(bundle.clone())
770                    );
771                } else {
772                    assert!(
773                        ThreadReplica::open(repo.heddle_dir(), id).is_err(),
774                        "public sibling proof must not install its native branch"
775                    );
776                }
777            }
778        }
779    }
780    #[test]
781    fn late_disclosure_failure_rolls_back_pack_owner_pin_spool_and_trust() {
782        let scratch = tempfile::tempdir().expect("scratch");
783        let directory = tempfile::tempdir().expect("receiver");
784        let repo = Repository::init(directory.path()).expect("repository");
785        let before = artifacts(repo.heddle_dir());
786        let (staged, root, pinned) = source(scratch.path(), false);
787        let bundle = staged.import_authority().expect("history").clone();
788        let history = AcceptedHistory::from_selected_spool(
789            &bundle,
790            &pinned,
791            1350,
792            heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
793                .expect("limits"),
794        )
795        .expect("selected history");
796        select_root(repo.heddle_dir(), &root).expect("root");
797        select_spool(
798            repo.heddle_dir(),
799            pinned.owner_genesis().spool_uuid(),
800            *history.genesis(),
801            *history.initial_owner(),
802        )
803        .expect("Spool");
804        let trust =
805            HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock).expect("trust");
806        let calls = AtomicUsize::new(0);
807        let authority = SelectedAuthority::new(
808            history,
809            bundle,
810            |_: &ImportPublicProofBundleV1,
811             _: i64,
812             _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
813                if calls.fetch_add(1, Ordering::SeqCst) > 1 {
814                    assert!(
815                        repo.heddle_dir().join("owner-authorization.bin").exists(),
816                        "failure must occur after staged artifacts were installed"
817                    );
818                    assert!(repo.heddle_dir().join("spool-id").exists());
819                    return Err(repo::thread_replication::Error::Hybrid(
820                        api::hybrid_codec::Reject::Expired,
821                    ));
822                }
823                Ok(())
824            },
825        );
826        assert!(
827            staged
828                .install_hosted(&repo, &trust, &authority, 1350)
829                .is_err()
830        );
831        assert_eq!(
832            calls.load(Ordering::SeqCst),
833            3,
834            "commit must recheck current disclosure after the file callback"
835        );
836        assert_eq!(
837            artifacts(repo.heddle_dir()),
838            before,
839            "late rejection must restore exact destination artifacts"
840        );
841        assert!(
842            trust
843                .snapshot()
844                .expect("rolled back trust")
845                .previous
846                .is_none()
847        );
848    }
849
850    #[test]
851    fn hosted_install_surfaces_typed_hybrid_rejection() {
852        for reason in [
853            api::hybrid_codec::Reject::StaleContext,
854            api::hybrid_codec::Reject::HighWater,
855        ] {
856            let scratch = tempfile::tempdir().expect("scratch");
857            let directory = tempfile::tempdir().expect("receiver");
858            let repo = Repository::init(directory.path()).expect("repository");
859            let before = artifacts(repo.heddle_dir());
860            let (staged, root, pinned) = source(scratch.path(), true);
861            let bundle = staged.import_authority().expect("history").clone();
862            let history = AcceptedHistory::from_selected_spool(
863                &bundle,
864                &pinned,
865                1350,
866                heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
867                    .expect("limits"),
868            )
869            .expect("selected history");
870            select_root(repo.heddle_dir(), &root).expect("root");
871            select_spool(
872                repo.heddle_dir(),
873                pinned.owner_genesis().spool_uuid(),
874                *history.genesis(),
875                *history.initial_owner(),
876            )
877            .expect("Spool");
878            let trust = HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock)
879                .expect("trust");
880            let authority = SelectedAuthority::new(
881                history,
882                bundle,
883                move |_: &ImportPublicProofBundleV1,
884                      _: i64,
885                      _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
886                    Err(repo::thread_replication::Error::Hybrid(reason))
887                },
888            );
889            let result = staged.install_hosted(&repo, &trust, &authority, 1350);
890            assert!(
891                matches!(result, Err(Error::Hybrid(actual)) if actual == reason),
892                "repository {reason:?} must surface as fetch::Error::Hybrid: {result:?}"
893            );
894            assert_eq!(
895                artifacts(repo.heddle_dir()),
896                before,
897                "typed rejection must not install artifacts"
898            );
899            assert!(
900                trust
901                    .snapshot()
902                    .expect("rolled back trust")
903                    .previous
904                    .is_none()
905            );
906        }
907    }
908
909    #[tokio::test]
910    async fn hosted_native_relay_retains_bundle_and_rechecks_revoked_durable_context() {
911        use std::sync::Arc;
912
913        use crate::replication::{
914            native::LocalReplica,
915            store::{ReceivedOperation, ReplicaStore},
916        };
917        let scratch = tempfile::tempdir().expect("scratch");
918        let directory = tempfile::tempdir().expect("receiver");
919        let repo = Repository::init(directory.path()).expect("repository");
920        let (staged, root, pinned) = source(scratch.path(), true);
921        let bundle = staged.import_authority().expect("history").clone();
922        let history = AcceptedHistory::from_selected_spool(
923            &bundle,
924            &pinned,
925            1350,
926            heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
927                .expect("limits"),
928        )
929        .expect("selected history");
930        select_root(repo.heddle_dir(), &root).expect("root");
931        select_spool(
932            repo.heddle_dir(),
933            pinned.owner_genesis().spool_uuid(),
934            *history.genesis(),
935            *history.initial_owner(),
936        )
937        .expect("Spool");
938        let trust = Arc::new(
939            HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock).expect("trust"),
940        );
941        let authority = Arc::new(SelectedAuthority::new(
942            history,
943            bundle.clone(),
944            |_: &ImportPublicProofBundleV1,
945             _: i64,
946             _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| Ok(()),
947        ));
948        staged
949            .install_hosted(&repo, &trust, authority.as_ref(), 1350)
950            .expect("genesis-only receiver");
951        let fixture: serde_json::Value =
952            serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
953                .expect("vectors");
954        let original = crate::replication::decode_record(record(&fixture, "converted_main"))
955            .expect("original conversion");
956        let operation = original.verify().expect("original signature");
957        let id = operation.id().expect("ID");
958        let replica =
959            ThreadReplica::open(repo.heddle_dir(), operation.thread).expect("selected replica");
960        let local = LocalReplica::new(replica, Arc::new(repo.store().clone()));
961        let relay = local.clone().with_hosted_authority(
962            repo.heddle_dir().to_path_buf(),
963            trust.clone(),
964            authority,
965        );
966        assert!(
967            relay
968                .receive(ReceivedOperation {
969                    native_authority: None,
970                    original: original.clone(),
971                    authority_admission: None,
972                    import_authority: None,
973                })
974                .await
975                .is_err(),
976            "receive cannot strip the required import history"
977        );
978        assert_eq!(
979            relay
980                .receive(ReceivedOperation {
981                    native_authority: None,
982                    original: original.clone(),
983                    authority_admission: None,
984                    import_authority: Some(Arc::new(bundle.clone()))
985                })
986                .await
987                .expect("independently verified native receive"),
988            heddle_object_model::object::thread_replication::Admission::Accepted
989        );
990        assert!(
991            local.operation(id).await.is_err(),
992            "an unconfigured relay must never strip retained HYBRID authority"
993        );
994        let (received, _) = relay
995            .operation(id)
996            .await
997            .expect("fresh relay admission")
998            .expect("original");
999        assert_eq!(received.original, original);
1000        assert_eq!(received.import_authority.as_deref(), Some(&bundle));
1001        let revoked = record(&fixture, "revoked_set");
1002        trust
1003            .mutate(&revoked, |_| Ok(()))
1004            .expect("independently persist N+1 revocation");
1005        assert!(
1006            matches!(
1007                relay.operation(id).await,
1008                Err(crate::replication::native::Error::Store(
1009                    repo::thread_replication::Error::Hybrid(api::hybrid_codec::Reject::HighWater)
1010                ))
1011            ),
1012            "export cannot revive retained N after durable N+1 revocation"
1013        );
1014    }
1015}