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