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    #[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        // Keep exact causal ancestors too. Native LocalKey history has no P2
134        // sidecar per capture; its admitted claim covers the earlier frontier.
135        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    /// Retained bytes are structural evidence, never installation authority.
233    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    /// Refresh witness-set/proof metadata after retirement or freshness renewal
252    /// while keeping the downloaded pack and every original signed byte.
253    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    /// Partial source remains read-only until a verified full closure arrives.
269    pub fn is_complete(&self) -> bool {
270        self.ready.full_closure_available
271    }
272}
273impl<R: MessageReader<Error = transport::Error>> Download<R> {
274    /// Consume one complete source download into bounded temporary files and
275    /// validate its original causal ancestry and exact selected source closure.
276    pub async fn stage(self, scratch: &Path) -> Result<StagedSource, Error> {
277        self.stage_inner(scratch, None).await
278    }
279    /// Use independently authenticated import certificates for original Git
280    /// ancestry. Exact native binding is checked after the originals arrive.
281    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}
494/// Structurally verified original source and actual artifact closure. This is
495/// not an author, audience, executor, or sharing-policy admission decision.
496pub 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    /// Retain the publication's exact originals and source as a hosted staged
517    /// install. The receiver supplies independently verified owner observation.
518    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        // Match immutable claims and signatures only. This self-described key
841        // is not enrolled here; the receiver must independently pin the issuer.
842        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            // Exact foreign references cut only structural staging. Installation
869            // must resolve this original through its verified retained prefix
870            // under current receiver trust, just like NativeClosure's resolver.
871            // Its parents belong to that origin's content closure.
872            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    // A claim cutoff is an additional signed causal barrier. Ancestors retain
958    // their original local author; work outside it must follow the claim.
959    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    // Dependency-first installation makes foreign source authority available
1045    // before admitting a local integration. Cycles cannot settle this graph.
1046    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;