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