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