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