1use 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 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 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 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 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 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 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 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}