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