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
63pub 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 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 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
495pub(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}