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