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},
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 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 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 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 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 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 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}