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 pub(super) partial_trees: Vec<heddle_object_model::object::PartialTree>,
31 pub(super) authority_admissions:
32 BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
33}
34impl StagedSource {
35 pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
36 [
37 self.directory.path().join("source.pack"),
38 self.directory.path().join("source.idx"),
39 ]
40 }
41 pub fn operations(&self) -> &[SignedOperation] {
42 &self.operations
43 }
44 pub fn dependency_geneses(&self) -> &[ThreadGenesisRecord] {
45 &self.dependencies
46 }
47 pub fn ready(&self) -> &TransferReady {
48 &self.ready
49 }
50 pub fn import_authority(&self) -> Option<&crate::contract::ImportPublicProofBundleV1> {
52 self.ready.import_authority.as_ref()
53 }
54 pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
55 self.ready.native_authority.as_ref()
56 }
57 pub fn refresh_native_authority(
58 &mut self,
59 refreshed: NativePublicProofBundleV1,
60 ) -> Result<(), Error> {
61 let original = self
62 .ready
63 .native_authority
64 .as_mut()
65 .ok_or(Error::HostedTrustRequired)?;
66 crate::hybrid::history::replace_native_receiver_metadata(original, refreshed)?;
67 Ok(())
68 }
69 pub fn refresh_import_authority(
72 &mut self,
73 refreshed: crate::contract::ImportPublicProofBundleV1,
74 ) -> Result<(), Error> {
75 let original = self
76 .ready
77 .import_authority
78 .as_mut()
79 .ok_or(Error::HostedTrustRequired)?;
80 crate::hybrid::history::replace_receiver_metadata(original, refreshed)?;
81 Ok(())
82 }
83 pub fn state(&self) -> &State {
84 &self.state
85 }
86 pub fn is_complete(&self) -> bool {
88 self.ready.full_closure_available
89 }
90}
91impl<R: MessageReader<Error = transport::Error>> Download<R> {
92 pub async fn stage(self, scratch: &Path) -> Result<StagedSource, Error> {
95 self.stage_inner(scratch, None).await
96 }
97 pub async fn stage_with_import_carriers(
100 self,
101 scratch: &Path,
102 carriers: crypto::import_authority::VerifiedImportCarriers,
103 ) -> Result<StagedSource, Error> {
104 self.stage_inner(scratch, Some(carriers)).await
105 }
106 async fn stage_inner(
107 mut self,
108 scratch: &Path,
109 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
110 ) -> Result<StagedSource, Error> {
111 if let Some(carriers) = &carriers {
112 let mut original = self
113 .state
114 .ready
115 .import_authority
116 .clone()
117 .ok_or(Error::HostedTrustRequired)?;
118 crate::hybrid::history::replace_receiver_metadata(
119 &mut original,
120 carriers.bundle().clone(),
121 )
122 .map_err(preparation)?;
123 }
124 if self.state.facets != [SharedFacet::Source as i32] {
125 return Err(Error::Invalid("staging requires the source facet alone"));
126 }
127 let total = self
128 .state
129 .ready
130 .packs
131 .iter()
132 .try_fold(0u64, |sum, extent| sum.checked_add(extent.length))
133 .ok_or(Error::Invalid("source artifact length overflow"))?;
134 if total > SOURCE_BYTES {
135 return Err(Error::Invalid("staged source exceeds 256 MiB"));
136 }
137 self.state.limits.max_operations = self.state.limits.max_operations.min(10_000);
138 let directory = tempfile::Builder::new()
139 .prefix("thread-download-")
140 .tempdir_in(scratch)?;
141 let mut files = [
142 tokio::fs::File::create(directory.path().join("source.pack")).await?,
143 tokio::fs::File::create(directory.path().join("source.idx")).await?,
144 ];
145 let mut operations = Vec::new();
146 let mut receipt_records = Vec::new();
147 let mut dependencies = Vec::new();
148 let mut metadata_bytes = 0usize;
149 let mut complete = false;
150 while let Some(item) = self.next().await? {
151 match item {
152 Item::Pack(chunk) => {
153 let kind = chunk
154 .extent
155 .as_ref()
156 .ok_or(Error::Invalid("chunk extent absent"))?
157 .kind;
158 let index = match pack_extent::Kind::try_from(kind) {
159 Ok(pack_extent::Kind::NativePack) => 0,
160 Ok(pack_extent::Kind::NativeIndex) => 1,
161 _ => return Err(Error::Invalid("native source artifacts required")),
162 };
163 files[index].write_all(&chunk.data).await?;
164 }
165 Item::Operations(batch) => {
166 metadata_bytes = metadata_bytes
167 .checked_add(batch.encoded_len())
168 .ok_or(Error::Invalid("source metadata length overflow"))?;
169 if metadata_bytes > METADATA_BYTES {
170 return Err(Error::Invalid("staged source metadata exceeds 16 MiB"));
171 }
172 for received in crate::authority_admission::match_batch(&batch)? {
173 operations.push(received.original);
174 receipt_records.extend(received.authority_admission);
175 }
176 }
177 Item::ThreadGenesis(record) => {
178 metadata_bytes = metadata_bytes
179 .checked_add(record.encoded_len())
180 .ok_or(Error::Invalid("source metadata length overflow"))?;
181 if metadata_bytes > METADATA_BYTES || dependencies.len() >= 127 {
182 return Err(Error::Invalid("dependency metadata exceeds bounds"));
183 }
184 dependencies.push(record);
185 }
186 Item::Complete(_) => complete = true,
187 Item::Sidecar(_) => return Err(Error::Invalid("source staging excludes sidecars")),
188 }
189 }
190 if !complete {
191 return Err(Error::Invalid("source staging requires Complete"));
192 }
193 for file in &mut files {
194 file.flush().await?;
195 file.sync_all().await?;
196 }
197 drop(files);
198 let mut ready = self.state.ready;
199 if let Some(carriers) = &carriers {
200 ready.import_authority = Some(carriers.bundle().clone());
201 }
202 tokio::task::spawn_blocking(move || {
203 validate_with_receipts_and_carriers(
204 directory,
205 ready,
206 operations,
207 dependencies,
208 receipt_records,
209 carriers,
210 )
211 })
212 .await
213 .map_err(|error| Error::Preparation(error.to_string()))?
214 }
215}
216#[cfg(test)]
217fn validate(
218 directory: tempfile::TempDir,
219 ready: TransferReady,
220 operations: Vec<SignedOperation>,
221 dependencies: Vec<ThreadGenesisRecord>,
222) -> Result<StagedSource, Error> {
223 validate_with_receipts(directory, ready, operations, dependencies, Vec::new())
224}
225
226struct DisclosureInput {
227 directory: tempfile::TempDir,
228 operations: Vec<SignedOperation>,
229 dependency_records: Vec<ThreadGenesisRecord>,
230 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
231 allow_partial: bool,
232 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
233}
234
235#[cfg(test)]
236pub(super) fn validate_with_receipts(
237 directory: tempfile::TempDir,
238 ready: TransferReady,
239 operations: Vec<SignedOperation>,
240 dependencies: Vec<ThreadGenesisRecord>,
241 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
242) -> Result<StagedSource, Error> {
243 validate_with_receipts_and_carriers(
244 directory,
245 ready,
246 operations,
247 dependencies,
248 receipt_records,
249 None,
250 )
251}
252pub(super) fn validate_with_receipts_and_carriers(
253 directory: tempfile::TempDir,
254 ready: TransferReady,
255 operations: Vec<SignedOperation>,
256 dependencies: Vec<ThreadGenesisRecord>,
257 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
258 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
259) -> Result<StagedSource, Error> {
260 if carriers
261 .as_ref()
262 .is_some_and(|c| ready.import_authority.as_ref() != Some(c.bundle()))
263 {
264 return Err(Error::HostedTrustRequired);
265 }
266 let value = validate_disclosure_artifacts(
267 ready
268 .thread
269 .as_ref()
270 .ok_or(Error::Invalid("Thread absent"))?,
271 ready
272 .current
273 .as_ref()
274 .ok_or(Error::Invalid("revision absent"))?,
275 ready
276 .thread_genesis
277 .as_ref()
278 .ok_or(Error::Invalid("original genesis absent"))?,
279 DisclosureInput {
280 directory,
281 operations,
282 dependency_records: dependencies,
283 receipt_records,
284 allow_partial: !ready.full_closure_available,
285 carriers,
286 },
287 )?;
288 Ok(StagedSource {
289 directory: value.directory,
290 ready,
291 operations: value.operations,
292 dependencies: value.dependencies,
293 state: value.state,
294 partial_trees: value.partial_trees,
295 authority_admissions: value.authority_admissions,
296 })
297}
298pub struct ValidatedSourceArtifacts {
301 pub(crate) native_authority: Option<NativePublicProofBundleV1>,
302 pub(crate) import_authority: Option<ImportPublicProofBundleV1>,
303 directory: tempfile::TempDir,
304 operations: Vec<SignedOperation>,
305 genesis: ThreadGenesisRecord,
306 dependencies: Vec<ThreadGenesisRecord>,
307 state: State,
308 partial_trees: Vec<heddle_object_model::object::PartialTree>,
309 authority_admissions:
310 BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
311}
312impl ValidatedSourceArtifacts {
313 pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
314 self.native_authority.as_ref()
315 }
316 pub fn import_authority(&self) -> Option<&ImportPublicProofBundleV1> {
317 self.import_authority.as_ref()
318 }
319
320 pub fn into_hosted_source(self, mut ready: TransferReady) -> Result<StagedSource, Error> {
323 if ready.import_authority != self.import_authority
324 || ready.native_authority != self.native_authority
325 || (self.import_authority.is_none() && self.native_authority.is_none())
326 {
327 return Err(Error::HostedTrustRequired);
328 }
329 let reference = ready
330 .thread
331 .as_ref()
332 .ok_or(Error::Invalid("Thread absent"))?;
333 super::verify_origin(&self.genesis, reference)?;
334 let revision = ready
335 .current
336 .as_ref()
337 .ok_or(Error::Invalid("revision absent"))?;
338 if revision.spool != reference.spool
339 || revision.revision
340 != Some(revision_ref::Revision::State(
341 api::heddle::api::common::StateId {
342 value: self.state.id().as_bytes().to_vec(),
343 },
344 ))
345 || !ready.full_closure_available
346 {
347 return Err(Error::Invalid(
348 "publication source differs from hosted install selection",
349 ));
350 }
351 ready.thread_genesis = Some(self.genesis);
352 Ok(StagedSource {
353 directory: self.directory,
354 ready,
355 operations: self.operations,
356 dependencies: self.dependencies,
357 state: self.state,
358 partial_trees: self.partial_trees,
359 authority_admissions: self.authority_admissions,
360 })
361 }
362
363 pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
364 [
365 self.directory.path().join("source.pack"),
366 self.directory.path().join("source.idx"),
367 ]
368 }
369 pub fn operations(&self) -> &[SignedOperation] {
370 &self.operations
371 }
372 pub fn geneses(&self) -> impl Iterator<Item = &ThreadGenesisRecord> {
373 std::iter::once(&self.genesis).chain(&self.dependencies)
374 }
375 pub fn state(&self) -> &State {
376 &self.state
377 }
378 pub fn authority_admissions(
379 &self,
380 ) -> &BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission> {
381 &self.authority_admissions
382 }
383}
384#[allow(clippy::too_many_arguments)]
385pub(crate) fn validate_artifacts(
386 directory: tempfile::TempDir,
387 thread: &ThreadRef,
388 revision: &RevisionRef,
389 original: &ThreadGenesisRecord,
390 operations: Vec<SignedOperation>,
391 dependency_records: Vec<ThreadGenesisRecord>,
392 receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
393 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
394) -> Result<ValidatedSourceArtifacts, Error> {
395 validate_disclosure_artifacts(
396 thread,
397 revision,
398 original,
399 DisclosureInput {
400 directory,
401 operations,
402 dependency_records,
403 receipt_records,
404 allow_partial: false,
405 carriers,
406 },
407 )
408}
409
410fn validate_disclosure_artifacts(
411 thread: &ThreadRef,
412 revision: &RevisionRef,
413 original: &ThreadGenesisRecord,
414 input: DisclosureInput,
415) -> Result<ValidatedSourceArtifacts, Error> {
416 let DisclosureInput {
417 directory,
418 operations,
419 dependency_records,
420 receipt_records,
421 allow_partial,
422 carriers,
423 } = input;
424 if operations.len() > 10_000
425 || dependency_records.len() >= 128
426 || receipt_records.len() > operations.len()
427 {
428 return Err(Error::Invalid("source original graph exceeds bounds"));
429 }
430 let mut metadata = original.encoded_len();
431 for record in &dependency_records {
432 metadata = metadata.saturating_add(record.encoded_len());
433 }
434 for operation in &operations {
435 metadata = metadata.saturating_add(operation.canonical.len() + operation.signature.len());
436 }
437 for receipt in &receipt_records {
438 metadata = metadata.saturating_add(receipt.canonical.len() + receipt.signature.len());
439 }
440 let mut evidence_ids = BTreeSet::new();
441 for wrapper in std::iter::once(original).chain(&dependency_records) {
442 for record in &wrapper.boundary_acceptances {
443 evidence_ids.insert(heddle_object_model::object::ContentHash::compute_typed(
444 heddle_object_model::object::original_boundary_acceptance::FORMAT,
445 &record.canonical_record,
446 ));
447 if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
448 return Err(Error::Invalid("boundary evidence count exceeded"));
449 }
450 }
451 }
452 for receipt in &receipt_records {
453 if let Some(evidence) = &receipt.boundary_acceptance {
454 if evidence_ids.insert(
455 evidence
456 .verify_signature()
457 .map_err(preparation)?
458 .id()
459 .map_err(preparation)?,
460 ) {
461 metadata =
462 metadata.saturating_add(evidence.canonical.len() + evidence.signature.len());
463 }
464 if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
465 return Err(Error::Invalid("boundary evidence count exceeded"));
466 }
467 }
468 }
469 if metadata > METADATA_BYTES {
470 return Err(Error::Invalid("source metadata exceeds 16 MiB"));
471 }
472 if revision.spool != thread.spool {
473 return Err(Error::Invalid("source revision crosses Spool"));
474 }
475 let genesis = super::verify_origin(original, thread)?;
476 let Some(revision_ref::Revision::State(selected)) = revision.revision.as_ref() else {
477 return Err(Error::Invalid("exact native State required"));
478 };
479 let selected_thread = genesis.id().map_err(preparation)?;
480 if operations.is_empty() {
481 if !dependency_records.is_empty() || !receipt_records.is_empty() {
482 return Err(Error::Invalid(
483 "initial source cannot carry dependency originals",
484 ));
485 }
486 let state =
487 heddle_object_model::object::thread_replication::initial_base::synthetic_initial_base()
488 .map_err(preparation)?;
489 let canonical = state.encode_current_msgpack().map_err(preparation)?;
490 heddle_object_model::object::thread_replication::initial_base::initial_base_state(
491 &genesis, &canonical,
492 )
493 .map_err(preparation)?;
494 if selected.value.as_slice() != state.id().as_bytes() {
495 return Err(Error::Invalid(
496 "selected initial source differs from canonical seed",
497 ));
498 }
499 PackReader::open(
500 &directory.path().join("source.pack"),
501 &directory.path().join("source.idx"),
502 )
503 .map_err(preparation)?
504 .validate_source_closure_with_metadata(&state, &[], None, SOURCE_OBJECTS, SOURCE_BYTES)
505 .map_err(preparation)?;
506 return Ok(ValidatedSourceArtifacts {
507 import_authority: None,
508 native_authority: None,
509 directory,
510 operations,
511 genesis: original.clone(),
512 dependencies: Vec::new(),
513 state,
514 partial_trees: Vec::new(),
515 authority_admissions: BTreeMap::new(),
516 });
517 }
518 let mut geneses = BTreeMap::from([(selected_thread, genesis)]);
519 let mut dependencies = Vec::new();
520 for wrapper in dependency_records {
521 let record = wrapper
522 .genesis
523 .as_ref()
524 .ok_or(Error::Invalid("dependency signed genesis absent"))?;
525 let candidate = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
526 &record.canonical_record,
527 )
528 .map_err(preparation)?;
529 let reference = ThreadRef {
530 spool: thread.spool.clone(),
531 id: Some(ThreadId {
532 value: candidate.id().map_err(preparation)?.as_bytes().to_vec(),
533 }),
534 };
535 let candidate = super::verify_origin(&wrapper, &reference)?;
536 let id = candidate.id().map_err(preparation)?;
537 if geneses.len() >= 128 || geneses.insert(id, candidate).is_some() {
538 return Err(Error::Invalid(
539 "duplicate or oversized dependency genesis set",
540 ));
541 }
542 dependencies.push(wrapper);
543 }
544 let mut claim_frontiers = BTreeMap::new();
545 for wrapper in std::iter::once(original).chain(&dependencies) {
546 let signed = wrapper
547 .genesis
548 .as_ref()
549 .ok_or(Error::Invalid("claim genesis absent"))?;
550 let genesis = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
551 &signed.canonical_record,
552 )
553 .map_err(preparation)?;
554 let mut frontier = BTreeSet::new();
555 let claims = crate::replication::ownership::verify_claims(wrapper, &genesis)?;
556 let resolutions = crate::replication::ownership::verify_resolutions(wrapper, &genesis)?;
557 if claims.is_empty() && resolutions.is_empty() {
558 continue;
559 }
560 for claim in claims {
561 frontier.extend(
562 claim
563 .original
564 .verify()
565 .map_err(preparation)?
566 .source_frontier,
567 );
568 }
569 for resolution in resolutions {
570 frontier.extend(
571 heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&resolution.original.canonical)
572 .map_err(preparation)?.frontier,
573 );
574 }
575 claim_frontiers.insert(genesis.id().map_err(preparation)?, frontier);
576 }
577 let mut originals = BTreeMap::new();
578 let mut decoded = BTreeMap::<ContentHash, ThreadOperation>::new();
579 let mut selected_operation = None;
580 let mut source_thread = selected_thread;
581 let mut inherited_bases = BTreeSet::new();
582 for _ in 0..128 {
583 let current = geneses
584 .get(&source_thread)
585 .ok_or(Error::Invalid("fork base source genesis absent"))?;
586 if selected.value.as_slice() != current.base.as_bytes() {
587 break;
588 }
589 let parent = current.parent.ok_or(Error::Invalid(
590 "non-system base has no original parent source",
591 ))?;
592 if !inherited_bases.insert(source_thread) {
593 return Err(Error::Invalid("fork base parent cycle"));
594 }
595 let ancestor = geneses
596 .get(&parent)
597 .ok_or(Error::Invalid("fork base parent original absent"))?;
598 if ancestor.spool != current.spool {
599 return Err(Error::Invalid("fork base crosses Spool"));
600 }
601 source_thread = parent;
602 }
603 if selected.value.as_slice()
604 == geneses
605 .get(&source_thread)
606 .ok_or(Error::Invalid("fork base source genesis absent"))?
607 .base
608 .as_bytes()
609 {
610 return Err(Error::Invalid("fork base source chain exceeds bound"));
611 }
612 for signed in &operations {
613 let operation = signed.verify().map_err(preparation)?;
614 let id = operation.id().map_err(preparation)?;
615 let state = operation
616 .source_state()
617 .map_err(preparation)?
618 .ok_or(Error::Invalid("non-source operation in source ancestry"))?;
619 if operation.thread == source_thread
620 && state.id().as_bytes().as_slice() == selected.value
621 && selected_operation.replace((id, state)).is_some()
622 {
623 return Err(Error::Invalid("ambiguous selected source proof"));
624 }
625 originals.insert(id, signed.clone());
626 if decoded.insert(id, operation).is_some() {
627 return Err(Error::Invalid("duplicate source proof"));
628 }
629 }
630 let mut authority_admissions = BTreeMap::new();
631 for receipt in receipt_records {
632 let statement = receipt.verify_signature().map_err(preparation)?;
633 let operation_id = statement.subject.operation_id().ok_or(Error::Invalid(
634 "source batch cannot carry ownership claim admission",
635 ))?;
636 let original = originals
637 .get(&operation_id)
638 .ok_or(Error::Invalid("unmatched source authority receipt"))?;
639 receipt.verify(original, &heddle_object_model::object::thread_replication::integration::TrustedHostedExecutor {
642 spool: statement.spool, spool_genesis: statement.spool_genesis, executor: statement.executor,
643 }).map_err(preparation)?;
644 if authority_admissions.insert(operation_id, receipt).is_some() {
645 return Err(Error::Invalid("duplicate source authority receipt"));
646 }
647 }
648 let (selected_id, state) =
649 selected_operation.ok_or(Error::Invalid("selected source proof absent"))?;
650 let mut pending = BTreeSet::from([selected_id]);
651 let mut seen = BTreeSet::new();
652 let mut used_threads = BTreeSet::new();
653 let mut edges = BTreeMap::new();
654 while let Some(id) = pending.pop_first() {
655 if !seen.insert(id) {
656 continue;
657 }
658 let operation = decoded
659 .get(&id)
660 .ok_or(Error::Invalid("incomplete source ancestry"))?;
661 let parents = operation
662 .parents
663 .iter()
664 .map(|id| {
665 decoded
666 .get(id)
667 .cloned()
668 .ok_or(Error::Invalid("incomplete source ancestry"))
669 })
670 .collect::<Result<Vec<_>, _>>()?;
671 let genesis = geneses
672 .get(&operation.thread)
673 .ok_or(Error::Invalid("source dependency genesis absent"))?;
674 if used_threads.insert(operation.thread)
675 && let Some(frontier) = claim_frontiers.get(&operation.thread)
676 {
677 for head in frontier {
678 if decoded
679 .get(head)
680 .is_none_or(|source| source.thread != operation.thread)
681 {
682 return Err(Error::Invalid(
683 "ownership claim cutoff source proof absent or foreign",
684 ));
685 }
686 }
687 pending.extend(frontier);
688 }
689 let imported = carriers
690 .as_ref()
691 .map(|c| c.bind(genesis, operation, &parents))
692 .transpose()
693 .map_err(preparation)?
694 .flatten();
695 match imported {
696 Some(bound) => bound.validate_parents(genesis, &parents),
697 None => operation.validate_parents(genesis, &parents),
698 }
699 .map_err(preparation)?;
700 let mut required = operation.parents.clone();
701 if let Some(receipt) = operation.local_integration().map_err(preparation)? {
702 let source = decoded
703 .get(&receipt.source_operation)
704 .ok_or(Error::Invalid(
705 "local integration original source proof absent",
706 ))?;
707 receipt.validate_source(source).map_err(preparation)?;
708 required.insert(receipt.source_operation);
709 pending.insert(receipt.source_operation);
710 }
711 if let Some(receipt) = operation.integration().map_err(preparation)? {
712 let source = decoded
713 .get(&receipt.source_operation)
714 .ok_or(Error::Invalid(
715 "hosted integration original source proof absent",
716 ))?;
717 receipt.validate_source(source).map_err(preparation)?;
718 required.insert(receipt.source_operation);
719 pending.insert(receipt.source_operation);
720 }
721 edges.insert(id, required);
722 pending.extend(
723 operation
724 .parents
725 .iter()
726 .filter(|id| !seen.contains(id))
727 .copied(),
728 );
729 }
730 if seen.len() != decoded.len()
731 || used_threads
732 .union(&inherited_bases)
733 .copied()
734 .collect::<BTreeSet<_>>()
735 != geneses.keys().copied().collect()
736 {
737 return Err(Error::Invalid("unselected source proofs"));
738 }
739 for (thread, frontier) in &claim_frontiers {
742 let mut history = BTreeSet::new();
743 let mut pending = frontier.clone();
744 while let Some(id) = pending.pop_first() {
745 if !history.insert(id) {
746 continue;
747 }
748 let operation = decoded
749 .get(&id)
750 .ok_or(Error::Invalid("claim cutoff ancestry absent"))?;
751 if operation.thread != *thread {
752 return Err(Error::Invalid("claim cutoff crosses Thread"));
753 }
754 pending.extend(&operation.parents);
755 }
756 {
757 for (id, operation) in &decoded {
758 if operation.thread == *thread && !history.contains(id) {
759 if matches!(
760 operation.source_author().map_err(preparation)?,
761 Some(
762 heddle_object_model::object::thread_replication::SourceAuthor::LocalKey
763 )
764 ) {
765 return Err(Error::Invalid(
766 "new local source lies outside signed ownership cutoff",
767 ));
768 }
769 edges
770 .get_mut(id)
771 .ok_or(Error::Invalid("source topology entry absent"))?
772 .extend(frontier);
773 }
774 }
775 }
776 }
777 let references = decoded
778 .values()
779 .map(|operation| {
780 operation
781 .reference_proof(
782 geneses
783 .get(&operation.thread)
784 .ok_or(Error::Invalid("dependency genesis absent"))?,
785 )
786 .map_err(preparation)
787 })
788 .collect::<Result<Vec<_>, _>>()?
789 .into_iter()
790 .flatten()
791 .collect::<Vec<_>>();
792 let capture = decoded
793 .get(&selected_id)
794 .ok_or(Error::Invalid("selected source operation absent"))?
795 .source_result()
796 .map_err(preparation)?
797 .ok_or(Error::Invalid("selected operation has no source result"))?;
798 let pack = PackReader::open(
799 &directory.path().join("source.pack"),
800 &directory.path().join("source.idx"),
801 )
802 .map_err(preparation)?;
803 let partial_trees = if allow_partial {
804 pack.validate_visible_source_closure(&state, SOURCE_OBJECTS, SOURCE_BYTES)
805 .map_err(preparation)?
806 .partial_trees
807 } else {
808 pack.validate_source_closure_with_metadata(
809 &state,
810 &references,
811 capture.visibility.as_ref(),
812 SOURCE_OBJECTS,
813 SOURCE_BYTES,
814 )
815 .map_err(preparation)?;
816 Vec::new()
817 };
818 let mut ready_ids: BTreeSet<_> = edges
821 .iter()
822 .filter(|(_, parents)| parents.is_empty())
823 .map(|(id, _)| *id)
824 .collect();
825 let mut children: BTreeMap<ContentHash, Vec<ContentHash>> = BTreeMap::new();
826 for (child, parents) in &edges {
827 for parent in parents {
828 children.entry(*parent).or_default().push(*child);
829 }
830 }
831 let mut ordered = Vec::new();
832 while let Some(id) = ready_ids.pop_first() {
833 ordered.push(
834 originals
835 .remove(&id)
836 .ok_or(Error::Invalid("duplicate source topology identity"))?,
837 );
838 if let Some(dependants) = children.get(&id) {
839 for child in dependants {
840 let parents = edges
841 .get_mut(child)
842 .ok_or(Error::Invalid("incomplete source topology"))?;
843 parents.remove(&id);
844 if parents.is_empty() {
845 ready_ids.insert(*child);
846 }
847 }
848 }
849 }
850 if !originals.is_empty() {
851 return Err(Error::Invalid("source dependency cycle"));
852 }
853 Ok(ValidatedSourceArtifacts {
854 import_authority: None,
855 native_authority: None,
856 directory,
857 genesis: original.clone(),
858 operations: ordered,
859 authority_admissions,
860 dependencies,
861 state,
862 partial_trees,
863 })
864}
865fn preparation(error: impl std::fmt::Display) -> Error {
866 Error::Preparation(error.to_string())
867}
868
869#[cfg(test)]
870#[path = "staging_tests.rs"]
871mod tests;