1#[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#[allow(clippy::large_enum_variant)] pub 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 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 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
498pub(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}