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            #[cfg(test)]
422            crash_point();
423        }
424    }
425    Ok(())
426}
427#[cfg(test)]
428thread_local! {
429    /// Test-only: abort the process, as a crash would, once this many staged
430    /// files have been published. No destructor, rollback or unlock runs.
431    static CRASH_AFTER_PUBLISHED: std::cell::Cell<Option<usize>> =
432        const { std::cell::Cell::new(None) };
433}
434#[cfg(test)]
435fn crash_point() {
436    CRASH_AFTER_PUBLISHED.with(|remaining| match remaining.get() {
437        Some(1) => {
438            eprintln!("hosted publish crash point reached");
439            std::process::abort();
440        }
441        Some(n) => remaining.set(Some(n - 1)),
442        None => {}
443    });
444}
445fn preparation(error: impl std::fmt::Display) -> Error {
446    Error::Preparation(error.to_string())
447}
448
449#[cfg(test)]
450pub(crate) mod tests {
451    use std::{
452        collections::BTreeMap,
453        sync::atomic::{AtomicUsize, Ordering},
454    };
455
456    use objects::{
457        object::{State, Tree},
458        store::{
459            ObjectStore,
460            pack::{ObjectType, PackBuilder, PackObjectId},
461        },
462    };
463    use repo::thread_replication::hosted_trust::*;
464
465    use super::*;
466    use crate::{
467        contract::*,
468        hybrid::authority::{
469            AcceptedHistory, SelectedAuthority,
470            tests::{bundle, selected},
471        },
472    };
473
474    pub(super) struct ReceiverClock;
475    impl Clock for ReceiverClock {
476        fn now_millis(&self) -> repo::thread_replication::Result<i64> {
477            Ok(1_350_000)
478        }
479        fn elapsed_millis(&self) -> repo::thread_replication::Result<u64> {
480            Ok(0)
481        }
482    }
483    pub(super) fn record<T: Message + Default>(fixture: &serde_json::Value, name: &str) -> T {
484        let vector = fixture["wire_vectors"]
485            .get(name)
486            .or_else(|| fixture["signed_vectors"].get(name))
487            .expect("published vector");
488        T::decode(
489            hex::decode(vector["wire_hex"].as_str().expect("wire bytes"))
490                .expect("hex")
491                .as_slice(),
492        )
493        .expect("original record")
494    }
495    pub(crate) fn source(
496        scratch: &Path,
497        seed_only: bool,
498    ) -> (
499        StagedSource,
500        RootSelection,
501        heddleco_capability_verifier::VerifiedCloneKeyring,
502    ) {
503        let fixture: serde_json::Value =
504            serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
505                .expect("release fixture");
506        let mut bundle = bundle();
507        bundle.history_proofs = [
508            "genesis_proof",
509            "genesis_dev_proof",
510            "publication_proof",
511            "dev_publication_proof",
512        ]
513        .map(|name| record(&fixture, name))
514        .to_vec();
515        let limits = heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
516            .expect("limits");
517        let pinned = selected(&bundle, limits);
518        let original: SignedRecord = record(&fixture, "converted_main");
519        let signed = crate::replication::decode_record(original).expect("converted original");
520        let operation = signed.verify().expect("original signature");
521        let genesis_record = bundle
522            .original_geneses
523            .iter()
524            .find(|record| {
525                crypto::thread_operation::SignedGenesis {
526                    canonical: record.canonical_record.clone(),
527                    signature: record.signatures[0].signature.clone(),
528                }
529                .verify()
530                .is_ok_and(|g| g.id().expect("genesis ID") == operation.thread)
531            })
532            .expect("selected genesis")
533            .clone();
534        let genesis = crypto::thread_operation::SignedGenesis {
535            canonical: genesis_record.canonical_record.clone(),
536            signature: genesis_record.signatures[0].signature.clone(),
537        }
538        .verify()
539        .expect("genesis signature");
540        let state: State = if seed_only {
541            objects::object::thread_replication::initial_base::synthetic_initial_base()
542                .expect("empty seed")
543        } else {
544            operation.source_state().expect("source").expect("State")
545        };
546        let tree = Tree::new();
547        assert_eq!(
548            tree.hash(),
549            state.tree,
550            "published native conversion has the empty source tree"
551        );
552        let mut builder = PackBuilder::for_repack(Default::default(), 0);
553        builder.add_id(
554            PackObjectId::StateId(state.id()),
555            ObjectType::State,
556            state.encode_current_msgpack().expect("State"),
557        );
558        builder.add_id(
559            PackObjectId::Hash(tree.hash()),
560            ObjectType::Tree,
561            tree.encode_canonical().expect("Tree"),
562        );
563        let (pack, index, _) = builder.build().expect("source pack");
564        let directory = tempfile::tempdir_in(scratch).expect("staging");
565        std::fs::write(directory.path().join("source.pack"), pack).expect("pack");
566        std::fs::write(directory.path().join("source.idx"), index).expect("index");
567        let spool = SpoolRef {
568            id: genesis.spool.clone(),
569        };
570        let owner = pinned.owner_state();
571        let history = AcceptedHistory::from_selected_spool(&bundle, &pinned, 1350, limits)
572            .expect("verified import owner history");
573        let authority = SelectedAuthority::new(
574            history,
575            bundle.clone(),
576            |_: &ImportPublicProofBundleV1, _: i64, _: &TrustTransaction<'_>| Ok(()),
577        );
578        let pin = api::import_authority::ImportWitnessRootPin {
579            authority: "https://weft.example.test".into(),
580            root_id: "descriptor-root-1".into(),
581            public_key: hex::decode(
582                fixture["keys"]["root"]["public_key_hex"]
583                    .as_str()
584                    .expect("root"),
585            )
586            .expect("hex"),
587            epoch: 1,
588        };
589        let carriers = repo::thread_replication::delegated_import::authenticate_import_carriers(
590            &bundle,
591            &authority,
592            &pin,
593            1_350_000,
594            &[],
595            &[],
596            |_| Ok(()),
597        )
598        .expect("independently authenticated import carriers");
599        let ready = TransferReady {
600            thread: Some(ThreadRef {
601                spool: Some(spool.clone()),
602                id: Some(ThreadId {
603                    value: operation.thread.as_bytes().to_vec(),
604                }),
605            }),
606            current: Some(RevisionRef {
607                spool: Some(spool),
608                revision: Some(revision_ref::Revision::State(
609                    api::heddle::api::common::StateId {
610                        value: state.id().as_bytes().to_vec(),
611                    },
612                )),
613            }),
614            thread_genesis: Some(ThreadGenesisRecord {
615                creator_authority: bundle
616                    .genesis_witnesses
617                    .iter()
618                    .find(|p| p.original_genesis.as_ref() == Some(&genesis_record))
619                    .expect("selected creator authority")
620                    .creator_authority_envelope
621                    .clone(),
622                genesis: Some(genesis_record),
623                ..Default::default()
624            }),
625            owner_genesis: bundle.owner_genesis.clone(),
626            ownership: Some(OwnerState {
627                owner: Some(PrincipalRef {
628                    id: uuid::Uuid::from_bytes(
629                        owner
630                            .signed_root()
631                            .root
632                            .as_ref()
633                            .expect("root")
634                            .account_uuid
635                            .as_slice()
636                            .try_into()
637                            .expect("UUID"),
638                    )
639                    .to_string(),
640                }),
641                root: Some(owner.signed_root().clone()),
642                accepted_transitions: pinned.wire().accepted_transitions.clone(),
643                version: owner.state_hash().to_vec(),
644                resource_keyring: Some(pinned.wire().clone()),
645                ..Default::default()
646            }),
647            full_closure_available: true,
648            import_authority: Some(bundle),
649            protocol: Some(crate::hybrid::protocol()),
650            ..Default::default()
651        };
652        let staged = super::super::staging::validate_with_receipts_and_carriers(
653            directory,
654            ready,
655            if seed_only { vec![] } else { vec![signed] },
656            vec![],
657            vec![],
658            Some(carriers),
659            Default::default(),
660        )
661        .expect("selected structural source");
662        let root = RootSelection {
663            authority: "https://weft.example.test".into(),
664            root_id: "descriptor-root-1".into(),
665            public_key: hex::decode(
666                fixture["keys"]["root"]["public_key_hex"]
667                    .as_str()
668                    .expect("root"),
669            )
670            .expect("hex")
671            .try_into()
672            .expect("key"),
673        };
674        (staged, root, pinned)
675    }
676    fn artifacts(path: &Path) -> BTreeMap<std::path::PathBuf, Vec<u8>> {
677        fn walk(root: &Path, path: &Path, values: &mut BTreeMap<std::path::PathBuf, Vec<u8>>) {
678            if !path.exists() {
679                return;
680            }
681            for entry in std::fs::read_dir(path).expect("directory") {
682                let entry = entry.expect("entry");
683                if entry.file_type().expect("type").is_dir() {
684                    walk(root, &entry.path(), values);
685                } else {
686                    values.insert(
687                        entry
688                            .path()
689                            .strip_prefix(root)
690                            .expect("relative")
691                            .to_path_buf(),
692                        std::fs::read(entry.path()).expect("bytes"),
693                    );
694                }
695            }
696        }
697        let mut values = BTreeMap::new();
698        for name in ["objects", "packs"] {
699            walk(path, &path.join(name), &mut values);
700        }
701        for name in ["owner-authorization.bin", "spool-id"] {
702            if let Ok(bytes) = std::fs::read(path.join(name)) {
703                values.insert(name.into(), bytes);
704            }
705        }
706        values
707    }
708    #[test]
709    fn selected_hosted_capture_and_genesis_only_install_without_sibling_branch() {
710        for seed_only in [false, true] {
711            let scratch = tempfile::tempdir().expect("scratch");
712            let directory = tempfile::tempdir().expect("receiver");
713            let repo = Repository::init(directory.path()).expect("repository");
714            let (staged, root, pinned) = source(scratch.path(), seed_only);
715            let state = staged.state().id();
716            let selected_thread = staged
717                .ready()
718                .thread
719                .as_ref()
720                .expect("Thread")
721                .id
722                .as_ref()
723                .expect("ID")
724                .value
725                .clone();
726            let bundle = staged.import_authority().expect("public history").clone();
727            let history = AcceptedHistory::from_selected_spool(
728                &bundle,
729                &pinned,
730                1350,
731                heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
732                    .expect("limits"),
733            )
734            .expect("selected history");
735            select_root(repo.heddle_dir(), &root).expect("root pin");
736            select_spool(
737                repo.heddle_dir(),
738                pinned.owner_genesis().spool_uuid(),
739                *history.genesis(),
740                *history.initial_owner(),
741            )
742            .expect("Spool selection");
743            let trust = HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock)
744                .expect("trust");
745            let authority = SelectedAuthority::new(
746                history,
747                bundle.clone(),
748                |_: &ImportPublicProofBundleV1,
749                 _: i64,
750                 _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
751                    Ok(())
752                },
753            );
754            assert_eq!(
755                staged
756                    .install_hosted(&repo, &trust, &authority, 1350)
757                    .expect("selected hosted install"),
758                state
759            );
760            assert!(repo.heddle_dir().join("spool-id").exists());
761            assert!(repo.heddle_dir().join("owner-authorization.bin").exists());
762            let (_, retained_owner) = repo
763                .pinned_owner_observation(1350)
764                .expect("independently retained owner observation");
765            assert_eq!(
766                retained_owner.owner_genesis().signed(),
767                pinned.owner_genesis().signed()
768            );
769            assert!(
770                repo.store()
771                    .get_state(&state)
772                    .expect("stored State")
773                    .is_some()
774            );
775            for record in &bundle.original_geneses {
776                let genesis = crypto::thread_operation::SignedGenesis {
777                    canonical: record.canonical_record.clone(),
778                    signature: record.signatures[0].signature.clone(),
779                }
780                .verify()
781                .expect("original genesis");
782                let id = genesis.id().expect("ID");
783                if id.as_bytes().as_slice() == selected_thread {
784                    assert_eq!(
785                        ThreadReplica::open(repo.heddle_dir(), id)
786                            .expect("selected replica")
787                            .hybrid_import_bundle()
788                            .expect("retained bundle"),
789                        Some(bundle.clone())
790                    );
791                } else {
792                    assert!(
793                        ThreadReplica::open(repo.heddle_dir(), id).is_err(),
794                        "public sibling proof must not install its native branch"
795                    );
796                }
797            }
798        }
799    }
800    #[test]
801    fn late_disclosure_failure_rolls_back_pack_owner_pin_spool_and_trust() {
802        let scratch = tempfile::tempdir().expect("scratch");
803        let directory = tempfile::tempdir().expect("receiver");
804        let repo = Repository::init(directory.path()).expect("repository");
805        let before = artifacts(repo.heddle_dir());
806        let (staged, root, pinned) = source(scratch.path(), false);
807        let bundle = staged.import_authority().expect("history").clone();
808        let history = AcceptedHistory::from_selected_spool(
809            &bundle,
810            &pinned,
811            1350,
812            heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
813                .expect("limits"),
814        )
815        .expect("selected history");
816        select_root(repo.heddle_dir(), &root).expect("root");
817        select_spool(
818            repo.heddle_dir(),
819            pinned.owner_genesis().spool_uuid(),
820            *history.genesis(),
821            *history.initial_owner(),
822        )
823        .expect("Spool");
824        let trust =
825            HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock).expect("trust");
826        let calls = AtomicUsize::new(0);
827        let authority = SelectedAuthority::new(
828            history,
829            bundle,
830            |_: &ImportPublicProofBundleV1,
831             _: i64,
832             _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
833                if calls.fetch_add(1, Ordering::SeqCst) > 1 {
834                    assert!(
835                        repo.heddle_dir().join("owner-authorization.bin").exists(),
836                        "failure must occur after staged artifacts were installed"
837                    );
838                    assert!(repo.heddle_dir().join("spool-id").exists());
839                    return Err(repo::thread_replication::Error::Hybrid(
840                        api::hybrid_codec::Reject::Expired,
841                    ));
842                }
843                Ok(())
844            },
845        );
846        assert!(
847            staged
848                .install_hosted(&repo, &trust, &authority, 1350)
849                .is_err()
850        );
851        assert_eq!(
852            calls.load(Ordering::SeqCst),
853            3,
854            "commit must recheck current disclosure after the file callback"
855        );
856        assert_eq!(
857            artifacts(repo.heddle_dir()),
858            before,
859            "late rejection must restore exact destination artifacts"
860        );
861        assert!(
862            trust
863                .snapshot()
864                .expect("rolled back trust")
865                .previous
866                .is_none()
867        );
868    }
869
870    #[test]
871    fn hosted_install_surfaces_typed_hybrid_rejection() {
872        for reason in [
873            api::hybrid_codec::Reject::StaleContext,
874            api::hybrid_codec::Reject::HighWater,
875        ] {
876            let scratch = tempfile::tempdir().expect("scratch");
877            let directory = tempfile::tempdir().expect("receiver");
878            let repo = Repository::init(directory.path()).expect("repository");
879            let before = artifacts(repo.heddle_dir());
880            let (staged, root, pinned) = source(scratch.path(), true);
881            let bundle = staged.import_authority().expect("history").clone();
882            let history = AcceptedHistory::from_selected_spool(
883                &bundle,
884                &pinned,
885                1350,
886                heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
887                    .expect("limits"),
888            )
889            .expect("selected history");
890            select_root(repo.heddle_dir(), &root).expect("root");
891            select_spool(
892                repo.heddle_dir(),
893                pinned.owner_genesis().spool_uuid(),
894                *history.genesis(),
895                *history.initial_owner(),
896            )
897            .expect("Spool");
898            let trust = HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock)
899                .expect("trust");
900            let authority = SelectedAuthority::new(
901                history,
902                bundle,
903                move |_: &ImportPublicProofBundleV1,
904                      _: i64,
905                      _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| {
906                    Err(repo::thread_replication::Error::Hybrid(reason))
907                },
908            );
909            let result = staged.install_hosted(&repo, &trust, &authority, 1350);
910            assert!(
911                matches!(result, Err(Error::Hybrid(actual)) if actual == reason),
912                "repository {reason:?} must surface as fetch::Error::Hybrid: {result:?}"
913            );
914            assert_eq!(
915                artifacts(repo.heddle_dir()),
916                before,
917                "typed rejection must not install artifacts"
918            );
919            assert!(
920                trust
921                    .snapshot()
922                    .expect("rolled back trust")
923                    .previous
924                    .is_none()
925            );
926        }
927    }
928
929    #[tokio::test]
930    async fn hosted_native_relay_retains_bundle_and_rechecks_revoked_durable_context() {
931        use std::sync::Arc;
932
933        use crate::replication::{
934            native::LocalReplica,
935            store::{ReceivedOperation, ReplicaStore},
936        };
937        let scratch = tempfile::tempdir().expect("scratch");
938        let directory = tempfile::tempdir().expect("receiver");
939        let repo = Repository::init(directory.path()).expect("repository");
940        let (staged, root, pinned) = source(scratch.path(), true);
941        let bundle = staged.import_authority().expect("history").clone();
942        let history = AcceptedHistory::from_selected_spool(
943            &bundle,
944            &pinned,
945            1350,
946            heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
947                .expect("limits"),
948        )
949        .expect("selected history");
950        select_root(repo.heddle_dir(), &root).expect("root");
951        select_spool(
952            repo.heddle_dir(),
953            pinned.owner_genesis().spool_uuid(),
954            *history.genesis(),
955            *history.initial_owner(),
956        )
957        .expect("Spool");
958        let trust = Arc::new(
959            HostedTrust::open(repo.heddle_dir(), &root.authority, ReceiverClock).expect("trust"),
960        );
961        let authority = Arc::new(SelectedAuthority::new(
962            history,
963            bundle.clone(),
964            |_: &ImportPublicProofBundleV1,
965             _: i64,
966             _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| Ok(()),
967        ));
968        staged
969            .install_hosted(&repo, &trust, authority.as_ref(), 1350)
970            .expect("genesis-only receiver");
971        let fixture: serde_json::Value =
972            serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
973                .expect("vectors");
974        let original = crate::replication::decode_record(record(&fixture, "converted_main"))
975            .expect("original conversion");
976        let operation = original.verify().expect("original signature");
977        let id = operation.id().expect("ID");
978        let replica =
979            ThreadReplica::open(repo.heddle_dir(), operation.thread).expect("selected replica");
980        let local = LocalReplica::new(replica, Arc::new(repo.store().clone()));
981        let relay = local.clone().with_hosted_authority(
982            repo.heddle_dir().to_path_buf(),
983            trust.clone(),
984            authority,
985        );
986        assert!(
987            relay
988                .receive(ReceivedOperation {
989                    native_authority: None,
990                    original: original.clone(),
991                    authority_admission: None,
992                    import_authority: None,
993                })
994                .await
995                .is_err(),
996            "receive cannot strip the required import history"
997        );
998        assert_eq!(
999            relay
1000                .receive(ReceivedOperation {
1001                    native_authority: None,
1002                    original: original.clone(),
1003                    authority_admission: None,
1004                    import_authority: Some(Arc::new(bundle.clone()))
1005                })
1006                .await
1007                .expect("independently verified native receive"),
1008            heddle_object_model::object::thread_replication::Admission::Accepted
1009        );
1010        assert!(
1011            local.operation(id).await.is_err(),
1012            "an unconfigured relay must never strip retained HYBRID authority"
1013        );
1014        let (received, _) = relay
1015            .operation(id)
1016            .await
1017            .expect("fresh relay admission")
1018            .expect("original");
1019        assert_eq!(received.original, original);
1020        assert_eq!(received.import_authority.as_deref(), Some(&bundle));
1021        let revoked = record(&fixture, "revoked_set");
1022        trust
1023            .mutate(&revoked, |_| Ok(()))
1024            .expect("independently persist N+1 revocation");
1025        assert!(
1026            matches!(
1027                relay.operation(id).await,
1028                Err(crate::replication::native::Error::Store(
1029                    repo::thread_replication::Error::Hybrid(api::hybrid_codec::Reject::HighWater)
1030                ))
1031            ),
1032            "export cannot revive retained N after durable N+1 revocation"
1033        );
1034    }
1035}