Skip to main content

heddle_thread_api/
fetch.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Native source downloads retain original Thread identities and causal proofs.
3//! Pack chunks are staging bytes: install only after the verified Complete frame.
4#[cfg(feature = "native")]
5pub mod hosted;
6#[cfg(feature = "native")]
7mod native;
8#[cfg(feature = "native")]
9pub use native::OwnedDeviceBinding;
10mod provider;
11pub use provider::{
12    Candidate as ProviderCandidate, ProviderConsentSigner, ProviderDownload, ProviderFetch,
13    ProviderPlanSession,
14};
15mod staging;
16use api::v2::client::{ClientError, MessageReader, Messages, RpcTransport};
17use prost::Message;
18pub(crate) use staging::validate_artifacts;
19#[cfg(test)]
20pub(crate) use staging::validate_with_receipts;
21pub use staging::{StagedSource, ValidatedSourceArtifacts};
22
23use crate::{Remote, contract::*, replication, rpc, transport};
24
25#[derive(Debug, thiserror::Error)]
26pub enum Error {
27    #[cfg(feature = "native")]
28    #[error("foreign prefix {limit_name} exceeds limit {limit}")]
29    ForeignPrefixLimitExceeded {
30        limit_name: &'static str,
31        limit: usize,
32    },
33    #[error(transparent)]
34    Client(#[from] ClientError<transport::Error>),
35    #[error(transparent)]
36    Transport(#[from] transport::Error),
37    #[error(transparent)]
38    Replication(#[from] replication::Error),
39    #[error("source download I/O: {0}")]
40    Io(#[from] std::io::Error),
41    #[error("source preparation: {0}")]
42    Preparation(String),
43    #[error("invalid source download: {0}")]
44    Invalid(&'static str),
45    #[error("hosted history requires independently selected root, owner and fresh witness trust")]
46    HostedTrustRequired,
47    #[error("HYBRID authority rejected: {0}")]
48    Hybrid(#[from] api::hybrid_codec::Reject),
49}
50
51impl crate::reopen::ReopenRetryable for Error {
52    fn is_reopen_retryable(&self) -> bool {
53        match self {
54            Error::Client(error) => crate::reopen::client_error_is_reopen_retryable(error),
55            Error::Transport(error) => crate::reopen::error_is_reopen_retryable(error),
56            _ => false,
57        }
58    }
59}
60
61#[cfg(feature = "native")]
62impl From<repo::thread_replication::Error> for Error {
63    fn from(error: repo::thread_replication::Error) -> Self {
64        match error {
65            repo::thread_replication::Error::ForeignPrefixLimitExceeded { limit_name, limit } => {
66                Self::ForeignPrefixLimitExceeded { limit_name, limit }
67            }
68            error => Self::Preparation(error.to_string()),
69        }
70    }
71}
72
73#[derive(Clone, Copy)]
74pub struct Limits {
75    pub max_artifact_bytes: u64,
76    pub max_total_bytes: u64,
77    pub max_operations: usize,
78}
79impl Default for Limits {
80    fn default() -> Self {
81        Self {
82            max_artifact_bytes: 512 * 1024 * 1024,
83            max_total_bytes: 1024 * 1024 * 1024,
84            max_operations: 100_000,
85        }
86    }
87}
88
89/// Each item has passed its frame, scope and cryptographic checks. An operation
90/// can still have missing causal parents; the durable replica decides admission.
91#[allow(clippy::large_enum_variant)] // heddle-api's inline import-authority bundle; boxing adds a heap hop per frame
92pub enum Item {
93    Pack(PackChunk),
94    Operations(ReplicationOperations),
95    ThreadGenesis(ThreadGenesisRecord),
96    Sidecar(TransferSidecar),
97    Complete(FetchComplete),
98}
99
100pub struct Download<R: MessageReader<Error = transport::Error>> {
101    messages: Messages<R, FetchServerFrame>,
102    state: Validation,
103}
104impl<R: MessageReader<Error = transport::Error>> Download<R> {
105    pub fn ready(&self) -> &TransferReady {
106        &self.state.ready
107    }
108    pub async fn next(&mut self) -> Result<Option<Item>, Error> {
109        if self.state.done {
110            return Ok(None);
111        }
112        let frame = self
113            .messages
114            .next()
115            .await?
116            .ok_or(Error::Invalid("stream ended before Complete"))?;
117        let item = match self.state.accept(frame) {
118            Ok(item) => item,
119            Err(error) => {
120                self.messages.cancel();
121                self.state.done = true;
122                return Err(error);
123            }
124        };
125        if self.state.done {
126            self.messages.cancel();
127        }
128        Ok(Some(item))
129    }
130}
131
132impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
133    /// A single exact Thread/revision request returns clone admission, immutable
134    /// operations, and bounded source artifacts. No reference update is implied.
135    pub async fn fetch_content(
136        &self,
137        open: FetchOpen,
138        limits: Limits,
139    ) -> Result<Download<T::Reader>, Error> {
140        if open.delivery == fetch_open::Delivery::ProviderPreferred as i32 {
141            return Err(Error::Invalid(
142                "provider delivery requires negotiated Fetch",
143            ));
144        }
145        if open.checkpoint.is_some() {
146            return Err(Error::Invalid(
147                "a fresh download requires an empty transfer checkpoint",
148            ));
149        }
150        crate::reopen::retry(|| self.fetch_content_once(open.clone(), limits)).await
151    }
152
153    async fn fetch_content_once(
154        &self,
155        open: FetchOpen,
156        limits: Limits,
157    ) -> Result<Download<T::Reader>, Error> {
158        crate::hybrid::fetch_open(&open).map_err(Error::Invalid)?;
159        if open.protocol.is_some() {
160            api::import_authority::require_hybrid_peer(self.description.protocol.as_ref())?;
161        }
162        let (sender, mut messages) = self
163            .api
164            .exchange::<rpc::SyncServiceFetch>(&FetchClientFrame {
165                body: Some(fetch_client_frame::Body::Open(open.clone())),
166            })
167            .await?;
168        // This direct-hosted transfer needs no provider negotiation. Closing the
169        // request half leaves response flow control and cancellation independent.
170        sender.finish().await?;
171        let frame = messages
172            .next()
173            .await?
174            .ok_or(Error::Invalid("Ready required"))?;
175        let Some(fetch_server_frame::Body::Ready(ready)) = frame.body else {
176            return Err(Error::Invalid("first response must be Ready"));
177        };
178        let state = Validation::new(open, ready, self.description.endpoint.as_ref(), limits)?;
179        Ok(Download { messages, state })
180    }
181}
182
183struct Validation {
184    ready: TransferReady,
185    facets: Vec<i32>,
186    frame_bytes: usize,
187    limits: Limits,
188    artifact: usize,
189    offset: u64,
190    digest: blake3::Hasher,
191    received: u64,
192    metadata_bytes: u64,
193    operations: usize,
194    threads: std::collections::BTreeSet<heddle_object_model::object::ContentHash>,
195    done: bool,
196}
197impl Validation {
198    fn new(
199        open: FetchOpen,
200        ready: TransferReady,
201        endpoint: Option<&EndpointRef>,
202        limits: Limits,
203    ) -> Result<Self, Error> {
204        crate::hybrid::transfer_ready(&ready).map_err(Error::Invalid)?;
205        crate::hybrid::negotiated(open.protocol.as_ref(), ready.protocol.as_ref())
206            .map_err(Error::Invalid)?;
207        let thread = open
208            .thread
209            .as_ref()
210            .ok_or(Error::Invalid("Thread required"))?;
211        if ready.endpoint.as_ref() != endpoint
212            || endpoint.is_none()
213            || ready.thread.as_ref() != Some(thread)
214            || ready
215                .current
216                .as_ref()
217                .is_none_or(|r| r.spool != thread.spool)
218            || open
219                .revision
220                .as_ref()
221                .is_some_and(|r| Some(r) != ready.current.as_ref())
222        {
223            return Err(Error::Invalid(
224                "admission does not match requested endpoint and revision",
225            ));
226        }
227        let genesis = ready
228            .thread_genesis
229            .as_ref()
230            .ok_or(Error::Invalid("original Thread genesis required"))?;
231        let verified_genesis = verify_origin(genesis, thread)?;
232        let spool = thread
233            .spool
234            .as_ref()
235            .ok_or(Error::Invalid("spool required"))?;
236        let id =
237            uuid::Uuid::parse_str(&spool.id).map_err(|_| Error::Invalid("invalid spool UUID"))?;
238        if endpoint.is_some_and(|endpoint| endpoint.kind == EndpointKind::Weft as i32) {
239            if ready
240                .owner_genesis
241                .as_ref()
242                .and_then(|g| g.genesis.as_ref())
243                .is_none_or(|g| g.spool_uuid != id.as_bytes())
244            {
245                return Err(Error::Invalid(
246                    "original owner genesis must bind this spool",
247                ));
248            }
249            if ready.ownership.is_none() {
250                return Err(Error::Invalid("portable owner history required"));
251            }
252        } else if endpoint.is_none_or(|endpoint| endpoint.kind != EndpointKind::Device as i32) {
253            return Err(Error::Invalid("source endpoint must be Weft or Heddle"));
254        }
255        let budget = ready
256            .budget
257            .ok_or(Error::Invalid("download budget required"))?;
258        if !(1024..=512 * 1024).contains(&budget.max_frame_bytes) || limits.max_operations == 0 {
259            return Err(Error::Invalid("invalid download limits"));
260        }
261        if ready.encoded_len() > budget.max_frame_bytes as usize {
262            return Err(Error::Invalid("admission exceeds frame budget"));
263        }
264        let provider = open.delivery == fetch_open::Delivery::ProviderPreferred as i32;
265        if (provider && (!ready.packs.is_empty() || !ready.full_closure_available))
266            || (!provider
267                && (ready.packs.len() != 2
268                    || ready.packs[0].kind != pack_extent::Kind::NativePack as i32
269                    || ready.packs[1].kind != pack_extent::Kind::NativeIndex as i32))
270        {
271            return Err(Error::Invalid("ordered native pack and index required"));
272        }
273        let mut total = 0_u64;
274        for extent in &ready.packs {
275            let address = extent
276                .pack
277                .as_ref()
278                .ok_or(Error::Invalid("artifact address required"))?;
279            if address.algorithm != "blake3"
280                || address.digest.len() != 32
281                || extent.offset != 0
282                || extent.length == 0
283                || extent.length > limits.max_artifact_bytes
284                || extent.extent_digest.as_ref() != Some(address)
285            {
286                return Err(Error::Invalid("invalid whole artifact extent"));
287            }
288            total = total
289                .checked_add(extent.length)
290                .ok_or(Error::Invalid("artifact size overflow"))?;
291        }
292        if total > limits.max_total_bytes {
293            return Err(Error::Invalid("download exceeds source budget"));
294        }
295        let checkpoint = ready
296            .checkpoint
297            .as_ref()
298            .ok_or(Error::Invalid("transfer checkpoint required"))?;
299        if checkpoint.transfer_id.is_empty()
300            || checkpoint.plan_digest.len() != 32
301            || checkpoint.committed_bytes != 0
302        {
303            return Err(Error::Invalid("invalid fresh transfer checkpoint"));
304        }
305        if !ready.full_closure_available
306            && !open.selection.as_ref().is_some_and(|s| s.allow_partial)
307        {
308            return Err(Error::Invalid(
309                "partial source disclosure was not requested",
310            ));
311        }
312        let facets = open.selection.map(|s| s.facets).unwrap_or_default();
313        if !facets.contains(&(SharedFacet::Source as i32))
314            || facets.iter().any(|f| {
315                !matches!(
316                    SharedFacet::try_from(*f),
317                    Ok(SharedFacet::Source | SharedFacet::Collaboration)
318                )
319            })
320        {
321            return Err(Error::Invalid(
322                "explicit supported source/discussion facets required",
323            ));
324        }
325        Ok(Self {
326            ready,
327            facets,
328            frame_bytes: budget.max_frame_bytes as usize,
329            limits,
330            artifact: 0,
331            offset: 0,
332            digest: blake3::Hasher::new(),
333            received: 0,
334            metadata_bytes: 0,
335            operations: 0,
336            threads: std::collections::BTreeSet::from([verified_genesis
337                .id()
338                .map_err(|_| Error::Invalid("invalid Thread genesis identity"))?]),
339            done: false,
340        })
341    }
342    fn accept(&mut self, frame: FetchServerFrame) -> Result<Item, Error> {
343        if self.done || frame.encoded_len() > self.frame_bytes {
344            return Err(Error::Invalid("frame exceeds active download budget"));
345        }
346        let body = frame.body.ok_or(Error::Invalid("empty download frame"))?;
347        match body {
348            fetch_server_frame::Body::Pack(chunk) => {
349                let expected = self
350                    .ready
351                    .packs
352                    .get(self.artifact)
353                    .ok_or(Error::Invalid("unexpected artifact"))?;
354                let extent = chunk
355                    .extent
356                    .as_ref()
357                    .ok_or(Error::Invalid("chunk extent required"))?;
358                let digest = ObjectAddress {
359                    algorithm: "blake3".into(),
360                    digest: blake3::hash(&chunk.data).as_bytes().to_vec(),
361                };
362                if chunk.data.is_empty()
363                    || extent.pack != expected.pack
364                    || extent.kind != expected.kind
365                    || extent.offset != self.offset
366                    || extent.length != chunk.data.len() as u64
367                    || extent.length > expected.length.saturating_sub(self.offset)
368                    || extent.extent_digest.as_ref() != Some(&digest)
369                {
370                    return Err(Error::Invalid(
371                        "chunk does not match the declared artifact extent",
372                    ));
373                }
374                if extent.length
375                    > self
376                        .limits
377                        .max_total_bytes
378                        .saturating_sub(self.received)
379                        .saturating_sub(self.metadata_bytes)
380                {
381                    return Err(Error::Invalid(
382                        "source and metadata exceed shared download budget",
383                    ));
384                }
385                self.digest.update(&chunk.data);
386                self.offset += extent.length;
387                self.received += extent.length;
388                if self.offset == expected.length {
389                    if expected
390                        .pack
391                        .as_ref()
392                        .is_none_or(|p| p.digest != self.digest.finalize().as_bytes())
393                    {
394                        return Err(Error::Invalid("whole artifact hash mismatch"));
395                    }
396                    self.artifact += 1;
397                    self.offset = 0;
398                    self.digest = blake3::Hasher::new();
399                }
400                Ok(Item::Pack(chunk))
401            }
402            fetch_server_frame::Body::Operations(batch) => {
403                crate::hybrid::operations(&batch).map_err(Error::Invalid)?;
404                if batch.native_authority.is_some()
405                    && batch.native_authority != self.ready.native_authority
406                {
407                    return Err(Error::Invalid(
408                        "native authority differs from transfer Ready",
409                    ));
410                }
411                if batch.import_authority.is_some()
412                    && batch.import_authority != self.ready.import_authority
413                {
414                    return Err(Error::Invalid(
415                        "operation proof bundle differs from negotiated source closure",
416                    ));
417                }
418                self.operations = self
419                    .operations
420                    .checked_add(batch.operations.len())
421                    .ok_or(Error::Invalid("operation count overflow"))?;
422                self.metadata_bytes = self
423                    .metadata_bytes
424                    .checked_add(batch.encoded_len() as u64)
425                    .ok_or(Error::Invalid("metadata size overflow"))?;
426                if self.operations > self.limits.max_operations
427                    || self.metadata_bytes
428                        > self.limits.max_total_bytes.saturating_sub(self.received)
429                {
430                    return Err(Error::Invalid("causal metadata exceeds download budget"));
431                }
432                for received in crate::authority_admission::match_batch(&batch)? {
433                    let operation = received
434                        .original
435                        .verify()
436                        .map_err(|_| Error::Invalid("invalid original operation signature"))?;
437                    if !self.threads.contains(&operation.thread)
438                        || !self
439                            .facets
440                            .contains(&replication::wire_facet(operation.facet()))
441                    {
442                        return Err(Error::Invalid("operation crosses Thread or selected facet"));
443                    }
444                }
445                Ok(Item::Operations(batch))
446            }
447            fetch_server_frame::Body::ThreadGenesis(record) => {
448                self.metadata_bytes = self
449                    .metadata_bytes
450                    .checked_add(record.encoded_len() as u64)
451                    .ok_or(Error::Invalid("metadata size overflow"))?;
452                if self.threads.len() >= 128
453                    || self.metadata_bytes
454                        > self.limits.max_total_bytes.saturating_sub(self.received)
455                {
456                    return Err(Error::Invalid("dependency genesis exceeds download budget"));
457                }
458                let genesis =
459                    heddle_object_model::object::thread_replication::ThreadGenesis::decode(
460                        &record
461                            .genesis
462                            .as_ref()
463                            .ok_or(Error::Invalid("dependency signed genesis absent"))?
464                            .canonical_record,
465                    )
466                    .map_err(|_| Error::Invalid("invalid dependency genesis"))?;
467                let spool = self
468                    .ready
469                    .thread
470                    .as_ref()
471                    .and_then(|thread| thread.spool.as_ref())
472                    .ok_or(Error::Invalid("Spool absent"))?;
473                if genesis.spool != spool.id {
474                    return Err(Error::Invalid("dependency genesis crosses Spool"));
475                }
476                let id = genesis
477                    .id()
478                    .map_err(|_| Error::Invalid("invalid dependency identity"))?;
479                let thread = ThreadRef {
480                    spool: Some(spool.clone()),
481                    id: Some(ThreadId {
482                        value: id.as_bytes().to_vec(),
483                    }),
484                };
485                verify_origin(&record, &thread)?;
486                if !self.threads.insert(id) {
487                    return Err(Error::Invalid("duplicate dependency genesis"));
488                }
489                Ok(Item::ThreadGenesis(record))
490            }
491            fetch_server_frame::Body::Complete(complete) => {
492                let original = self
493                    .ready
494                    .checkpoint
495                    .as_ref()
496                    .ok_or(Error::Invalid("checkpoint absent"))?;
497                let checkpoint = complete
498                    .checkpoint
499                    .as_ref()
500                    .ok_or(Error::Invalid("final checkpoint required"))?;
501                if complete.revision != self.ready.current
502                    || self.artifact != self.ready.packs.len()
503                    || checkpoint.transfer_id != original.transfer_id
504                    || checkpoint.plan_digest != original.plan_digest
505                    || checkpoint.committed_bytes != self.received
506                    || complete.closure
507                        != if self.ready.full_closure_available {
508                            Coverage::Complete as i32
509                        } else {
510                            Coverage::Partial as i32
511                        }
512                    || !complete.missing.is_empty()
513                {
514                    return Err(Error::Invalid(
515                        "download does not match its exact declared source coverage",
516                    ));
517                }
518                self.done = true;
519                Ok(Item::Complete(complete))
520            }
521            fetch_server_frame::Body::Sidecar(_) => Err(Error::Invalid(
522                "sidecar facet was not explicitly negotiated",
523            )),
524            fetch_server_frame::Body::ProviderPlan(_) => Err(Error::Invalid(
525                "provider transfer requires explicit client consent",
526            )),
527            fetch_server_frame::Body::ProviderInline(_) => Err(Error::Invalid(
528                "provider inline record requires an admitted provider plan",
529            )),
530            fetch_server_frame::Body::ProviderOffer(_) => Err(Error::Invalid(
531                "provider offer requires explicit client delivery negotiation",
532            )),
533            fetch_server_frame::Body::Ready(_) => {
534                Err(Error::Invalid("duplicate download admission"))
535            }
536        }
537    }
538}
539
540#[cfg(test)]
541#[path = "fetch_tests.rs"]
542mod tests;
543
544/// Verify original signatures and exact receipt bindings without deriving trust
545/// from the carried executor key. Installation pins executor authority separately.
546pub(crate) fn verify_origin(
547    record: &ThreadGenesisRecord,
548    thread: &ThreadRef,
549) -> Result<heddle_object_model::object::thread_replication::ThreadGenesis, Error> {
550    use heddle_object_model::object::thread_replication::integration::TrustedHostedExecutor;
551    let genesis = replication::opening::verify_genesis_record(record, thread)?;
552    if let Some(signed) = crate::boundary_acceptance::genesis_admission(record)? {
553        let value = signed
554            .verify_signature()
555            .map_err(|_| Error::Invalid("invalid genesis receipt"))?;
556        let trust = TrustedHostedExecutor {
557            spool: value.spool,
558            spool_genesis: value.spool_genesis,
559            executor: value.executor,
560        };
561        let evidence = signed
562            .boundary_acceptance
563            .as_ref()
564            .map(|value| value.verify_signature())
565            .transpose()
566            .map_err(|_| Error::Invalid("invalid boundary evidence"))?;
567        value
568            .authorize_with_acceptance(
569                &genesis,
570                &record.creator_authority,
571                &trust,
572                evidence.as_ref(),
573            )
574            .map_err(|_| Error::Invalid("genesis admission differs from original proof"))?;
575    }
576    Ok(genesis)
577}