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