use std::{
collections::{BTreeMap, BTreeSet},
path::Path,
};
use api::v2::client::MessageReader;
use crypto::thread_operation::SignedOperation;
use heddle_object_model::object::{
ContentHash, State, StateId, thread_replication::ThreadOperation,
};
use heddle_pack::store::pack::PackReader;
use prost::Message;
use tokio::io::AsyncWriteExt;
pub(super) use super::ancestry::AncestryInput;
use super::{Download, Error, Item, ancestry};
use crate::{contract::*, transport};
const METADATA_BYTES: usize = 16 * 1024 * 1024;
const SOURCE_BYTES: u64 = 256 * 1024 * 1024;
const SOURCE_OBJECTS: usize = 100_000;
pub struct StagedSource {
pub(super) directory: tempfile::TempDir,
pub(super) ready: TransferReady,
pub(super) operations: Vec<SignedOperation>,
pub(super) dependencies: Vec<ThreadGenesisRecord>,
pub(super) state: State,
#[cfg(feature = "native")]
pub(super) prefix_original: Option<SignedRecord>,
pub(super) partial_trees: Vec<heddle_object_model::object::PartialTree>,
pub(super) authority_admissions:
BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
pub(super) ancestry: Vec<ancestry::VerifiedFloor>,
}
impl StagedSource {
#[cfg(feature = "native")]
pub fn select_prefix(&mut self, reference: &ForeignDependencyV1) -> Result<(), Error> {
let proof = match (self.import_authority(), self.native_authority()) {
(Some(b), None) => crate::hybrid::authority::PublicProof::from(b.clone()),
(None, Some(b)) => crate::hybrid::authority::PublicProof::from(b.clone()),
_ => return Err(Error::HostedTrustRequired),
};
self.prefix_original = match proof.prefix_original(reference) {
Ok(record) => Some(record.clone()),
Err(_) => self
.operations
.iter()
.map(|signed| {
let operation = signed.verify().map_err(preparation)?;
Ok(SignedRecord {
format: heddle_object_model::object::thread_replication::OPERATION_FORMAT
.into(),
canonical_record: signed.canonical.clone(),
signatures: vec![RecordSignature {
public_key: operation.publisher.to_vec(),
signature: signed.signature.clone(),
}],
})
})
.collect::<Result<Vec<_>, Error>>()?
.into_iter()
.find(|r| {
api::import_authority::signed_native_digest(r)
.is_ok_and(|d| d == reference.signed_native_digest)
}),
};
if self.prefix_original.is_none() {
return Err(api::hybrid_codec::Reject::Scope.into());
}
let projected = proof.prefix(reference)?;
let (authorities, landings, geneses, imported) = match &projected {
crate::hybrid::authority::PublicProof::Import(b) => (
&b.authority_witnesses,
&b.landing_witnesses,
b.genesis_witnesses
.iter()
.filter_map(|p| p.original_genesis.as_ref())
.collect::<Vec<_>>(),
Some(b.as_ref()),
),
crate::hybrid::authority::PublicProof::Native(b) => (
&b.authority_witnesses,
&b.landing_witnesses,
b.genesis_witnesses
.iter()
.filter_map(|p| p.original_genesis.as_ref())
.collect::<Vec<_>>(),
None,
),
};
let retained = authorities
.iter()
.flat_map(|p| p.original.iter().chain(&p.dependencies))
.chain(landings.iter().flat_map(|p| {
p.execution
.iter()
.chain(p.source_operation.iter())
.chain(&p.review_evidence)
}))
.chain(self.prefix_original.iter())
.collect::<Vec<_>>();
let mut operations = Vec::new();
for signed in &self.operations {
let operation = signed.verify().map_err(preparation)?;
let id = operation.id().map_err(preparation)?;
let published = if let Some(imported) = imported {
let frontier = api::import_authority::frontier_digest(&ImportFrontierV1 {
format_version: 1,
thread_id: operation.thread.as_bytes().to_vec(),
operation_ids: vec![id.as_bytes().to_vec()],
})?;
imported.operations.iter().any(|o| {
o.body
.as_ref()
.is_some_and(|b| b.resulting_frontier_digest == frontier)
})
} else {
false
};
if published
|| retained.iter().any(|r| {
r.canonical_record == signed.canonical
&& r.signatures.iter().any(|s| {
s.public_key.as_slice() == operation.publisher
&& s.signature == signed.signature
})
})
{
operations.push(signed.clone());
}
}
let available = self
.operations
.iter()
.map(|signed| {
let operation = signed.verify().map_err(preparation)?;
Ok((operation.id().map_err(preparation)?, (signed, operation)))
})
.collect::<Result<BTreeMap<_, _>, Error>>()?;
let mut pending = operations
.iter()
.map(|signed| {
signed
.verify()
.map_err(preparation)?
.id()
.map_err(preparation)
})
.collect::<Result<Vec<_>, Error>>()?;
let mut causal = BTreeSet::new();
while let Some(id) = pending.pop() {
if causal.insert(id) {
let (signed, op) = available
.get(&id)
.ok_or(Error::Invalid("prefix ancestor absent"))?;
if !foreign_endpoint(projected.foreign_dependencies(), op, signed)? {
pending.extend(op.parents.iter().copied());
}
}
}
self.operations = available
.into_iter()
.filter(|(id, _)| causal.contains(id))
.map(|(_, (signed, _))| signed.clone())
.collect();
let mut threads = BTreeSet::new();
for record in geneses {
threads.insert(
crypto::import_authority::verify_native_genesis(record)
.map_err(preparation)?
.1
.id()
.map_err(preparation)?,
);
}
for signed in &self.operations {
threads.insert(signed.verify().map_err(preparation)?.thread);
}
for r in projected.foreign_dependencies() {
threads.insert(ContentHash::from_bytes(
r.thread_genesis_digest
.as_slice()
.try_into()
.map_err(|_| Error::Hybrid(api::hybrid_codec::Reject::Scope))?,
));
}
self.dependencies.retain(|g| {
g.genesis.as_ref().is_some_and(|r| {
crypto::import_authority::verify_native_genesis(r)
.is_ok_and(|(_, g)| g.id().is_ok_and(|id| threads.contains(&id)))
})
});
for wrapper in self
.ready
.thread_genesis
.iter_mut()
.chain(&mut self.dependencies)
{
wrapper.ownership_claims.retain(|r| retained.contains(&r));
wrapper
.ownership_resolutions
.retain(|r| retained.contains(&r));
}
match projected {
crate::hybrid::authority::PublicProof::Import(b) => {
self.ready.import_authority = Some(*b)
}
crate::hybrid::authority::PublicProof::Native(b) => {
self.ready.native_authority = Some(*b)
}
}
Ok(())
}
pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
[
self.directory.path().join("source.pack"),
self.directory.path().join("source.idx"),
]
}
pub fn operations(&self) -> &[SignedOperation] {
&self.operations
}
pub fn dependency_geneses(&self) -> &[ThreadGenesisRecord] {
&self.dependencies
}
pub fn ready(&self) -> &TransferReady {
&self.ready
}
pub fn import_authority(&self) -> Option<&crate::contract::ImportPublicProofBundleV1> {
self.ready.import_authority.as_ref()
}
pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
self.ready.native_authority.as_ref()
}
pub fn refresh_native_authority(
&mut self,
refreshed: NativePublicProofBundleV1,
) -> Result<(), Error> {
let original = self
.ready
.native_authority
.as_mut()
.ok_or(Error::HostedTrustRequired)?;
crate::hybrid::history::replace_native_receiver_metadata(original, refreshed)?;
Ok(())
}
pub fn refresh_import_authority(
&mut self,
refreshed: crate::contract::ImportPublicProofBundleV1,
) -> Result<(), Error> {
let original = self
.ready
.import_authority
.as_mut()
.ok_or(Error::HostedTrustRequired)?;
crate::hybrid::history::replace_receiver_metadata(original, refreshed)?;
Ok(())
}
pub fn state(&self) -> &State {
&self.state
}
pub fn is_complete(&self) -> bool {
self.ready.full_closure_available
}
pub fn verified_import_floors(&self) -> impl Iterator<Item = (StateId, &BTreeSet<StateId>)> {
self.ancestry
.iter()
.filter(|floor| floor.coverage == import_ancestry_page::Coverage::Floor)
.map(|floor| (floor.tip, &floor.members))
}
pub(super) fn ancestry_paths(&self) -> Option<[std::path::PathBuf; 2]> {
let pack = self.directory.path().join("ancestry.pack");
pack.exists()
.then(|| [pack, self.directory.path().join("ancestry.idx")])
}
}
impl<R: MessageReader<Error = transport::Error>> Download<R> {
pub async fn stage(self, scratch: &Path) -> Result<StagedSource, Error> {
self.stage_inner(scratch, None).await
}
pub async fn stage_with_import_carriers(
self,
scratch: &Path,
carriers: crypto::import_authority::VerifiedImportCarriers,
) -> Result<StagedSource, Error> {
self.stage_inner(scratch, Some(carriers)).await
}
async fn stage_inner(
mut self,
scratch: &Path,
carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
) -> Result<StagedSource, Error> {
if let Some(carriers) = &carriers {
let mut original = self
.state
.ready
.import_authority
.clone()
.ok_or(Error::HostedTrustRequired)?;
crate::hybrid::history::replace_receiver_metadata(
&mut original,
carriers.bundle().clone(),
)
.map_err(preparation)?;
}
if self.state.facets != [SharedFacet::Source as i32] {
return Err(Error::Invalid("staging requires the source facet alone"));
}
let total = self
.state
.ready
.packs
.iter()
.try_fold(0u64, |sum, extent| sum.checked_add(extent.length))
.ok_or(Error::Invalid("source artifact length overflow"))?;
if total > SOURCE_BYTES {
return Err(Error::Invalid("staged source exceeds 256 MiB"));
}
self.state.limits.max_operations = self.state.limits.max_operations.min(10_000);
let directory = tempfile::Builder::new()
.prefix("thread-download-")
.tempdir_in(scratch)?;
let mut files = [
tokio::fs::File::create(directory.path().join("source.pack")).await?,
tokio::fs::File::create(directory.path().join("source.idx")).await?,
];
let mut operations = Vec::new();
let mut receipt_records = Vec::new();
let mut dependencies = Vec::new();
let mut ancestry = AncestryInput::new(self.state.excluded_tips.clone());
let mut metadata_bytes = 0usize;
let mut complete = false;
while let Some(item) = self.next().await? {
match item {
Item::Pack(chunk) => {
let kind = chunk
.extent
.as_ref()
.ok_or(Error::Invalid("chunk extent absent"))?
.kind;
let index = match pack_extent::Kind::try_from(kind) {
Ok(pack_extent::Kind::NativePack) => 0,
Ok(pack_extent::Kind::NativeIndex) => 1,
_ => return Err(Error::Invalid("native source artifacts required")),
};
files[index].write_all(&chunk.data).await?;
}
Item::Operations(batch) => {
metadata_bytes = metadata_bytes
.checked_add(batch.encoded_len())
.ok_or(Error::Invalid("source metadata length overflow"))?;
if metadata_bytes > METADATA_BYTES {
return Err(Error::Invalid("staged source metadata exceeds 16 MiB"));
}
for received in crate::authority_admission::match_batch(&batch)? {
operations.push(received.original);
receipt_records.extend(received.authority_admission);
}
}
Item::ThreadGenesis(record) => {
metadata_bytes = metadata_bytes
.checked_add(record.encoded_len())
.ok_or(Error::Invalid("source metadata length overflow"))?;
if metadata_bytes > METADATA_BYTES || dependencies.len() >= 127 {
return Err(Error::Invalid("dependency metadata exceeds bounds"));
}
dependencies.push(record);
}
Item::ImportAncestry(page) => ancestry.push(
page,
self.state
.ready
.thread
.as_ref()
.ok_or(Error::Invalid("Thread absent"))?,
directory.path(),
)?,
Item::Complete(_) => complete = true,
Item::Sidecar(_) => return Err(Error::Invalid("source staging excludes sidecars")),
}
}
if !complete {
return Err(Error::Invalid("source staging requires Complete"));
}
for file in &mut files {
file.flush().await?;
file.sync_all().await?;
}
drop(files);
ancestry.finish()?;
let mut ready = self.state.ready;
if let Some(carriers) = &carriers {
ready.import_authority = Some(carriers.bundle().clone());
}
tokio::task::spawn_blocking(move || {
validate_with_receipts_and_carriers(
directory,
ready,
operations,
dependencies,
receipt_records,
carriers,
ancestry,
)
})
.await
.map_err(|error| Error::Preparation(error.to_string()))?
}
}
#[cfg(test)]
fn validate(
directory: tempfile::TempDir,
ready: TransferReady,
operations: Vec<SignedOperation>,
dependencies: Vec<ThreadGenesisRecord>,
) -> Result<StagedSource, Error> {
validate_with_receipts(directory, ready, operations, dependencies, Vec::new())
}
struct DisclosureInput {
directory: tempfile::TempDir,
operations: Vec<SignedOperation>,
dependency_records: Vec<ThreadGenesisRecord>,
receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
allow_partial: bool,
carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
foreign: Vec<ForeignDependencyV1>,
ancestry: AncestryInput,
require_import_ancestry: bool,
}
#[cfg(test)]
pub(crate) fn validate_with_receipts(
directory: tempfile::TempDir,
ready: TransferReady,
operations: Vec<SignedOperation>,
dependencies: Vec<ThreadGenesisRecord>,
receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
) -> Result<StagedSource, Error> {
validate_with_receipts_and_carriers(
directory,
ready,
operations,
dependencies,
receipt_records,
None,
AncestryInput::default(),
)
}
pub(super) fn validate_with_receipts_and_carriers(
directory: tempfile::TempDir,
mut ready: TransferReady,
operations: Vec<SignedOperation>,
dependencies: Vec<ThreadGenesisRecord>,
receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
ancestry: AncestryInput,
) -> Result<StagedSource, Error> {
if carriers
.as_ref()
.is_some_and(|c| ready.import_authority.as_ref() != Some(c.bundle()))
{
return Err(Error::HostedTrustRequired);
}
let value = validate_disclosure_artifacts(
ready
.thread
.as_ref()
.ok_or(Error::Invalid("Thread absent"))?,
ready
.current
.as_ref()
.ok_or(Error::Invalid("revision absent"))?,
ready
.thread_genesis
.as_ref()
.ok_or(Error::Invalid("original genesis absent"))?,
DisclosureInput {
directory,
operations,
dependency_records: dependencies,
receipt_records,
allow_partial: !ready.full_closure_available,
foreign: ready
.import_authority
.as_ref()
.map(|b| b.foreign_dependencies.clone())
.or_else(|| {
ready
.native_authority
.as_ref()
.map(|b| b.foreign_dependencies.clone())
})
.unwrap_or_default(),
carriers,
ancestry,
require_import_ancestry: true,
},
)?;
ready.thread_genesis = Some(value.genesis.clone());
Ok(StagedSource {
directory: value.directory,
ready,
#[cfg(feature = "native")]
prefix_original: None,
operations: value.operations,
dependencies: value.dependencies,
state: value.state,
partial_trees: value.partial_trees,
authority_admissions: value.authority_admissions,
ancestry: value.ancestry,
})
}
pub struct ValidatedSourceArtifacts {
pub(crate) native_authority: Option<NativePublicProofBundleV1>,
pub(crate) import_authority: Option<ImportPublicProofBundleV1>,
directory: tempfile::TempDir,
operations: Vec<SignedOperation>,
genesis: ThreadGenesisRecord,
dependencies: Vec<ThreadGenesisRecord>,
state: State,
partial_trees: Vec<heddle_object_model::object::PartialTree>,
authority_admissions:
BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission>,
ancestry: Vec<ancestry::VerifiedFloor>,
}
impl ValidatedSourceArtifacts {
pub fn native_authority(&self) -> Option<&NativePublicProofBundleV1> {
self.native_authority.as_ref()
}
pub fn import_authority(&self) -> Option<&ImportPublicProofBundleV1> {
self.import_authority.as_ref()
}
pub fn into_hosted_source(self, mut ready: TransferReady) -> Result<StagedSource, Error> {
if ready.import_authority != self.import_authority
|| ready.native_authority != self.native_authority
|| (self.import_authority.is_none() && self.native_authority.is_none())
{
return Err(Error::HostedTrustRequired);
}
let reference = ready
.thread
.as_ref()
.ok_or(Error::Invalid("Thread absent"))?;
super::verify_origin(&self.genesis, reference)?;
let revision = ready
.current
.as_ref()
.ok_or(Error::Invalid("revision absent"))?;
if revision.spool != reference.spool
|| revision.revision
!= Some(revision_ref::Revision::State(
api::heddle::api::common::StateId {
value: self.state.id().as_bytes().to_vec(),
},
))
|| !ready.full_closure_available
{
return Err(Error::Invalid(
"publication source differs from hosted install selection",
));
}
ready.thread_genesis = Some(self.genesis);
Ok(StagedSource {
directory: self.directory,
ready,
#[cfg(feature = "native")]
prefix_original: None,
operations: self.operations,
dependencies: self.dependencies,
state: self.state,
partial_trees: self.partial_trees,
authority_admissions: self.authority_admissions,
ancestry: self.ancestry,
})
}
pub fn artifact_paths(&self) -> [std::path::PathBuf; 2] {
[
self.directory.path().join("source.pack"),
self.directory.path().join("source.idx"),
]
}
pub fn operations(&self) -> &[SignedOperation] {
&self.operations
}
pub fn geneses(&self) -> impl Iterator<Item = &ThreadGenesisRecord> {
std::iter::once(&self.genesis).chain(&self.dependencies)
}
pub fn state(&self) -> &State {
&self.state
}
pub fn authority_admissions(
&self,
) -> &BTreeMap<ContentHash, crypto::thread_authority_admission::SignedAuthorityAdmission> {
&self.authority_admissions
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn validate_artifacts(
directory: tempfile::TempDir,
thread: &ThreadRef,
revision: &RevisionRef,
original: &ThreadGenesisRecord,
operations: Vec<SignedOperation>,
dependency_records: Vec<ThreadGenesisRecord>,
receipt_records: Vec<crypto::thread_authority_admission::SignedAuthorityAdmission>,
carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
foreign: Vec<ForeignDependencyV1>,
) -> Result<ValidatedSourceArtifacts, Error> {
validate_disclosure_artifacts(
thread,
revision,
original,
DisclosureInput {
directory,
operations,
dependency_records,
receipt_records,
allow_partial: false,
foreign,
carriers,
ancestry: AncestryInput::default(),
require_import_ancestry: false,
},
)
}
fn validate_disclosure_artifacts(
thread: &ThreadRef,
revision: &RevisionRef,
original: &ThreadGenesisRecord,
input: DisclosureInput,
) -> Result<ValidatedSourceArtifacts, Error> {
let DisclosureInput {
directory,
operations,
dependency_records,
receipt_records,
allow_partial,
carriers,
foreign,
ancestry,
require_import_ancestry,
} = input;
if !ancestry.is_empty() && carriers.is_none() {
return Err(Error::Invalid(
"import ancestry requires independently authenticated import carriers",
));
}
if operations.len() > 10_000
|| dependency_records.len() >= 128
|| receipt_records.len() > operations.len()
{
return Err(Error::Invalid("source original graph exceeds bounds"));
}
let mut metadata = original.encoded_len();
for record in &dependency_records {
metadata = metadata.saturating_add(record.encoded_len());
}
for operation in &operations {
metadata = metadata.saturating_add(operation.canonical.len() + operation.signature.len());
}
for receipt in &receipt_records {
metadata = metadata.saturating_add(receipt.canonical.len() + receipt.signature.len());
}
let mut evidence_ids = BTreeSet::new();
for wrapper in std::iter::once(original).chain(&dependency_records) {
for record in &wrapper.boundary_acceptances {
evidence_ids.insert(heddle_object_model::object::ContentHash::compute_typed(
heddle_object_model::object::original_boundary_acceptance::FORMAT,
&record.canonical_record,
));
if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
return Err(Error::Invalid("boundary evidence count exceeded"));
}
}
}
for receipt in &receipt_records {
if let Some(evidence) = &receipt.boundary_acceptance {
if evidence_ids.insert(
evidence
.verify_signature()
.map_err(preparation)?
.id()
.map_err(preparation)?,
) {
metadata =
metadata.saturating_add(evidence.canonical.len() + evidence.signature.len());
}
if evidence_ids.len() > crate::boundary_acceptance::MAX_ACCEPTANCES {
return Err(Error::Invalid("boundary evidence count exceeded"));
}
}
}
if metadata > METADATA_BYTES {
return Err(Error::Invalid("source metadata exceeds 16 MiB"));
}
if revision.spool != thread.spool {
return Err(Error::Invalid("source revision crosses Spool"));
}
let genesis = super::verify_origin(original, thread)?;
let Some(revision_ref::Revision::State(selected)) = revision.revision.as_ref() else {
return Err(Error::Invalid("exact native State required"));
};
let selected_thread = genesis.id().map_err(preparation)?;
if operations.is_empty() {
if !dependency_records.is_empty() || !receipt_records.is_empty() {
return Err(Error::Invalid(
"initial source cannot carry dependency originals",
));
}
let state =
heddle_object_model::object::thread_replication::initial_base::synthetic_initial_base()
.map_err(preparation)?;
let canonical = state.encode_current_msgpack().map_err(preparation)?;
heddle_object_model::object::thread_replication::initial_base::initial_base_state(
&genesis, &canonical,
)
.map_err(preparation)?;
if selected.value.as_slice() != state.id().as_bytes() {
return Err(Error::Invalid(
"selected initial source differs from canonical seed",
));
}
PackReader::open(
&directory.path().join("source.pack"),
&directory.path().join("source.idx"),
)
.map_err(preparation)?
.validate_source_closure_with_metadata(&state, &[], None, SOURCE_OBJECTS, SOURCE_BYTES)
.map_err(preparation)?;
if !ancestry.is_empty() {
return Err(Error::Invalid(
"initial source cannot carry import ancestry",
));
}
return Ok(ValidatedSourceArtifacts {
import_authority: None,
native_authority: None,
directory,
operations,
genesis: original.clone(),
dependencies: Vec::new(),
state,
partial_trees: Vec::new(),
authority_admissions: BTreeMap::new(),
ancestry: Vec::new(),
});
}
let mut geneses = BTreeMap::from([(selected_thread, genesis)]);
let mut dependencies = Vec::new();
for wrapper in dependency_records {
let record = wrapper
.genesis
.as_ref()
.ok_or(Error::Invalid("dependency signed genesis absent"))?;
let candidate = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
&record.canonical_record,
)
.map_err(preparation)?;
let reference = ThreadRef {
spool: thread.spool.clone(),
id: Some(ThreadId {
value: candidate.id().map_err(preparation)?.as_bytes().to_vec(),
}),
};
let candidate = super::verify_origin(&wrapper, &reference)?;
let id = candidate.id().map_err(preparation)?;
if geneses.len() >= 128 || geneses.insert(id, candidate).is_some() {
return Err(Error::Invalid(
"duplicate or oversized dependency genesis set",
));
}
dependencies.push(wrapper);
}
let mut claim_frontiers = BTreeMap::new();
for wrapper in std::iter::once(original).chain(&dependencies) {
let signed = wrapper
.genesis
.as_ref()
.ok_or(Error::Invalid("claim genesis absent"))?;
let genesis = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
&signed.canonical_record,
)
.map_err(preparation)?;
let mut frontier = BTreeSet::new();
let claims = crate::replication::ownership::verify_claims(wrapper, &genesis)?;
let resolutions = crate::replication::ownership::verify_resolutions(wrapper, &genesis)?;
if claims.is_empty() && resolutions.is_empty() {
continue;
}
for claim in claims {
frontier.extend(
claim
.original
.verify()
.map_err(preparation)?
.source_frontier,
);
}
for resolution in resolutions {
frontier.extend(
heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&resolution.original.canonical)
.map_err(preparation)?.frontier,
);
}
claim_frontiers.insert(genesis.id().map_err(preparation)?, frontier);
}
let mut originals = BTreeMap::new();
let mut decoded = BTreeMap::<ContentHash, ThreadOperation>::new();
let mut selected_operation = None;
let mut source_thread = selected_thread;
let mut inherited_bases = BTreeSet::new();
for _ in 0..128 {
let current = geneses
.get(&source_thread)
.ok_or(Error::Invalid("fork base source genesis absent"))?;
if selected.value.as_slice() != current.base.as_bytes() {
break;
}
let parent = current.parent.ok_or(Error::Invalid(
"non-system base has no original parent source",
))?;
if !inherited_bases.insert(source_thread) {
return Err(Error::Invalid("fork base parent cycle"));
}
let ancestor = geneses
.get(&parent)
.ok_or(Error::Invalid("fork base parent original absent"))?;
if ancestor.spool != current.spool {
return Err(Error::Invalid("fork base crosses Spool"));
}
source_thread = parent;
}
if selected.value.as_slice()
== geneses
.get(&source_thread)
.ok_or(Error::Invalid("fork base source genesis absent"))?
.base
.as_bytes()
{
return Err(Error::Invalid("fork base source chain exceeds bound"));
}
for signed in &operations {
let operation = signed.verify().map_err(preparation)?;
let id = operation.id().map_err(preparation)?;
let state = operation
.source_state()
.map_err(preparation)?
.ok_or(Error::Invalid("non-source operation in source ancestry"))?;
if operation.thread == source_thread
&& state.id().as_bytes().as_slice() == selected.value
&& selected_operation.replace((id, state)).is_some()
{
return Err(Error::Invalid("ambiguous selected source proof"));
}
originals.insert(id, signed.clone());
if decoded.insert(id, operation).is_some() {
return Err(Error::Invalid("duplicate source proof"));
}
}
let mut authority_admissions = BTreeMap::new();
for receipt in receipt_records {
let statement = receipt.verify_signature().map_err(preparation)?;
let operation_id = statement.subject.operation_id().ok_or(Error::Invalid(
"source batch cannot carry ownership claim admission",
))?;
let original = originals
.get(&operation_id)
.ok_or(Error::Invalid("unmatched source authority receipt"))?;
receipt.verify(original, &heddle_object_model::object::thread_replication::integration::TrustedHostedExecutor {
spool: statement.spool, spool_genesis: statement.spool_genesis, executor: statement.executor,
}).map_err(preparation)?;
if authority_admissions.insert(operation_id, receipt).is_some() {
return Err(Error::Invalid("duplicate source authority receipt"));
}
}
let (selected_id, state, selected_in_floor) = match selected_operation {
Some((id, state)) => (id, state, false),
None => {
let selected_state = StateId::from_bytes(
selected
.value
.as_slice()
.try_into()
.map_err(|_| Error::Invalid("selected revision identity width"))?,
);
let (tip, canonical) = ancestry
.selected_page_tip(selected_state, directory.path())?
.ok_or(Error::Invalid("selected source proof absent"))?;
let owner = decoded
.iter()
.find(|(_, operation)| {
operation.thread == source_thread
&& operation
.source_state()
.ok()
.flatten()
.is_some_and(|state| state.id() == tip)
})
.map(|(id, _)| *id)
.ok_or(Error::Invalid(
"import ancestry tip is not a carried source operation",
))?;
let state = State::decode_current_msgpack(&canonical)
.map_err(|_| Error::Invalid("selected import ancestor is not canonical"))?;
if state.id() != selected_state {
return Err(Error::Invalid(
"selected import ancestor differs from its address",
));
}
(owner, state, true)
}
};
let mut import_floors = Vec::new();
let mut pending = BTreeSet::from([selected_id]);
let mut seen = BTreeSet::new();
let mut used_threads = BTreeSet::new();
let mut claim_threads = BTreeSet::new();
let mut foreign_endpoints = BTreeSet::new();
let mut edges = BTreeMap::new();
while let Some(id) = pending.pop_first() {
if !seen.insert(id) {
continue;
}
let operation = decoded
.get(&id)
.ok_or(Error::Invalid("incomplete source ancestry"))?;
let signed = originals
.get(&id)
.ok_or(Error::Invalid("original absent"))?;
if foreign_endpoint(&foreign, operation, signed)? {
used_threads.insert(operation.thread);
foreign_endpoints.insert(id);
edges.insert(id, BTreeSet::new());
continue;
}
let parents = operation
.parents
.iter()
.map(|id| {
decoded
.get(id)
.cloned()
.ok_or(Error::Invalid("incomplete source ancestry"))
})
.collect::<Result<Vec<_>, _>>()?;
let genesis = geneses
.get(&operation.thread)
.ok_or(Error::Invalid("source dependency genesis absent"))?;
used_threads.insert(operation.thread);
if claim_threads.insert(operation.thread)
&& let Some(frontier) = claim_frontiers.get(&operation.thread)
{
for head in frontier {
if decoded
.get(head)
.is_none_or(|source| source.thread != operation.thread)
{
return Err(Error::Invalid(
"ownership claim cutoff source proof absent or foreign",
));
}
}
pending.extend(frontier);
}
let imported = carriers
.as_ref()
.map(|c| c.bind(genesis, operation, &parents))
.transpose()
.map_err(preparation)?
.flatten();
match &imported {
Some(bound) => bound.validate_parents(genesis, &parents),
None => operation.validate_parents(genesis, &parents),
}
.map_err(preparation)?;
if let Some(bound) = imported {
let mut frontier = BTreeSet::new();
for parent in &parents {
if let Some(state) = parent.source_state().map_err(preparation)? {
frontier.insert(state.id());
}
}
import_floors.push(ancestry::ImportFloorInput {
digest: api::import_authority::signed_operation_digest(bound.signed())
.map_err(preparation)?,
tip: operation
.source_state()
.map_err(preparation)?
.ok_or(Error::Invalid("import operation has no source State"))?,
frontier,
});
}
let mut required = operation.parents.clone();
if let Some(receipt) = operation.local_integration().map_err(preparation)? {
let source = decoded
.get(&receipt.source_operation)
.ok_or(Error::Invalid(
"local integration original source proof absent",
))?;
receipt.validate_source(source).map_err(preparation)?;
required.insert(receipt.source_operation);
pending.insert(receipt.source_operation);
}
if let Some(receipt) = operation.integration().map_err(preparation)? {
let source = decoded
.get(&receipt.source_operation)
.ok_or(Error::Invalid(
"hosted integration original source proof absent",
))?;
receipt.validate_source(source).map_err(preparation)?;
required.insert(receipt.source_operation);
pending.insert(receipt.source_operation);
}
edges.insert(id, required);
pending.extend(
operation
.parents
.iter()
.filter(|id| !seen.contains(id))
.copied(),
);
}
if seen.len() != decoded.len()
|| used_threads
.union(&inherited_bases)
.copied()
.collect::<BTreeSet<_>>()
!= geneses.keys().copied().collect()
{
return Err(Error::Invalid("unselected source proofs"));
}
for (thread, frontier) in &claim_frontiers {
if !claim_threads.contains(thread) {
continue;
}
let mut history = BTreeSet::new();
let mut pending = frontier.clone();
while let Some(id) = pending.pop_first() {
if !history.insert(id) {
continue;
}
let operation = decoded
.get(&id)
.ok_or(Error::Invalid("claim cutoff ancestry absent"))?;
if operation.thread != *thread {
return Err(Error::Invalid("claim cutoff crosses Thread"));
}
if !foreign_endpoints.contains(&id) {
pending.extend(&operation.parents);
}
}
{
for (id, operation) in &decoded {
if operation.thread == *thread
&& !history.contains(id)
&& !foreign_endpoints.contains(id)
{
if matches!(
operation.source_author().map_err(preparation)?,
Some(
heddle_object_model::object::thread_replication::SourceAuthor::LocalKey
)
) {
return Err(Error::Invalid(
"new local source lies outside signed ownership cutoff",
));
}
edges
.get_mut(id)
.ok_or(Error::Invalid("source topology entry absent"))?
.extend(frontier);
}
}
}
}
let references = decoded
.values()
.map(|operation| {
operation
.reference_proof(
geneses
.get(&operation.thread)
.ok_or(Error::Invalid("dependency genesis absent"))?,
)
.map_err(preparation)
})
.collect::<Result<Vec<_>, _>>()?
.into_iter()
.flatten()
.collect::<Vec<_>>();
let capture = decoded
.get(&selected_id)
.ok_or(Error::Invalid("selected source operation absent"))?
.source_result()
.map_err(preparation)?
.ok_or(Error::Invalid("selected operation has no source result"))?;
let verified_ancestry = ancestry::verify(
&ancestry,
&import_floors,
thread,
selected_in_floor.then_some(state.id()),
require_import_ancestry,
)?;
let pack = PackReader::open(
&directory.path().join("source.pack"),
&directory.path().join("source.idx"),
)
.map_err(preparation)?;
let (references, visibility) = if selected_in_floor {
(Vec::new(), None)
} else {
(references, capture.visibility.as_ref())
};
let partial_trees = if allow_partial {
pack.validate_visible_source_closure(&state, SOURCE_OBJECTS, SOURCE_BYTES)
.map_err(preparation)?
.partial_trees
} else {
pack.validate_source_closure_with_metadata(
&state,
&references,
visibility,
SOURCE_OBJECTS,
SOURCE_BYTES,
)
.map_err(preparation)?;
Vec::new()
};
let mut ready_ids: BTreeSet<_> = edges
.iter()
.filter(|(_, parents)| parents.is_empty())
.map(|(id, _)| *id)
.collect();
let mut children: BTreeMap<ContentHash, Vec<ContentHash>> = BTreeMap::new();
for (child, parents) in &edges {
for parent in parents {
children.entry(*parent).or_default().push(*child);
}
}
let mut ordered = Vec::new();
while let Some(id) = ready_ids.pop_first() {
ordered.push(
originals
.remove(&id)
.ok_or(Error::Invalid("duplicate source topology identity"))?,
);
if let Some(dependants) = children.get(&id) {
for child in dependants {
let parents = edges
.get_mut(child)
.ok_or(Error::Invalid("incomplete source topology"))?;
parents.remove(&id);
if parents.is_empty() {
ready_ids.insert(*child);
}
}
}
}
if !originals.is_empty() {
return Err(Error::Invalid("source dependency cycle"));
}
let mut genesis_record = original.clone();
omit_unusable_ownership(&mut genesis_record, &decoded)?;
for dependency in &mut dependencies {
omit_unusable_ownership(dependency, &decoded)?;
}
Ok(ValidatedSourceArtifacts {
import_authority: None,
native_authority: None,
directory,
genesis: genesis_record,
operations: ordered,
authority_admissions,
dependencies,
state,
partial_trees,
ancestry: verified_ancestry.floors,
})
}
fn source_frontier_is_installable(
thread: ContentHash,
frontier: &BTreeSet<ContentHash>,
operations: &BTreeMap<ContentHash, ThreadOperation>,
) -> bool {
frontier.iter().all(|id| {
operations.get(id).is_some_and(|operation| {
operation.thread == thread && operation.source_state().ok().flatten().is_some()
})
})
}
fn omit_unusable_ownership(
wrapper: &mut ThreadGenesisRecord,
operations: &BTreeMap<ContentHash, ThreadOperation>,
) -> Result<(), Error> {
if wrapper.ownership_claims.is_empty() && wrapper.ownership_resolutions.is_empty() {
return Ok(());
}
let signed = wrapper
.genesis
.as_ref()
.ok_or(Error::Invalid("claim genesis absent"))?;
let genesis = heddle_object_model::object::thread_replication::ThreadGenesis::decode(
&signed.canonical_record,
)
.map_err(preparation)?;
let thread = genesis.id().map_err(preparation)?;
let mut kept_claims = BTreeSet::new();
let mut claims = Vec::new();
for record in wrapper.ownership_claims.drain(..) {
let crate::thread_ownership::ClaimProof::Complete(proof) =
crate::thread_ownership::decode(&record).map_err(preparation)?
else {
return Err(Error::Invalid(
"transferred claim requires both original signatures",
));
};
let claim = proof.verify().map_err(preparation)?;
if source_frontier_is_installable(thread, &claim.source_frontier, operations) {
kept_claims.insert(claim.id().map_err(preparation)?);
claims.push(record);
}
}
wrapper.ownership_claims = claims;
let mut claim_admissions = Vec::new();
for record in wrapper.ownership_claim_admissions.drain(..) {
let statement = admission_subject(&record)?;
let Some(id) = statement.subject.claim_id() else {
return Err(Error::Invalid("ownership claim admission subject"));
};
if kept_claims.contains(&id) {
claim_admissions.push(record);
}
}
wrapper.ownership_claim_admissions = claim_admissions;
let mut kept_resolutions = BTreeSet::new();
let mut resolutions = Vec::new();
for record in wrapper.ownership_resolutions.drain(..) {
let proof = crate::thread_ownership::decode_resolution(&record).map_err(preparation)?;
let resolution = heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&proof.canonical).map_err(preparation)?;
let claims_present = kept_claims.contains(&resolution.winning_claim)
&& resolution
.conflicting_claims
.iter()
.all(|id| kept_claims.contains(id));
if claims_present
&& source_frontier_is_installable(thread, &resolution.frontier, operations)
{
kept_resolutions.insert(resolution.id().map_err(preparation)?);
resolutions.push(record);
}
}
wrapper.ownership_resolutions = resolutions;
let mut resolution_admissions = Vec::new();
for record in wrapper.ownership_resolution_admissions.drain(..) {
let statement = admission_subject(&record)?;
if matches!(
statement.subject,
heddle_object_model::object::thread_authority_admission::OriginalAuthoritySubject::OwnershipResolution(id)
if kept_resolutions.contains(&id)
) {
resolution_admissions.push(record);
}
}
wrapper.ownership_resolution_admissions = resolution_admissions;
Ok(())
}
fn admission_subject(
record: &SignedRecord,
) -> Result<heddle_object_model::object::thread_authority_admission::ThreadAuthorityAdmission, Error>
{
crate::authority_admission::decode(record)
.map_err(preparation)?
.verify_signature()
.map_err(preparation)
}
fn foreign_endpoint(
foreign: &[ForeignDependencyV1],
operation: &ThreadOperation,
signed: &SignedOperation,
) -> Result<bool, Error> {
if !foreign
.iter()
.any(|r| r.thread_genesis_digest.as_slice() == operation.thread.as_bytes())
{
return Ok(false);
}
let original = SignedRecord {
format: heddle_object_model::object::thread_replication::OPERATION_FORMAT.into(),
canonical_record: signed.canonical.clone(),
signatures: vec![RecordSignature {
public_key: operation.publisher.to_vec(),
signature: signed.signature.clone(),
}],
};
let digest = api::import_authority::signed_native_digest(&original)?;
Ok(foreign.iter().any(|r| {
r.thread_genesis_digest.as_slice() == operation.thread.as_bytes()
&& r.signed_native_digest == digest
}))
}
fn preparation(error: impl std::fmt::Display) -> Error {
Error::Preparation(error.to_string())
}
#[cfg(test)]
#[path = "staging_tests.rs"]
mod tests;