1use std::{
4 collections::{BTreeMap, BTreeSet},
5 path::Path,
6};
7
8use api::v2::client::MessageReader;
9use crypto::thread_operation::SignedOperation;
10use heddle_object_model::object::{
11 ContentHash, State, StateId, thread_replication::ThreadOperation,
12};
13use heddle_pack::store::pack::PackReader;
14use prost::Message;
15use tokio::io::AsyncWriteExt;
16
17pub(super) use super::ancestry::AncestryInput;
18use super::{Download, Error, Item, ancestry};
19use crate::{contract::*, transport};
20
21const METADATA_BYTES: usize = 16 * 1024 * 1024;
22const SOURCE_BYTES: u64 = 256 * 1024 * 1024;
23
24pub struct StagedSource {
27 #[cfg(feature = "native")]
29 pub(super) partial_trees: Option<heddle_pack::store::pack::VisibleSourceClosure>,
30 _scratch_lease: heddle_pack::store::pack::ScratchLease,
31 pub(super) directory: tempfile::TempDir,
32 pub(super) ready: TransferReady,
33 pub(super) operations: Vec<SignedOperation>,
34 pub(super) dependencies: Vec<ThreadGenesisRecord>,
35 pub(super) state: State,
36 #[cfg(feature = "native")]
37 pub(super) prefix_original: Option<SignedRecord>,
38 #[cfg(feature = "native")]
39 pub(super) authority_admissions:
40 BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
41 pub(super) ancestry: Vec<ancestry::VerifiedFloor>,
44}
45impl StagedSource {
46 #[cfg(feature = "native")]
47 pub fn select_prefix(&mut self, reference: &ForeignDependencyV1) -> Result<(), Error> {
48 let proof = match (self.import_authority(), self.native_authority()) {
49 (Some(b), None) => crate::hybrid::authority::PublicProof::from(b.clone()),
50 (None, Some(b)) => crate::hybrid::authority::PublicProof::from(b.clone()),
51 _ => return Err(Error::HostedTrustRequired),
52 };
53 self.prefix_original = match proof.prefix_original(reference) {
54 Ok(record) => Some(record.clone()),
55 Err(_) => self
56 .operations
57 .iter()
58 .map(|signed| {
59 let operation = signed.verify().map_err(preparation)?;
60 Ok(SignedRecord {
61 format: heddle_object_model::object::thread_replication::OPERATION_FORMAT
62 .into(),
63 canonical_record: signed.canonical.clone(),
64 signatures: vec![RecordSignature {
65 public_key: operation.publisher.to_vec(),
66 signature: signed.signature.clone(),
67 }],
68 })
69 })
70 .collect::<Result<Vec<_>, Error>>()?
71 .into_iter()
72 .find(|r| {
73 api::import_authority::signed_native_digest(r)
74 .is_ok_and(|d| d == reference.signed_native_digest)
75 }),
76 };
77 if self.prefix_original.is_none() {
78 return Err(api::hybrid_codec::Reject::Scope.into());
79 }
80 let projected = proof.prefix(reference)?;
81 let (authorities, landings, geneses, imported) = match &projected {
82 crate::hybrid::authority::PublicProof::Import(b) => (
83 &b.authority_witnesses,
84 &b.landing_witnesses,
85 b.genesis_witnesses
86 .iter()
87 .filter_map(|p| p.original_genesis.as_ref())
88 .collect::<Vec<_>>(),
89 Some(b.as_ref()),
90 ),
91 crate::hybrid::authority::PublicProof::Native(b) => (
92 &b.authority_witnesses,
93 &b.landing_witnesses,
94 b.genesis_witnesses
95 .iter()
96 .filter_map(|p| p.original_genesis.as_ref())
97 .collect::<Vec<_>>(),
98 None,
99 ),
100 };
101 let retained = authorities
102 .iter()
103 .flat_map(|p| p.original.iter().chain(&p.dependencies))
104 .chain(landings.iter().flat_map(|p| {
105 p.execution
106 .iter()
107 .chain(p.source_operation.iter())
108 .chain(&p.review_evidence)
109 }))
110 .chain(self.prefix_original.iter())
111 .collect::<Vec<_>>();
112 let mut operations = Vec::new();
113 for signed in &self.operations {
114 let operation = signed.verify().map_err(preparation)?;
115 let id = operation.id().map_err(preparation)?;
116 let published = if let Some(imported) = imported {
117 let frontier = api::import_authority::frontier_digest(&ImportFrontierV1 {
118 format_version: 1,
119 thread_id: operation.thread.as_bytes().to_vec(),
120 operation_ids: vec![id.as_bytes().to_vec()],
121 })?;
122 imported.operations.iter().any(|o| {
123 o.body
124 .as_ref()
125 .is_some_and(|b| b.resulting_frontier_digest == frontier)
126 })
127 } else {
128 false
129 };
130 if published
131 || retained.iter().any(|r| {
132 r.canonical_record == signed.canonical
133 && r.signatures.iter().any(|s| {
134 s.public_key.as_slice() == operation.publisher
135 && s.signature == signed.signature
136 })
137 })
138 {
139 operations.push(signed.clone());
140 }
141 }
142 let available = self
145 .operations
146 .iter()
147 .map(|signed| {
148 let operation = signed.verify().map_err(preparation)?;
149 Ok((operation.id().map_err(preparation)?, (signed, operation)))
150 })
151 .collect::<Result<BTreeMap<_, _>, Error>>()?;
152 let mut pending = operations
153 .iter()
154 .map(|signed| {
155 signed
156 .verify()
157 .map_err(preparation)?
158 .id()
159 .map_err(preparation)
160 })
161 .collect::<Result<Vec<_>, Error>>()?;
162 let mut causal = BTreeSet::new();
163 while let Some(id) = pending.pop() {
164 if causal.insert(id) {
165 let (signed, op) = available
166 .get(&id)
167 .ok_or(Error::Invalid("prefix ancestor absent"))?;
168 if !foreign_endpoint(projected.foreign_dependencies(), op, signed)? {
169 pending.extend(op.parents.iter().copied());
170 }
171 }
172 }
173 self.operations = available
174 .into_iter()
175 .filter(|(id, _)| causal.contains(id))
176 .map(|(_, (signed, _))| signed.clone())
177 .collect();
178 let mut threads = BTreeSet::new();
179 for record in geneses {
180 threads.insert(
181 crypto::import_authority::verify_native_genesis(record)
182 .map_err(preparation)?
183 .1
184 .id()
185 .map_err(preparation)?,
186 );
187 }
188 for signed in &self.operations {
189 threads.insert(signed.verify().map_err(preparation)?.thread);
190 }
191 for r in projected.foreign_dependencies() {
192 threads.insert(ContentHash::from_bytes(
193 r.thread_genesis_digest
194 .as_slice()
195 .try_into()
196 .map_err(|_| Error::Hybrid(api::hybrid_codec::Reject::Scope))?,
197 ));
198 }
199 self.dependencies.retain(|g| {
200 g.genesis.as_ref().is_some_and(|r| {
201 crypto::import_authority::verify_native_genesis(r)
202 .is_ok_and(|(_, g)| g.id().is_ok_and(|id| threads.contains(&id)))
203 })
204 });
205 for wrapper in self
206 .ready
207 .thread_genesis
208 .iter_mut()
209 .chain(&mut self.dependencies)
210 {
211 wrapper.ownership_claims.retain(|r| retained.contains(&r));
212 wrapper
213 .ownership_resolutions
214 .retain(|r| retained.contains(&r));
215 }
216 match projected {
217 crate::hybrid::authority::PublicProof::Import(b) => {
218 self.ready.import_authority = Some(*b)
219 }
220 crate::hybrid::authority::PublicProof::Native(b) => {
221 self.ready.native_authority = Some(*b)
222 }
223 }
224 Ok(())
225 }
226 pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
227 [
228 self.directory.path().join("source.pack"),
229 self.directory.path().join("source.idx"),
230 ]
231 }
232 pub fn operations(&self) -> &[SignedOperation] {
233 &self.operations
234 }
235 pub fn dependency_geneses(&self) -> &[ThreadGenesisRecord] {
236 &self.dependencies
237 }
238 pub fn ready(&self) -> &TransferReady {
239 &self.ready
240 }
241 pub fn import_authority(&self) -> Option<&crate::contract::ImportPublicProofBundleV1> {
243 self.ready.import_authority.as_ref()
244 }
245 pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
246 self.ready.native_authority.as_ref()
247 }
248 pub fn refresh_native_authority(
249 &mut self,
250 refreshed: NativePublicProofBundleV1,
251 ) -> Result<(), Error> {
252 let original = self
253 .ready
254 .native_authority
255 .as_mut()
256 .ok_or(Error::HostedTrustRequired)?;
257 crate::hybrid::history::replace_native_receiver_metadata(original, refreshed)?;
258 Ok(())
259 }
260 pub fn refresh_import_authority(
263 &mut self,
264 refreshed: crate::contract::ImportPublicProofBundleV1,
265 ) -> Result<(), Error> {
266 let original = self
267 .ready
268 .import_authority
269 .as_mut()
270 .ok_or(Error::HostedTrustRequired)?;
271 crate::hybrid::history::replace_receiver_metadata(original, refreshed)?;
272 Ok(())
273 }
274 pub fn state(&self) -> &State {
275 &self.state
276 }
277 pub fn is_complete(&self) -> bool {
279 self.ready.full_closure_available
280 }
281 pub fn verified_import_floors(&self) -> impl Iterator<Item = (StateId, &BTreeSet<StateId>)> {
285 self.ancestry
286 .iter()
287 .filter(|floor| floor.coverage == import_ancestry_page::Coverage::Floor)
288 .map(|floor| (floor.tip, &floor.members))
289 }
290 #[cfg(feature = "native")]
291 pub(super) fn ancestry_paths(&self) -> Option<[std::path::PathBuf; 2]> {
292 let pack = self.directory.path().join("ancestry.pack");
293 pack.exists()
294 .then(|| [pack, self.directory.path().join("ancestry.idx")])
295 }
296}
297impl<R: MessageReader<Error = transport::Error>> Download<R> {
298 pub async fn stage(self, scratch: &Path) -> Result<StagedSource, Error> {
301 self.stage_inner(scratch, None).await
302 }
303 pub async fn stage_with_import_carriers(
306 self,
307 scratch: &Path,
308 carriers: crypto::import_authority::VerifiedImportCarriers,
309 ) -> Result<StagedSource, Error> {
310 self.stage_inner(scratch, Some(carriers)).await
311 }
312 async fn stage_inner(
313 mut self,
314 scratch: &Path,
315 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
316 ) -> Result<StagedSource, Error> {
317 if let Some(carriers) = &carriers {
318 let mut original = self
319 .state
320 .ready
321 .import_authority
322 .clone()
323 .ok_or(Error::HostedTrustRequired)?;
324 crate::hybrid::history::replace_receiver_metadata(
325 &mut original,
326 carriers.bundle().clone(),
327 )
328 .map_err(preparation)?;
329 }
330 if self.state.facets != [SharedFacet::Source as i32] {
331 return Err(Error::Invalid("staging requires the source facet alone"));
332 }
333 let total = self
334 .state
335 .ready
336 .packs
337 .iter()
338 .try_fold(0u64, |sum, extent| sum.checked_add(extent.length))
339 .ok_or(Error::Invalid("source artifact length overflow"))?;
340 if total > SOURCE_BYTES {
341 return Err(Error::Invalid("staged source exceeds 256 MiB"));
342 }
343 self.state.limits.max_operations = self.state.limits.max_operations.min(10_000);
344 let directory = tempfile::Builder::new()
345 .prefix("thread-download-")
346 .tempdir_in(scratch)?;
347 let _scratch_lease = heddle_pack::store::pack::ScratchLease::acquire(directory.path())?;
348 let mut files = [
349 tokio::fs::File::create(directory.path().join("source.pack")).await?,
350 tokio::fs::File::create(directory.path().join("source.idx")).await?,
351 ];
352 let mut operations = Vec::new();
353 let mut receipt_records = Vec::new();
354 let mut dependencies = Vec::new();
355 let mut ancestry = AncestryInput::new(self.state.excluded_tips.clone());
356 let mut metadata_bytes = 0usize;
357 let mut complete = false;
358 while let Some(item) = self.next().await? {
359 match item {
360 Item::Pack(chunk) => {
361 let kind = chunk
362 .extent
363 .as_ref()
364 .ok_or(Error::Invalid("chunk extent absent"))?
365 .kind;
366 let index = match pack_extent::Kind::try_from(kind) {
367 Ok(pack_extent::Kind::NativePack) => 0,
368 Ok(pack_extent::Kind::NativeIndex) => 1,
369 _ => return Err(Error::Invalid("native source artifacts required")),
370 };
371 files[index].write_all(&chunk.data).await?;
372 }
373 Item::Operations(batch) => {
374 metadata_bytes = metadata_bytes
375 .checked_add(batch.encoded_len())
376 .ok_or(Error::Invalid("source metadata length overflow"))?;
377 if metadata_bytes > METADATA_BYTES {
378 return Err(Error::Invalid("staged source metadata exceeds 16 MiB"));
379 }
380 for received in crate::authority_admission::match_batch(&batch)? {
381 operations.push(received.original);
382 receipt_records.extend(received.authority_admission);
383 }
384 }
385 Item::ThreadGenesis(record) => {
386 metadata_bytes = metadata_bytes
387 .checked_add(record.encoded_len())
388 .ok_or(Error::Invalid("source metadata length overflow"))?;
389 if metadata_bytes > METADATA_BYTES || dependencies.len() >= 127 {
390 return Err(Error::Invalid("dependency metadata exceeds bounds"));
391 }
392 dependencies.push(record);
393 }
394 Item::ImportAncestry(page) => ancestry.push(
396 page,
397 self.state
398 .ready
399 .thread
400 .as_ref()
401 .ok_or(Error::Invalid("Thread absent"))?,
402 directory.path(),
403 )?,
404 Item::Complete(_) => complete = true,
405 Item::Sidecar(_) => return Err(Error::Invalid("source staging excludes sidecars")),
406 }
407 }
408 if !complete {
409 return Err(Error::Invalid("source staging requires Complete"));
410 }
411 for file in &mut files {
412 file.flush().await?;
413 file.sync_all().await?;
414 }
415 drop(files);
416 ancestry.finish()?;
417 let mut ready = self.state.ready;
418 if let Some(carriers) = &carriers {
419 ready.import_authority = Some(carriers.bundle().clone());
420 }
421 tokio::task::spawn_blocking(move || {
422 validate_with_receipts_and_carriers(
423 directory,
424 ready,
425 operations,
426 dependencies,
427 receipt_records,
428 carriers,
429 ancestry,
430 )
431 })
432 .await
433 .map_err(|error| Error::Preparation(error.to_string()))?
434 }
435}
436#[cfg(all(test, feature = "native"))]
437fn validate(
438 directory: tempfile::TempDir,
439 ready: TransferReady,
440 operations: Vec<SignedOperation>,
441 dependencies: Vec<ThreadGenesisRecord>,
442) -> Result<StagedSource, Error> {
443 validate_with_receipts(directory, ready, operations, dependencies, Vec::new())
444}
445
446struct DisclosureInput {
447 directory: tempfile::TempDir,
448 operations: Vec<SignedOperation>,
449 dependency_records: Vec<ThreadGenesisRecord>,
450 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
451 allow_partial: bool,
452 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
453 foreign: Vec<ForeignDependencyV1>,
454 ancestry: AncestryInput,
455 require_import_ancestry: bool,
456}
457
458#[cfg(all(test, feature = "native"))]
459pub(crate) fn validate_with_receipts(
460 directory: tempfile::TempDir,
461 ready: TransferReady,
462 operations: Vec<SignedOperation>,
463 dependencies: Vec<ThreadGenesisRecord>,
464 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
465) -> Result<StagedSource, Error> {
466 validate_with_receipts_and_carriers(
467 directory,
468 ready,
469 operations,
470 dependencies,
471 receipt_records,
472 None,
473 AncestryInput::default(),
474 )
475}
476pub(super) fn validate_with_receipts_and_carriers(
477 directory: tempfile::TempDir,
478 mut ready: TransferReady,
479 operations: Vec<SignedOperation>,
480 dependencies: Vec<ThreadGenesisRecord>,
481 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
482 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
483 ancestry: AncestryInput,
484) -> Result<StagedSource, Error> {
485 if carriers
486 .as_ref()
487 .is_some_and(|c| ready.import_authority.as_ref() != Some(c.bundle()))
488 {
489 return Err(Error::HostedTrustRequired);
490 }
491 let value = validate_disclosure_artifacts(
492 ready
493 .thread
494 .as_ref()
495 .ok_or(Error::Invalid("Thread absent"))?,
496 ready
497 .current
498 .as_ref()
499 .ok_or(Error::Invalid("revision absent"))?,
500 ready
501 .thread_genesis
502 .as_ref()
503 .ok_or(Error::Invalid("original genesis absent"))?,
504 DisclosureInput {
505 directory,
506 operations,
507 dependency_records: dependencies,
508 receipt_records,
509 allow_partial: !ready.full_closure_available,
510 foreign: ready
511 .import_authority
512 .as_ref()
513 .map(|b| b.foreign_dependencies.clone())
514 .or_else(|| {
515 ready
516 .native_authority
517 .as_ref()
518 .map(|b| b.foreign_dependencies.clone())
519 })
520 .unwrap_or_default(),
521 carriers,
522 ancestry,
523 require_import_ancestry: true,
524 },
525 )?;
526 ready.thread_genesis = Some(value.genesis.clone());
528 Ok(StagedSource {
529 _scratch_lease: value._scratch_lease,
530 directory: value.directory,
531 ready,
532 #[cfg(feature = "native")]
533 prefix_original: None,
534 operations: value.operations,
535 dependencies: value.dependencies,
536 state: value.state,
537 #[cfg(feature = "native")]
538 partial_trees: value.partial_trees,
539 #[cfg(feature = "native")]
540 authority_admissions: value.authority_admissions,
541 ancestry: value.ancestry,
542 })
543}
544pub struct ValidatedSourceArtifacts {
547 pub(crate) native_authority: Option<NativePublicProofBundleV1>,
548 pub(crate) import_authority: Option<ImportPublicProofBundleV1>,
549 #[cfg(feature = "native")]
550 partial_trees: Option<heddle_pack::store::pack::VisibleSourceClosure>,
551 _scratch_lease: heddle_pack::store::pack::ScratchLease,
552 directory: tempfile::TempDir,
553 operations: Vec<SignedOperation>,
554 genesis: ThreadGenesisRecord,
555 dependencies: Vec<ThreadGenesisRecord>,
556 state: State,
557 authority_admissions:
558 BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
559 ancestry: Vec<ancestry::VerifiedFloor>,
560}
561impl ValidatedSourceArtifacts {
562 pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
563 self.native_authority.as_ref()
564 }
565 pub fn import_authority(&self) -> Option<&ImportPublicProofBundleV1> {
566 self.import_authority.as_ref()
567 }
568
569 pub fn into_hosted_source(self, mut ready: TransferReady) -> Result<StagedSource, Error> {
572 if ready.import_authority != self.import_authority
573 || ready.native_authority != self.native_authority
574 || (self.import_authority.is_none() && self.native_authority.is_none())
575 {
576 return Err(Error::HostedTrustRequired);
577 }
578 let reference = ready
579 .thread
580 .as_ref()
581 .ok_or(Error::Invalid("Thread absent"))?;
582 super::verify_origin(&self.genesis, reference)?;
583 let revision = ready
584 .current
585 .as_ref()
586 .ok_or(Error::Invalid("revision absent"))?;
587 if revision.spool != reference.spool
588 || revision.revision
589 != Some(revision_ref::Revision::State(
590 api::heddle::api::common::StateId {
591 value: self.state.id().as_bytes().to_vec(),
592 },
593 ))
594 || !ready.full_closure_available
595 {
596 return Err(Error::Invalid(
597 "publication source differs from hosted install selection",
598 ));
599 }
600 ready.thread_genesis = Some(self.genesis);
601 Ok(StagedSource {
602 _scratch_lease: self._scratch_lease,
603 directory: self.directory,
604 ready,
605 #[cfg(feature = "native")]
606 prefix_original: None,
607 operations: self.operations,
608 dependencies: self.dependencies,
609 state: self.state,
610 #[cfg(feature = "native")]
611 partial_trees: self.partial_trees,
612 #[cfg(feature = "native")]
613 authority_admissions: self.authority_admissions,
614 ancestry: self.ancestry,
615 })
616 }
617
618 pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
619 [
620 self.directory.path().join("source.pack"),
621 self.directory.path().join("source.idx"),
622 ]
623 }
624 pub fn operations(&self) -> &[SignedOperation] {
625 &self.operations
626 }
627 pub fn geneses(&self) -> impl Iterator<Item = &ThreadGenesisRecord> {
628 std::iter::once(&self.genesis).chain(&self.dependencies)
629 }
630 pub fn state(&self) -> &State {
631 &self.state
632 }
633 pub fn authority_admissions(
634 &self,
635 ) -> &BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission> {
636 &self.authority_admissions
637 }
638}
639#[allow(clippy::too_many_arguments)]
640pub(crate) fn validate_artifacts(
641 directory: tempfile::TempDir,
642 thread: &ThreadRef,
643 revision: &RevisionRef,
644 original: &ThreadGenesisRecord,
645 operations: Vec<SignedOperation>,
646 dependency_records: Vec<ThreadGenesisRecord>,
647 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
648 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
649 foreign: Vec<ForeignDependencyV1>,
650) -> Result<ValidatedSourceArtifacts, Error> {
651 validate_disclosure_artifacts(
652 thread,
653 revision,
654 original,
655 DisclosureInput {
656 directory,
657 operations,
658 dependency_records,
659 receipt_records,
660 allow_partial: false,
661 foreign,
662 carriers,
663 ancestry: AncestryInput::default(),
664 require_import_ancestry: false,
667 },
668 )
669}
670
671fn validate_disclosure_artifacts(
672 thread: &ThreadRef,
673 revision: &RevisionRef,
674 original: &ThreadGenesisRecord,
675 input: DisclosureInput,
676) -> Result<ValidatedSourceArtifacts, Error> {
677 let DisclosureInput {
678 directory,
679 operations,
680 dependency_records,
681 receipt_records,
682 allow_partial,
683 carriers,
684 foreign,
685 ancestry,
686 require_import_ancestry,
687 } = input;
688 let scratch_lease = heddle_pack::store::pack::ScratchLease::acquire(directory.path())?;
689 if !ancestry.is_empty() && carriers.is_none() {
690 return Err(Error::Invalid(
691 "import ancestry requires independently authenticated import carriers",
692 ));
693 }
694 if operations.len() > 10_000
695 || dependency_records.len() >= 128
696 || receipt_records.len() > operations.len()
697 {
698 return Err(Error::Invalid("source original graph exceeds bounds"));
699 }
700 let mut metadata = original.encoded_len();
701 for record in &dependency_records {
702 metadata = metadata.saturating_add(record.encoded_len());
703 }
704 for operation in &operations {
705 metadata = metadata.saturating_add(operation.canonical.len() + operation.signature.len());
706 }
707 for receipt in &receipt_records {
708 metadata = metadata.saturating_add(receipt.canonical.len() + receipt.signature.len());
709 }
710 let mut evidence_ids = BTreeSet::new();
711 for wrapper in std::iter::once(original).chain(&dependency_records) {
712 for record in &wrapper.boundary_acceptances {
713 evidence_ids.insert(heddle_object_model::object::ContentHash::compute_typed(
714 heddle_object_model::object::original_boundary_acceptance::FORMAT,
715 &record.canonical_record,
716 ));
717 if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
718 return Err(Error::Invalid("boundary evidence count exceeded"));
719 }
720 }
721 }
722 for receipt in &receipt_records {
723 if let Some(evidence) = &receipt.boundary_acceptance {
724 if evidence_ids.insert(
725 evidence
726 .verify_signature()
727 .map_err(preparation)?
728 .id()
729 .map_err(preparation)?,
730 ) {
731 metadata =
732 metadata.saturating_add(evidence.canonical.len() + evidence.signature.len());
733 }
734 if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
735 return Err(Error::Invalid("boundary evidence count exceeded"));
736 }
737 }
738 }
739 if metadata > METADATA_BYTES {
740 return Err(Error::Invalid("source metadata exceeds 16 MiB"));
741 }
742 if revision.spool != thread.spool {
743 return Err(Error::Invalid("source revision crosses Spool"));
744 }
745 let genesis = super::verify_origin(original, thread)?;
746 let Some(revision_ref::Revision::State(selected)) = revision.revision.as_ref() else {
747 return Err(Error::Invalid("exact native State required"));
748 };
749 let selected_thread = genesis.id().map_err(preparation)?;
750 if operations.is_empty() {
751 if !dependency_records.is_empty() || !receipt_records.is_empty() {
752 return Err(Error::Invalid(
753 "initial source cannot carry dependency originals",
754 ));
755 }
756 let state =
757 heddle_object_model::object::thread_replication::initial_base::synthetic_initial_base()
758 .map_err(preparation)?;
759 let canonical = state.encode_current_msgpack().map_err(preparation)?;
760 heddle_object_model::object::thread_replication::initial_base::initial_base_state(
761 &genesis, &canonical,
762 )
763 .map_err(preparation)?;
764 if selected.value.as_slice() != state.id().as_bytes() {
765 return Err(Error::Invalid(
766 "selected initial source differs from canonical seed",
767 ));
768 }
769 PackReader::open(
770 &directory.path().join("source.pack"),
771 &directory.path().join("source.idx"),
772 directory.path(),
773 )
774 .map_err(preparation)?
775 .validate_source_closure_with_metadata(&state, &[], None, SOURCE_BYTES)
776 .map_err(preparation)?;
777 if !ancestry.is_empty() {
778 return Err(Error::Invalid(
779 "initial source cannot carry import ancestry",
780 ));
781 }
782 return Ok(ValidatedSourceArtifacts {
783 import_authority: None,
784 native_authority: None,
785 _scratch_lease: scratch_lease,
786 directory,
787 operations,
788 genesis: original.clone(),
789 dependencies: Vec::new(),
790 state,
791 #[cfg(feature = "native")]
792 partial_trees: None,
793 authority_admissions: BTreeMap::new(),
794 ancestry: Vec::new(),
795 });
796 }
797 let mut geneses = BTreeMap::from([(selected_thread, genesis)]);
798 let mut dependencies = Vec::new();
799 for wrapper in dependency_records {
800 let record = wrapper
801 .genesis
802 .as_ref()
803 .ok_or(Error::Invalid("dependency signed genesis absent"))?;
804 let candidate = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
805 &record.canonical_record,
806 )
807 .map_err(preparation)?;
808 let reference = ThreadRef {
809 spool: thread.spool.clone(),
810 id: Some(ThreadId {
811 value: candidate.id().map_err(preparation)?.as_bytes().to_vec(),
812 }),
813 };
814 let candidate = super::verify_origin(&wrapper, &reference)?;
815 let id = candidate.id().map_err(preparation)?;
816 if geneses.len() >= 128 || geneses.insert(id, candidate).is_some() {
817 return Err(Error::Invalid(
818 "duplicate or oversized dependency genesis set",
819 ));
820 }
821 dependencies.push(wrapper);
822 }
823 let mut claim_frontiers = BTreeMap::new();
824 for wrapper in std::iter::once(original).chain(&dependencies) {
825 let signed = wrapper
826 .genesis
827 .as_ref()
828 .ok_or(Error::Invalid("claim genesis absent"))?;
829 let genesis = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
830 &signed.canonical_record,
831 )
832 .map_err(preparation)?;
833 let mut frontier = BTreeSet::new();
834 let claims = crate::replication::ownership::verify_claims(wrapper, &genesis)?;
835 let resolutions = crate::replication::ownership::verify_resolutions(wrapper, &genesis)?;
836 if claims.is_empty() && resolutions.is_empty() {
837 continue;
838 }
839 for claim in claims {
840 frontier.extend(
841 claim
842 .original
843 .verify()
844 .map_err(preparation)?
845 .source_frontier,
846 );
847 }
848 for resolution in resolutions {
849 frontier.extend(
850 heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&resolution.original.canonical)
851 .map_err(preparation)?.frontier,
852 );
853 }
854 claim_frontiers.insert(genesis.id().map_err(preparation)?, frontier);
855 }
856 let mut originals = BTreeMap::new();
857 let mut decoded = BTreeMap::<ContentHash, ThreadOperation>::new();
858 let mut selected_operation = None;
859 let mut source_thread = selected_thread;
860 let mut inherited_bases = BTreeSet::new();
861 for _ in 0..128 {
862 let current = geneses
863 .get(&source_thread)
864 .ok_or(Error::Invalid("fork base source genesis absent"))?;
865 if selected.value.as_slice() != current.base.as_bytes() {
866 break;
867 }
868 let parent = current.parent.ok_or(Error::Invalid(
869 "non-system base has no original parent source",
870 ))?;
871 if !inherited_bases.insert(source_thread) {
872 return Err(Error::Invalid("fork base parent cycle"));
873 }
874 let ancestor = geneses
875 .get(&parent)
876 .ok_or(Error::Invalid("fork base parent original absent"))?;
877 if ancestor.spool != current.spool {
878 return Err(Error::Invalid("fork base crosses Spool"));
879 }
880 source_thread = parent;
881 }
882 if selected.value.as_slice()
883 == geneses
884 .get(&source_thread)
885 .ok_or(Error::Invalid("fork base source genesis absent"))?
886 .base
887 .as_bytes()
888 {
889 return Err(Error::Invalid("fork base source chain exceeds bound"));
890 }
891 for signed in &operations {
892 let operation = signed.verify().map_err(preparation)?;
893 let id = operation.id().map_err(preparation)?;
894 let state = operation
895 .source_state()
896 .map_err(preparation)?
897 .ok_or(Error::Invalid("non-source operation in source ancestry"))?;
898 if operation.thread == source_thread
899 && state.id().as_bytes().as_slice() == selected.value
900 && selected_operation.replace((id, state)).is_some()
901 {
902 return Err(Error::Invalid("ambiguous selected source proof"));
903 }
904 originals.insert(id, signed.clone());
905 if decoded.insert(id, operation).is_some() {
906 return Err(Error::Invalid("duplicate source proof"));
907 }
908 }
909 let mut authority_admissions = BTreeMap::new();
910 for receipt in receipt_records {
911 let statement = receipt.verify_signature().map_err(preparation)?;
912 let operation_id = statement.subject.operation_id().ok_or(Error::Invalid(
913 "source batch cannot carry ownership claim admission",
914 ))?;
915 let original = originals
916 .get(&operation_id)
917 .ok_or(Error::Invalid("unmatched source authority receipt"))?;
918 receipt.verify(original, &heddle_object_model::object::thread_replication::integration::TrustedHostedExecutor {
921 spool: statement.spool, spool_genesis: statement.spool_genesis, executor: statement.executor,
922 }).map_err(preparation)?;
923 if authority_admissions.insert(operation_id, receipt).is_some() {
924 return Err(Error::Invalid("duplicate source authority receipt"));
925 }
926 }
927 let (selected_id, state, selected_in_floor) = match selected_operation {
931 Some((id, state)) => (id, state, false),
932 None => {
933 let selected_state = StateId::from_bytes(
934 selected
935 .value
936 .as_slice()
937 .try_into()
938 .map_err(|_| Error::Invalid("selected revision identity width"))?,
939 );
940 let (tip, canonical) = ancestry
941 .selected_page_tip(selected_state, directory.path())?
942 .ok_or(Error::Invalid("selected source proof absent"))?;
943 let owner = decoded
944 .iter()
945 .find(|(_, operation)| {
946 operation.thread == source_thread
947 && operation
948 .source_state()
949 .ok()
950 .flatten()
951 .is_some_and(|state| state.id() == tip)
952 })
953 .map(|(id, _)| *id)
954 .ok_or(Error::Invalid(
955 "import ancestry tip is not a carried source operation",
956 ))?;
957 let state = State::decode_current_msgpack(&canonical)
958 .map_err(|_| Error::Invalid("selected import ancestor is not canonical"))?;
959 if state.id() != selected_state {
960 return Err(Error::Invalid(
961 "selected import ancestor differs from its address",
962 ));
963 }
964 (owner, state, true)
965 }
966 };
967 let mut import_floors = Vec::new();
968 let mut pending = BTreeSet::from([selected_id]);
969 let mut seen = BTreeSet::new();
970 let mut used_threads = BTreeSet::new();
971 let mut claim_threads = BTreeSet::new();
972 let mut foreign_endpoints = BTreeSet::new();
973 let mut edges = BTreeMap::new();
974 while let Some(id) = pending.pop_first() {
975 if !seen.insert(id) {
976 continue;
977 }
978 let operation = decoded
979 .get(&id)
980 .ok_or(Error::Invalid("incomplete source ancestry"))?;
981 let signed = originals
982 .get(&id)
983 .ok_or(Error::Invalid("original absent"))?;
984 if foreign_endpoint(&foreign, operation, signed)? {
985 used_threads.insert(operation.thread);
990 foreign_endpoints.insert(id);
991 edges.insert(id, BTreeSet::new());
992 continue;
993 }
994 let parents = operation
995 .parents
996 .iter()
997 .map(|id| {
998 decoded
999 .get(id)
1000 .cloned()
1001 .ok_or(Error::Invalid("incomplete source ancestry"))
1002 })
1003 .collect::<Result<Vec<_>, _>>()?;
1004 let genesis = geneses
1005 .get(&operation.thread)
1006 .ok_or(Error::Invalid("source dependency genesis absent"))?;
1007 used_threads.insert(operation.thread);
1008 if claim_threads.insert(operation.thread)
1009 && let Some(frontier) = claim_frontiers.get(&operation.thread)
1010 {
1011 for head in frontier {
1012 if decoded
1013 .get(head)
1014 .is_none_or(|source| source.thread != operation.thread)
1015 {
1016 return Err(Error::Invalid(
1017 "ownership claim cutoff source proof absent or foreign",
1018 ));
1019 }
1020 }
1021 pending.extend(frontier);
1022 }
1023 let imported = carriers
1024 .as_ref()
1025 .map(|c| c.bind(genesis, operation, &parents))
1026 .transpose()
1027 .map_err(preparation)?
1028 .flatten();
1029 match &imported {
1030 Some(bound) => bound.validate_parents(genesis, &parents),
1031 None => operation.validate_parents(genesis, &parents),
1032 }
1033 .map_err(preparation)?;
1034 if let Some(bound) = imported {
1035 let mut frontier = BTreeSet::new();
1039 for parent in &parents {
1040 if let Some(state) = parent.source_state().map_err(preparation)? {
1041 frontier.insert(state.id());
1042 }
1043 }
1044 import_floors.push(ancestry::ImportFloorInput {
1045 digest: api::import_authority::signed_operation_digest(bound.signed())
1046 .map_err(preparation)?,
1047 tip: operation
1048 .source_state()
1049 .map_err(preparation)?
1050 .ok_or(Error::Invalid("import operation has no source State"))?,
1051 frontier,
1052 });
1053 }
1054
1055 let mut required = operation.parents.clone();
1056 if let Some(receipt) = operation.local_integration().map_err(preparation)? {
1057 let source = decoded
1058 .get(&receipt.source_operation)
1059 .ok_or(Error::Invalid(
1060 "local integration original source proof absent",
1061 ))?;
1062 receipt.validate_source(source).map_err(preparation)?;
1063 required.insert(receipt.source_operation);
1064 pending.insert(receipt.source_operation);
1065 }
1066 if let Some(receipt) = operation.integration().map_err(preparation)? {
1067 let source = decoded
1068 .get(&receipt.source_operation)
1069 .ok_or(Error::Invalid(
1070 "hosted integration original source proof absent",
1071 ))?;
1072 receipt.validate_source(source).map_err(preparation)?;
1073 required.insert(receipt.source_operation);
1074 pending.insert(receipt.source_operation);
1075 }
1076 edges.insert(id, required);
1077 pending.extend(
1078 operation
1079 .parents
1080 .iter()
1081 .filter(|id| !seen.contains(id))
1082 .copied(),
1083 );
1084 }
1085 if seen.len() != decoded.len()
1086 || used_threads
1087 .union(&inherited_bases)
1088 .copied()
1089 .collect::<BTreeSet<_>>()
1090 != geneses.keys().copied().collect()
1091 {
1092 return Err(Error::Invalid("unselected source proofs"));
1093 }
1094 for (thread, frontier) in &claim_frontiers {
1097 if !claim_threads.contains(thread) {
1098 continue;
1099 }
1100 let mut history = BTreeSet::new();
1101 let mut pending = frontier.clone();
1102 while let Some(id) = pending.pop_first() {
1103 if !history.insert(id) {
1104 continue;
1105 }
1106 let operation = decoded
1107 .get(&id)
1108 .ok_or(Error::Invalid("claim cutoff ancestry absent"))?;
1109 if operation.thread != *thread {
1110 return Err(Error::Invalid("claim cutoff crosses Thread"));
1111 }
1112 if !foreign_endpoints.contains(&id) {
1113 pending.extend(&operation.parents);
1114 }
1115 }
1116 {
1117 for (id, operation) in &decoded {
1118 if operation.thread == *thread
1119 && !history.contains(id)
1120 && !foreign_endpoints.contains(id)
1121 {
1122 if matches!(
1123 operation.source_author().map_err(preparation)?,
1124 Some(
1125 heddle_object_model::object::thread_replication::SourceAuthor::LocalKey
1126 )
1127 ) {
1128 return Err(Error::Invalid(
1129 "new local source lies outside signed ownership cutoff",
1130 ));
1131 }
1132 edges
1133 .get_mut(id)
1134 .ok_or(Error::Invalid("source topology entry absent"))?
1135 .extend(frontier);
1136 }
1137 }
1138 }
1139 }
1140 let references = decoded
1141 .values()
1142 .map(|operation| {
1143 operation
1144 .reference_proof(
1145 geneses
1146 .get(&operation.thread)
1147 .ok_or(Error::Invalid("dependency genesis absent"))?,
1148 )
1149 .map_err(preparation)
1150 })
1151 .collect::<Result<Vec<_>, _>>()?
1152 .into_iter()
1153 .flatten()
1154 .collect::<Vec<_>>();
1155 let capture = decoded
1156 .get(&selected_id)
1157 .ok_or(Error::Invalid("selected source operation absent"))?
1158 .source_result()
1159 .map_err(preparation)?
1160 .ok_or(Error::Invalid("selected operation has no source result"))?;
1161 let verified_ancestry = ancestry::verify(
1165 &ancestry,
1166 &import_floors,
1167 thread,
1168 selected_in_floor.then_some(state.id()),
1169 require_import_ancestry,
1170 )?;
1171 let pack = PackReader::open(
1172 &directory.path().join("source.pack"),
1173 &directory.path().join("source.idx"),
1174 directory.path(),
1175 )
1176 .map_err(preparation)?;
1177 let (references, visibility) = if selected_in_floor {
1180 (Vec::new(), None)
1181 } else {
1182 (references, capture.visibility.as_ref())
1183 };
1184 let _partial_trees = if allow_partial {
1185 Some(
1186 pack.validate_visible_source_closure(&state, SOURCE_BYTES)
1187 .map_err(preparation)?,
1188 )
1189 } else {
1190 pack.validate_source_closure_with_metadata(&state, &references, visibility, SOURCE_BYTES)
1191 .map_err(preparation)?;
1192 None
1193 };
1194 let mut ready_ids: BTreeSet<_> = edges
1197 .iter()
1198 .filter(|(_, parents)| parents.is_empty())
1199 .map(|(id, _)| *id)
1200 .collect();
1201 let mut children: BTreeMap<ContentHash, Vec<ContentHash>> = BTreeMap::new();
1202 for (child, parents) in &edges {
1203 for parent in parents {
1204 children.entry(*parent).or_default().push(*child);
1205 }
1206 }
1207 let mut ordered = Vec::new();
1208 while let Some(id) = ready_ids.pop_first() {
1209 ordered.push(
1210 originals
1211 .remove(&id)
1212 .ok_or(Error::Invalid("duplicate source topology identity"))?,
1213 );
1214 if let Some(dependants) = children.get(&id) {
1215 for child in dependants {
1216 let parents = edges
1217 .get_mut(child)
1218 .ok_or(Error::Invalid("incomplete source topology"))?;
1219 parents.remove(&id);
1220 if parents.is_empty() {
1221 ready_ids.insert(*child);
1222 }
1223 }
1224 }
1225 }
1226 if !originals.is_empty() {
1227 return Err(Error::Invalid("source dependency cycle"));
1228 }
1229 let mut genesis_record = original.clone();
1234 omit_unusable_ownership(&mut genesis_record, &decoded)?;
1235 for dependency in &mut dependencies {
1236 omit_unusable_ownership(dependency, &decoded)?;
1237 }
1238 Ok(ValidatedSourceArtifacts {
1239 import_authority: None,
1240 native_authority: None,
1241 _scratch_lease: scratch_lease,
1242 directory,
1243 genesis: genesis_record,
1244 operations: ordered,
1245 authority_admissions,
1246 dependencies,
1247 state,
1248 #[cfg(feature = "native")]
1249 partial_trees: _partial_trees,
1250 ancestry: verified_ancestry.floors,
1251 })
1252}
1253fn source_frontier_is_installable(
1254 thread: ContentHash,
1255 frontier: &BTreeSet<ContentHash>,
1256 operations: &BTreeMap<ContentHash, ThreadOperation>,
1257) -> bool {
1258 frontier.iter().all(|id| {
1259 operations.get(id).is_some_and(|operation| {
1260 operation.thread == thread && operation.source_state().ok().flatten().is_some()
1261 })
1262 })
1263}
1264fn omit_unusable_ownership(
1265 wrapper: &mut ThreadGenesisRecord,
1266 operations: &BTreeMap<ContentHash, ThreadOperation>,
1267) -> Result<(), Error> {
1268 if wrapper.ownership_claims.is_empty() && wrapper.ownership_resolutions.is_empty() {
1269 return Ok(());
1270 }
1271 let signed = wrapper
1272 .genesis
1273 .as_ref()
1274 .ok_or(Error::Invalid("claim genesis absent"))?;
1275 let genesis = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
1276 &signed.canonical_record,
1277 )
1278 .map_err(preparation)?;
1279 let thread = genesis.id().map_err(preparation)?;
1280 let mut kept_claims = BTreeSet::new();
1281 let mut claims = Vec::new();
1282 for record in wrapper.ownership_claims.drain(..) {
1283 let crate::thread_ownership::ClaimProof::Complete(proof) =
1284 crate::thread_ownership::decode(&record).map_err(preparation)?
1285 else {
1286 return Err(Error::Invalid(
1287 "transferred claim requires both original signatures",
1288 ));
1289 };
1290 let claim = proof.verify().map_err(preparation)?;
1291 if source_frontier_is_installable(thread, &claim.source_frontier, operations) {
1292 kept_claims.insert(claim.id().map_err(preparation)?);
1293 claims.push(record);
1294 }
1295 }
1296 wrapper.ownership_claims = claims;
1297 let mut claim_admissions = Vec::new();
1298 for record in wrapper.ownership_claim_admissions.drain(..) {
1299 let statement = admission_subject(&record)?;
1300 let Some(id) = statement.subject.claim_id() else {
1301 return Err(Error::Invalid("ownership claim admission subject"));
1302 };
1303 if kept_claims.contains(&id) {
1304 claim_admissions.push(record);
1305 }
1306 }
1307 wrapper.ownership_claim_admissions = claim_admissions;
1308 let mut kept_resolutions = BTreeSet::new();
1309 let mut resolutions = Vec::new();
1310 for record in wrapper.ownership_resolutions.drain(..) {
1311 let proof = crate::thread_ownership::decode_resolution(&record).map_err(preparation)?;
1312 let resolution = heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&proof.canonical).map_err(preparation)?;
1313 let claims_present = kept_claims.contains(&resolution.winning_claim)
1314 && resolution
1315 .conflicting_claims
1316 .iter()
1317 .all(|id| kept_claims.contains(id));
1318 if claims_present
1319 && source_frontier_is_installable(thread, &resolution.frontier, operations)
1320 {
1321 kept_resolutions.insert(resolution.id().map_err(preparation)?);
1322 resolutions.push(record);
1323 }
1324 }
1325 wrapper.ownership_resolutions = resolutions;
1326 let mut resolution_admissions = Vec::new();
1327 for record in wrapper.ownership_resolution_admissions.drain(..) {
1328 let statement = admission_subject(&record)?;
1329 if matches!(
1330 statement.subject,
1331 heddle_object_model::object::thread_authority_admission::OriginalAuthoritySubject::OwnershipResolution(id)
1332 if kept_resolutions.contains(&id)
1333 ) {
1334 resolution_admissions.push(record);
1335 }
1336 }
1337 wrapper.ownership_resolution_admissions = resolution_admissions;
1338 Ok(())
1339}
1340fn admission_subject(
1341 record: &SignedRecord,
1342) -> Result<heddle_object_model::object::thread_authority_admission::ThreadAuthorityAdmission, Error>
1343{
1344 crate::authority_admission::decode(record)
1345 .map_err(preparation)?
1346 .verify_signature()
1347 .map_err(preparation)
1348}
1349fn foreign_endpoint(
1350 foreign: &[ForeignDependencyV1],
1351 operation: &ThreadOperation,
1352 signed: &SignedOperation,
1353) -> Result<bool, Error> {
1354 if !foreign
1355 .iter()
1356 .any(|r| r.thread_genesis_digest.as_slice() == operation.thread.as_bytes())
1357 {
1358 return Ok(false);
1359 }
1360 let original = SignedRecord {
1361 format: heddle_object_model::object::thread_replication::OPERATION_FORMAT.into(),
1362 canonical_record: signed.canonical.clone(),
1363 signatures: vec![RecordSignature {
1364 public_key: operation.publisher.to_vec(),
1365 signature: signed.signature.clone(),
1366 }],
1367 };
1368 let digest = api::import_authority::signed_native_digest(&original)?;
1369 Ok(foreign.iter().any(|r| {
1370 r.thread_genesis_digest.as_slice() == operation.thread.as_bytes()
1371 && r.signed_native_digest == digest
1372 }))
1373}
1374
1375fn preparation(error: impl std::fmt::Display) -> Error {
1376 Error::Preparation(error.to_string())
1377}
1378
1379#[cfg(all(test, feature = "native"))]
1380#[path = "staging_tests.rs"]
1381mod tests;