Skip to main content

heddle_thread_api/fetch/
staging.rs

1//! Source bytes remain temporary until the terminal receipt and exact closure
2//! have been checked. Staging never changes a repository or a checkout.
3use 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
22/// Structurally checked source artifacts and original proofs. These bytes grant
23/// no authority. Dropping this value removes its temporary files.
24pub 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    /// Retained bytes are structural evidence, never installation authority.
51    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    /// Refresh witness-set/proof metadata after retirement or freshness renewal
70    /// while keeping the downloaded pack and every original signed byte.
71    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    /// Partial source remains read-only until a verified full closure arrives.
87    pub fn is_complete(&self) -> bool {
88        self.ready.full_closure_available
89    }
90}
91impl<R: MessageReader<Error = transport::Error>> Download<R> {
92    /// Consume one complete source download into bounded temporary files and
93    /// validate its original causal ancestry and exact selected source closure.
94    pub async fn stage(self, scratch: &Path) -> Result<StagedSource, Error> {
95        self.stage_inner(scratch, None).await
96    }
97    /// Use independently authenticated import certificates for original Git
98    /// ancestry. Exact native binding is checked after the originals arrive.
99    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}
298/// Structurally verified original source and actual artifact closure. This is
299/// not an author, audience, executor, or sharing-policy admission decision.
300pub 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    /// Retain the publication's exact originals and source as a hosted staged
321    /// install. The receiver supplies independently verified owner observation.
322    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        // Match immutable claims and signatures only. This self-described key
640        // is not enrolled here; the receiver must independently pin the issuer.
641        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    // A claim cutoff is an additional signed causal barrier. Ancestors retain
740    // their original local author; work outside it must follow the claim.
741    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    // Dependency-first installation makes foreign source authority available
819    // before admitting a local integration. Cycles cannot settle this graph.
820    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;