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