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