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