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