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