1use std::{io::Read, path::Path};
7
8use api::{
9 heddle::api::v1alpha2::{
10 EndpointRef, FetchClientFrame, FetchOpen, FetchServerFrame, ProviderConsent, ProviderOffer,
11 ProviderPlan, ProviderPlanChallenge, ReadProviderExtentRequest, RecordSignature,
12 SignedRecord, fetch_client_frame, fetch_open, fetch_server_frame, provider_assembly_record,
13 provider_extent_event,
14 },
15 provider_v2::{
16 PROVIDER_CONSENT_FORMAT, provider_consent_signing_bytes, validate_plan_for_offer,
17 validate_provider_offer,
18 },
19 v2::client::{MessageReader, MessageWriter, Messages, RpcTransport, Sender},
20};
21use heddle_object_model::object::{ContentHash, StateId};
22use heddle_pack::store::pack::PackObjectId;
23use wire::{ProviderPackExtent, ProviderPackIndexEntry, ProviderPackManifest, ProviderPackSpool};
24
25use super::{Download, Error, Item, Limits, StagedSource, Validation};
26use crate::{Remote, contract::TransferReady, rpc, transport};
27
28pub trait ProviderConsentSigner {
32 fn verified_subject(&self) -> Result<String, Error>;
33 fn client_endpoint(&self) -> Result<EndpointRef, Error>;
34 fn public_key(&self) -> &[u8];
35 fn sign(&self, canonical: &[u8]) -> Result<Vec<u8>, Error>;
36}
37
38pub struct ProviderDownload<
41 W: MessageWriter<Error = transport::Error>,
42 R: MessageReader<Error = transport::Error>,
43> {
44 sender: Sender<W, FetchClientFrame>,
45 messages: Messages<R, FetchServerFrame>,
46 state: Validation,
47 open: FetchOpen,
48 issuer: EndpointRef,
49}
50
51pub enum ProviderFetch<
55 W: MessageWriter<Error = transport::Error>,
56 R: MessageReader<Error = transport::Error>,
57> {
58 Direct(Box<Download<R>>),
59 Provider(Box<ProviderDownload<W, R>>),
60}
61
62pub struct ProviderPlanSession<
64 W: MessageWriter<Error = transport::Error>,
65 R: MessageReader<Error = transport::Error>,
66> {
67 plan: ProviderPlan,
68 ready: TransferReady,
69 originals: Vec<Item>,
70 ancestry: super::staging::AncestryInput,
71 ancestry_directory: tempfile::TempDir,
72 sender: Sender<W, FetchClientFrame>,
73 messages: Messages<R, FetchServerFrame>,
74 state: Validation,
75 spool: Option<ProviderPackSpool>,
76}
77
78impl<W: MessageWriter<Error = transport::Error>, R: MessageReader<Error = transport::Error>>
79 ProviderPlanSession<W, R>
80{
81 pub fn plan(&self) -> &ProviderPlan {
83 &self.plan
84 }
85
86 pub fn ready(&self) -> &TransferReady {
88 &self.ready
89 }
90
91 pub async fn receive_inline(&mut self, scratch: &Path) -> Result<(), Error> {
96 if self.spool.is_some() {
97 return Err(Error::Invalid("provider inline delivery already started"));
98 }
99 if self.plan.output_pack_length > 256 * 1024 * 1024
100 || self.plan.output_pack_length > self.state.limits.max_artifact_bytes
101 || self.plan.output_pack_length
102 > self
103 .state
104 .limits
105 .max_total_bytes
106 .saturating_sub(self.state.metadata_bytes)
107 {
108 return Err(Error::Invalid(
109 "provider source pack exceeds staging budget",
110 ));
111 }
112 let manifest = pack_manifest(&self.plan)?;
113 let scratch = scratch.to_path_buf();
114 let spool =
115 tokio::task::spawn_blocking(move || ProviderPackSpool::new_in(&scratch, manifest))
116 .await
117 .map_err(|error| Error::Preparation(error.to_string()))?
118 .map_err(|error| Error::Preparation(error.to_string()))?;
119 let writer = spool.writer();
120 let mut offsets = vec![0_u64; self.plan.records.len()];
121 let inline_count = self
122 .plan
123 .records
124 .iter()
125 .filter(|record| {
126 matches!(
127 record.source,
128 Some(provider_assembly_record::Source::Inline(_))
129 )
130 })
131 .count();
132 let mut complete = 0_usize;
133 while complete < inline_count {
134 let frame = self.messages.next().await?.ok_or(Error::Invalid(
135 "provider inline records ended before complete coverage",
136 ))?;
137 let Some(fetch_server_frame::Body::ProviderInline(chunk)) = frame.body else {
138 return Err(Error::Invalid(
139 "unexpected frame during provider inline delivery",
140 ));
141 };
142 let index = usize::try_from(chunk.record_index)
143 .map_err(|_| Error::Invalid("provider inline record index"))?;
144 let record = self
145 .plan
146 .records
147 .get(index)
148 .ok_or(Error::Invalid("provider inline record absent"))?;
149 if !matches!(
150 record.source,
151 Some(provider_assembly_record::Source::Inline(_))
152 ) || chunk.assembly_digest != self.plan.assembly_digest
153 || chunk.data.is_empty()
154 || chunk.data.len() > 1024 * 1024
155 || chunk.offset != offsets[index]
156 {
157 return Err(Error::Invalid("provider inline chunk differs from plan"));
158 }
159 let next = chunk
160 .offset
161 .checked_add(chunk.data.len() as u64)
162 .ok_or(Error::Invalid("provider inline offset overflow"))?;
163 if next > record.encoded_length {
164 return Err(Error::Invalid("provider inline chunk exceeds record"));
165 }
166 let finished = next == record.encoded_length;
167 write_chunk(
168 writer.clone(),
169 index,
170 chunk.offset,
171 chunk.data,
172 record.clone(),
173 finished,
174 )
175 .await?;
176 offsets[index] = next;
177 if finished {
178 complete += 1;
179 }
180 }
181 self.spool = Some(spool);
182 Ok(())
183 }
184
185 pub async fn receive_provider_ranges<T: RpcTransport<Error = transport::Error>>(
189 &mut self,
190 providers: &[Remote<T>],
191 ) -> Result<(), Error> {
192 let spool = self
193 .spool
194 .as_ref()
195 .ok_or(Error::Invalid("provider inline stage required"))?;
196 let writer = spool.writer();
197 let mut grouped = vec![Vec::new(); self.plan.extents.len()];
198 for (index, record) in self.plan.records.iter().enumerate() {
199 if let Some(provider_assembly_record::Source::Provider(source)) = &record.source {
200 grouped
201 .get_mut(source.extent_index as usize)
202 .ok_or(Error::Invalid("provider record extent absent"))?
203 .push((source.source_offset, index, record));
204 }
205 }
206 for records in &mut grouped {
207 records.sort_by_key(|(offset, _, _)| *offset);
208 }
209 for (extent_index, extent) in self.plan.extents.iter().enumerate() {
210 let endpoint = extent
211 .provider
212 .as_ref()
213 .ok_or(Error::Invalid("provider endpoint absent"))?;
214 let remote = providers
215 .iter()
216 .find(|remote| remote.description.endpoint.as_ref() == Some(endpoint))
217 .ok_or(Error::Invalid("selected provider endpoint unavailable"))?;
218 let range = extent
219 .range
220 .as_ref()
221 .ok_or(Error::Invalid("provider range absent"))?;
222 let ticket = extent
223 .ticket
224 .as_ref()
225 .ok_or(Error::Invalid("provider ticket absent"))?;
226 let records = grouped
227 .get(extent_index)
228 .ok_or(Error::Invalid("provider extent group absent"))?;
229 if records.is_empty() || records[0].0 != 0 {
230 return Err(Error::Invalid("provider range has no tiled records"));
231 }
232 let mut messages = remote
233 .api
234 .observe::<rpc::SyncServiceReadProviderExtent>(&ReadProviderExtentRequest {
235 ticket: Some(ticket.clone()),
236 extent_set_digest: self.plan.extent_set_digest.clone(),
237 range: Some(range.clone()),
238 })
239 .await?;
240 let first = messages
241 .next()
242 .await?
243 .ok_or(Error::Invalid("provider extent Ready absent"))?;
244 if !matches!(first.body, Some(provider_extent_event::Body::Ready(ref ready)) if ready == range)
245 {
246 return Err(Error::Invalid("provider extent Ready differs from ticket"));
247 }
248 let mut offset = 0_u64;
249 let mut record_pos = 0_usize;
250 loop {
251 let event = messages
252 .next()
253 .await?
254 .ok_or(Error::Invalid("provider extent ended without Complete"))?;
255 match event.body {
256 Some(provider_extent_event::Body::Chunk(chunk)) => {
257 if chunk.offset != offset
258 || chunk.data.is_empty()
259 || chunk.data.len() > 1024 * 1024
260 {
261 return Err(Error::Invalid("provider extent chunk offset or size"));
262 }
263 let end = offset
264 .checked_add(chunk.data.len() as u64)
265 .ok_or(Error::Invalid("provider extent offset overflow"))?;
266 if end > range.length {
267 return Err(Error::Invalid("provider extent exceeds ticket"));
268 }
269 let mut used = 0_usize;
270 while used < chunk.data.len() {
271 let (start, index, record) = records
272 .get(record_pos)
273 .ok_or(Error::Invalid("provider extra bytes after records"))?;
274 let relative = offset
275 .checked_sub(*start)
276 .ok_or(Error::Invalid("provider record gap"))?;
277 let remaining = record
278 .encoded_length
279 .checked_sub(relative)
280 .ok_or(Error::Invalid("provider record overflow"))?;
281 let take =
282 usize::try_from(remaining.min((chunk.data.len() - used) as u64))
283 .map_err(|_| Error::Invalid("provider record size"))?;
284 let finished = relative + take as u64 == record.encoded_length;
285 write_chunk(
286 writer.clone(),
287 *index,
288 relative,
289 chunk.data[used..used + take].to_vec(),
290 (*record).clone(),
291 finished,
292 )
293 .await?;
294 offset += take as u64;
295 used += take;
296 if finished {
297 record_pos += 1;
298 }
299 }
300 }
301 Some(provider_extent_event::Body::Complete(checkpoint)) => {
302 if offset != range.length
303 || record_pos != records.len()
304 || checkpoint.committed_bytes != range.length
305 {
306 return Err(Error::Invalid("provider extent incomplete"));
307 }
308 break;
309 }
310 _ => return Err(Error::Invalid("unexpected provider extent frame")),
311 }
312 }
313 }
314 Ok(())
315 }
316
317 pub async fn complete(self, scratch: &Path) -> Result<StagedSource, Error> {
321 self.complete_inner(scratch, None).await
322 }
323 pub async fn complete_with_import_carriers(
324 self,
325 scratch: &Path,
326 carriers: crypto::import_authority::VerifiedImportCarriers,
327 ) -> Result<StagedSource, Error> {
328 self.complete_inner(scratch, Some(carriers)).await
329 }
330 async fn complete_inner(
331 mut self,
332 scratch: &Path,
333 carriers: Option<crypto::import_authority::VerifiedImportCarriers>,
334 ) -> Result<StagedSource, Error> {
335 let spool = self
336 .spool
337 .take()
338 .ok_or(Error::Invalid("provider pack stage required"))?;
339 let completed = tokio::task::spawn_blocking(move || spool.finish())
340 .await
341 .map_err(|error| Error::Preparation(error.to_string()))?
342 .map_err(|error| Error::Preparation(error.to_string()))?;
343 let (pack, index) = completed.artifact_paths();
344 let directory = tempfile::Builder::new()
345 .prefix("provider-download-")
346 .tempdir_in(scratch)?;
347 if !self.ancestry.is_empty() {
348 for name in ["ancestry.pack", "ancestry.idx"] {
349 std::fs::copy(
350 self.ancestry_directory.path().join(name),
351 directory.path().join(name),
352 )?;
353 }
354 }
355 std::fs::hard_link(pack, directory.path().join("source.pack"))?;
356 std::fs::hard_link(index, directory.path().join("source.idx"))?;
357 let mut operations = Vec::new();
358 let mut receipt_records = Vec::new();
359 let mut dependencies = Vec::new();
360 let ancestry = self.ancestry;
361 for item in self.originals {
362 match item {
363 Item::Operations(batch) => {
364 for received in crate::authority_admission::match_batch(&batch)? {
365 operations.push(received.original);
366 receipt_records.extend(received.authority_admission);
367 }
368 }
369 Item::ThreadGenesis(record) => dependencies.push(record),
370 _ => {
371 return Err(Error::Invalid(
372 "provider source originals differ from Ready",
373 ));
374 }
375 }
376 }
377 let mut ready = self.ready.clone();
378 if let Some(carriers) = &carriers {
379 let original = ready
380 .import_authority
381 .as_mut()
382 .ok_or(Error::HostedTrustRequired)?;
383 crate::hybrid::history::replace_receiver_metadata(original, carriers.bundle().clone())
384 .map_err(|e| Error::Preparation(e.to_string()))?;
385 }
386 let staged = tokio::task::spawn_blocking(move || {
387 super::staging::validate_with_receipts_and_carriers(
388 directory,
389 ready,
390 operations,
391 dependencies,
392 receipt_records,
393 carriers,
394 ancestry,
395 )
396 })
397 .await
398 .map_err(|error| Error::Preparation(error.to_string()))??;
399 let pack_path = pack.to_path_buf();
400 let (bytes, digest) = tokio::task::spawn_blocking(move || -> Result<_, Error> {
401 let mut file = std::fs::File::open(pack_path)?;
402 let mut hasher = blake3::Hasher::new();
403 let mut bytes = 0_u64;
404 let mut buffer = [0_u8; 64 * 1024];
405 loop {
406 let read = file.read(&mut buffer)?;
407 if read == 0 {
408 break;
409 }
410 bytes = bytes
411 .checked_add(read as u64)
412 .ok_or(Error::Invalid("provider output length overflow"))?;
413 hasher.update(&buffer[..read]);
414 }
415 Ok((bytes, hasher.finalize()))
416 })
417 .await
418 .map_err(|error| Error::Preparation(error.to_string()))??;
419 if bytes != self.plan.output_pack_length {
420 return Err(Error::Invalid("provider output length differs from plan"));
421 }
422 let result = api::heddle::api::v1alpha2::ProviderResult {
423 extent_set_digest: self.plan.extent_set_digest.clone(),
424 assembly_digest: self.plan.assembly_digest.clone(),
425 verified_range_commitments: self
426 .plan
427 .extents
428 .iter()
429 .map(|extent| {
430 extent
431 .range
432 .as_ref()
433 .map(|range| range.record_set_commitment.clone())
434 .ok_or(Error::Invalid("provider range absent"))
435 })
436 .collect::<Result<Vec<_>, _>>()?,
437 assembled_pack: Some(api::heddle::api::v1alpha2::ObjectAddress {
438 algorithm: "blake3".into(),
439 digest: digest.as_bytes().to_vec(),
440 }),
441 };
442 self.sender
443 .send(&FetchClientFrame {
444 body: Some(fetch_client_frame::Body::ProviderResult(result)),
445 })
446 .await?;
447 self.sender.finish().await?;
448 let frame = self
449 .messages
450 .next()
451 .await?
452 .ok_or(Error::Invalid("provider terminal Complete absent"))?;
453 let Some(fetch_server_frame::Body::Complete(complete)) = frame.body else {
454 return Err(Error::Invalid("provider terminal Complete required"));
455 };
456 let initial = self
457 .state
458 .ready
459 .checkpoint
460 .as_ref()
461 .ok_or(Error::Invalid("provider initial checkpoint absent"))?;
462 let final_checkpoint = complete
463 .checkpoint
464 .as_ref()
465 .ok_or(Error::Invalid("provider final checkpoint absent"))?;
466 if complete.revision != self.state.ready.current
467 || complete.closure != api::heddle::api::v1alpha2::Coverage::Complete as i32
468 || !complete.missing.is_empty()
469 || final_checkpoint.transfer_id != initial.transfer_id
470 || final_checkpoint.plan_digest != initial.plan_digest
471 || final_checkpoint.committed_bytes != bytes
472 {
473 return Err(Error::Invalid(
474 "provider completion differs from verified source",
475 ));
476 }
477 self.messages.cancel();
478 Ok(staged)
479 }
480}
481
482async fn write_chunk(
483 writer: wire::ProviderPackWriter,
484 index: usize,
485 offset: u64,
486 data: Vec<u8>,
487 record: api::heddle::api::v1alpha2::ProviderAssemblyRecord,
488 finished: bool,
489) -> Result<(), Error> {
490 tokio::task::spawn_blocking(move || {
491 writer
492 .write_extent_chunk(index, offset, &data)
493 .map_err(|error| Error::Preparation(error.to_string()))?;
494 if finished {
495 verify_record(&writer, index, &record)?;
496 }
497 Ok(())
498 })
499 .await
500 .map_err(|error| Error::Preparation(error.to_string()))?
501}
502
503fn verify_record(
504 writer: &wire::ProviderPackWriter,
505 index: usize,
506 record: &api::heddle::api::v1alpha2::ProviderAssemblyRecord,
507) -> Result<(), Error> {
508 let mut hasher = blake3::Hasher::new();
509 writer
510 .hash_extent_prefix(index, record.encoded_length, &mut hasher)
511 .map_err(|error| Error::Preparation(error.to_string()))?;
512 let digest = record
513 .encoded_digest
514 .as_ref()
515 .ok_or(Error::Invalid("provider record digest absent"))?;
516 if hasher.finalize().as_bytes() != digest.digest.as_slice() {
517 return Err(Error::Invalid("provider encoded record digest differs"));
518 }
519 writer
520 .mark_verified(index)
521 .map_err(|error| Error::Preparation(error.to_string()))
522}
523
524impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
525 pub async fn begin_provider_fetch(
526 &self,
527 open: FetchOpen,
528 limits: Limits,
529 ) -> Result<ProviderFetch<T::Writer, T::Reader>, Error> {
530 if open.delivery != fetch_open::Delivery::ProviderPreferred as i32
531 || open.checkpoint.is_some()
532 {
533 return Err(Error::Invalid("fresh preferred provider Fetch required"));
534 }
535 let issuer = self
536 .description
537 .endpoint
538 .clone()
539 .ok_or(Error::Invalid("issuer endpoint required"))?;
540 let (sender, mut messages) = self
541 .api
542 .exchange::<rpc::SyncServiceFetch>(&FetchClientFrame {
543 body: Some(fetch_client_frame::Body::Open(open.clone())),
544 })
545 .await?;
546 let frame = messages
547 .next()
548 .await?
549 .ok_or(Error::Invalid("provider Ready required"))?;
550 let Some(fetch_server_frame::Body::Ready(ready)) = frame.body else {
551 return Err(Error::Invalid("first provider response must be Ready"));
552 };
553 if !ready.packs.is_empty() {
554 let mut direct = open.clone();
555 direct.delivery = fetch_open::Delivery::Direct as i32;
556 let state = Validation::new(direct, ready, Some(&issuer), limits)?;
557 sender.finish().await?;
558 return Ok(ProviderFetch::Direct(Box::new(Download {
559 messages,
560 state,
561 })));
562 }
563 let state = Validation::new(open.clone(), ready, Some(&issuer), limits)?;
564 Ok(ProviderFetch::Provider(Box::new(ProviderDownload {
565 sender,
566 messages,
567 state,
568 open,
569 issuer,
570 })))
571 }
572}
573
574impl<W: MessageWriter<Error = transport::Error>, R: MessageReader<Error = transport::Error>>
575 ProviderDownload<W, R>
576{
577 pub async fn negotiate(
578 mut self,
579 signer: &impl ProviderConsentSigner,
580 ) -> Result<ProviderPlanSession<W, R>, Error> {
581 let mut originals = Vec::new();
582 let ancestry_directory = tempfile::Builder::new()
583 .prefix("provider-download-")
584 .tempdir()?;
585 let mut ancestry = super::staging::AncestryInput::new(self.state.excluded_tips.clone());
586 let offer = loop {
587 let frame = self
588 .messages
589 .next()
590 .await?
591 .ok_or(Error::Invalid("provider Offer required"))?;
592 match frame.body {
593 Some(fetch_server_frame::Body::Operations(_))
594 | Some(fetch_server_frame::Body::ThreadGenesis(_))
595 | Some(fetch_server_frame::Body::ImportAncestry(_)) => {
596 match self.state.accept(frame)? {
597 Item::ImportAncestry(page) => ancestry.push(
598 page,
599 self.state
600 .ready
601 .thread
602 .as_ref()
603 .ok_or(Error::Invalid("Thread absent"))?,
604 ancestry_directory.path(),
605 )?,
606 original => originals.push(original),
607 }
608 }
609 Some(fetch_server_frame::Body::ProviderOffer(offer)) => break offer,
610 _ => return Err(Error::Invalid("unexpected frame before provider Offer")),
611 }
612 };
613 if self
614 .state
615 .ready
616 .checkpoint
617 .as_ref()
618 .is_none_or(|checkpoint| checkpoint.plan_digest != offer.assembly_digest)
619 {
620 return Err(Error::Invalid(
621 "provider Offer differs from Ready checkpoint",
622 ));
623 }
624 let candidate =
625 Candidate::new(&self.open, &self.issuer, &signer.client_endpoint()?, offer)?;
626 if candidate.challenge()?.revision.as_ref() != self.state.ready.current.as_ref() {
627 return Err(Error::Invalid(
628 "provider offer differs from admitted revision",
629 ));
630 }
631 self.sender
632 .send(&FetchClientFrame {
633 body: Some(fetch_client_frame::Body::Consent(
634 candidate.consent(signer)?,
635 )),
636 })
637 .await?;
638 let issued = self
639 .messages
640 .next()
641 .await?
642 .ok_or(Error::Invalid("issued provider Plan required"))?;
643 let Some(fetch_server_frame::Body::ProviderPlan(plan)) = issued.body else {
644 return Err(Error::Invalid(
645 "first frame after consent must be provider Plan",
646 ));
647 };
648 candidate.admit(&plan)?;
649 ancestry.finish()?;
650 Ok(ProviderPlanSession {
651 ready: self.state.ready.clone(),
652 plan,
653 originals,
654 ancestry,
655 ancestry_directory,
656 sender: self.sender,
657 messages: self.messages,
658 state: self.state,
659 spool: None,
660 })
661 }
662}
663
664fn pack_manifest(plan: &ProviderPlan) -> Result<ProviderPackManifest, Error> {
668 api::provider_v2::validate_provider_plan(plan)
669 .map_err(|_| Error::Invalid("invalid issued provider plan"))?;
670 let header: [u8; 16] = plan
671 .pack_header
672 .as_slice()
673 .try_into()
674 .map_err(|_| Error::Invalid("invalid provider pack header"))?;
675 let extents = plan
676 .records
677 .iter()
678 .map(|record| {
679 let object = record
680 .object
681 .as_ref()
682 .ok_or(Error::Invalid("provider object absent"))?;
683 let address = object
684 .address
685 .as_ref()
686 .ok_or(Error::Invalid("provider object address absent"))?;
687 let digest = record
688 .encoded_digest
689 .as_ref()
690 .ok_or(Error::Invalid("provider encoded digest absent"))?;
691 let object_hash: [u8; 32] = address
692 .digest
693 .as_slice()
694 .try_into()
695 .map_err(|_| Error::Invalid("provider object digest length"))?;
696 let encoded_hash: [u8; 32] = digest
697 .digest
698 .as_slice()
699 .try_into()
700 .map_err(|_| Error::Invalid("provider encoded digest length"))?;
701 let id = match object.kind.as_str() {
702 "blob" | "tree" => PackObjectId::Hash(ContentHash::from_bytes(object_hash)),
703 "state" => PackObjectId::StateId(StateId::from_bytes(object_hash)),
704 _ => return Err(Error::Invalid("provider object kind")),
705 };
706 Ok(ProviderPackExtent {
707 output_offset: record.output_offset,
708 length: record.encoded_length,
709 digest: encoded_hash,
710 objects: vec![ProviderPackIndexEntry {
711 id,
712 output_offset: record.output_offset,
713 }],
714 })
715 })
716 .collect::<Result<Vec<_>, Error>>()?;
717 Ok(ProviderPackManifest {
718 header,
719 output_pack_length: plan.output_pack_length,
720 extents,
721 })
722}
723
724pub struct Candidate {
727 offer: ProviderOffer,
728}
729
730impl Candidate {
731 pub fn new(
732 open: &FetchOpen,
733 issuer: &EndpointRef,
734 client: &EndpointRef,
735 offer: ProviderOffer,
736 ) -> Result<Self, Error> {
737 if open.delivery != fetch_open::Delivery::ProviderPreferred as i32 {
738 return Err(Error::Invalid("provider offer requires preferred delivery"));
739 }
740 validate_provider_offer(&offer).map_err(|_| Error::Invalid("invalid provider offer"))?;
741 let challenge = offer
742 .challenge
743 .as_ref()
744 .ok_or(Error::Invalid("provider challenge required"))?;
745 if challenge.thread != open.thread
746 || open
747 .revision
748 .as_ref()
749 .is_some_and(|revision| challenge.revision.as_ref() != Some(revision))
750 || challenge.issuer.as_ref() != Some(issuer)
751 || challenge.client.as_ref() != Some(client)
752 {
753 return Err(Error::Invalid(
754 "provider offer differs from selected source or peer",
755 ));
756 }
757 let expiry = challenge
758 .expires_at
759 .as_ref()
760 .ok_or(Error::Invalid("provider expiry required"))?;
761 let now = std::time::SystemTime::now()
762 .duration_since(std::time::UNIX_EPOCH)
763 .map_err(|_| Error::Invalid("provider clock unavailable"))?;
764 if expiry.seconds
765 <= i64::try_from(now.as_secs())
766 .map_err(|_| Error::Invalid("provider clock overflow"))?
767 {
768 return Err(Error::Invalid("provider offer expired"));
769 }
770 Ok(Self { offer })
771 }
772
773 pub fn challenge(&self) -> Result<&ProviderPlanChallenge, Error> {
774 self.offer
775 .challenge
776 .as_ref()
777 .ok_or(Error::Invalid("validated provider challenge absent"))
778 }
779
780 pub fn consent(&self, signer: &impl ProviderConsentSigner) -> Result<ProviderConsent, Error> {
781 let identity = format!("principal:{}", signer.verified_subject()?);
782 let canonical = provider_consent_signing_bytes(self.challenge()?, &identity)
783 .map_err(|_| Error::Invalid("invalid provider consent challenge"))?;
784 let key = signer.public_key();
785 if key.len() != 32 {
786 return Err(Error::Invalid("provider consent key must be Ed25519"));
787 }
788 let signature = signer.sign(&canonical)?;
789 if signature.len() != 64 {
790 return Err(Error::Invalid("provider consent signature length"));
791 }
792 Ok(ProviderConsent {
793 extent_set_digest: self.offer.extent_set_digest.clone(),
794 exact_plan_consent: Some(SignedRecord {
795 format: PROVIDER_CONSENT_FORMAT.into(),
796 canonical_record: canonical,
797 signatures: vec![RecordSignature {
798 public_key: key.to_vec(),
799 signature,
800 }],
801 }),
802 assembly_digest: self.offer.assembly_digest.clone(),
803 })
804 }
805
806 pub fn admit(&self, plan: &ProviderPlan) -> Result<(), Error> {
807 validate_plan_for_offer(&self.offer, plan)
808 .map_err(|_| Error::Invalid("issued provider plan differs from consented offer"))
809 }
810}
811
812#[cfg(test)]
813mod tests {
814 use std::{
815 collections::VecDeque,
816 sync::{
817 Arc,
818 atomic::{AtomicBool, AtomicUsize, Ordering},
819 },
820 };
821
822 use api::{
823 heddle::api::v1alpha2::{
824 Coverage, DescribeEndpointResponse, EndpointKind, FetchComplete, ObjectAddress,
825 PackChunk, ProviderAssemblyRecord, ProviderDialRoute, ProviderExtent,
826 ProviderExtentEvent, ProviderInlineChunk, ProviderInlineSource, ProviderOffer,
827 ProviderOfferExtent, ProviderPhysicalRange, ProviderPlan, ProviderPlanChallenge,
828 ProviderRangeChunk, ProviderRangeSource, ProviderReadTicket, ReadProviderExtentRequest,
829 SharedFacet, ThreadRef, TransferCheckpoint, TransferObject, provider_assembly_record,
830 provider_dial_route, provider_extent_event,
831 },
832 provider_v2::{
833 provider_assembly_digest, provider_extent_set_digest, provider_record_set_commitment,
834 },
835 v2::{
836 MethodDescriptor,
837 client::{Client, Rpc, RpcTransport},
838 },
839 };
840 use crypto::{Ed25519Signer, Signer};
841 use prost::Message;
842
843 use super::*;
844
845 struct Reader(VecDeque<Vec<u8>>);
846 impl MessageReader for Reader {
847 type Error = transport::Error;
848 async fn next(&mut self) -> Result<Option<Vec<u8>>, Self::Error> {
849 Ok(self.0.pop_front())
850 }
851 fn cancel(&mut self) {
852 self.0.clear();
853 }
854 }
855 struct Writer(Arc<AtomicBool>);
856 impl MessageWriter for Writer {
857 type Error = transport::Error;
858 async fn send(&mut self, _: Vec<u8>) -> Result<(), Self::Error> {
859 Err(transport::Error::Protocol("direct fallback sent consent"))
860 }
861 async fn finish(&mut self) -> Result<(), Self::Error> {
862 self.0.store(true, Ordering::Release);
863 Ok(())
864 }
865 fn abort(&mut self) {}
866 }
867 struct Peer {
868 frames: Vec<Vec<u8>>,
869 finished: Arc<AtomicBool>,
870 }
871 impl RpcTransport for Peer {
872 type Error = transport::Error;
873 type Reader = Reader;
874 type Writer = Writer;
875 async fn unary(
876 &self,
877 _: &'static MethodDescriptor,
878 _: Vec<u8>,
879 ) -> Result<Vec<u8>, Self::Error> {
880 Err(transport::Error::Protocol("unused"))
881 }
882 async fn observe(
883 &self,
884 _: &'static MethodDescriptor,
885 _: Vec<u8>,
886 ) -> Result<Reader, Self::Error> {
887 Err(transport::Error::Protocol("unused"))
888 }
889 async fn exchange(
890 &self,
891 _: &'static MethodDescriptor,
892 opening: Vec<u8>,
893 ) -> Result<(Writer, Reader), Self::Error> {
894 let frame = FetchClientFrame::decode(opening.as_slice())
895 .map_err(|_| transport::Error::Protocol("bad opening"))?;
896 assert!(matches!(
897 frame.body,
898 Some(fetch_client_frame::Body::Open(FetchOpen {
899 delivery,
900 routes,
901 ..
902 })) if delivery == fetch_open::Delivery::ProviderPreferred as i32
903 && matches!(
904 routes.as_slice(),
905 [ProviderDialRoute {
906 provider: Some(EndpointRef { kind, .. }),
907 address: Some(provider_dial_route::Address::RelayUrl(_)),
908 }] if *kind == EndpointKind::Provider as i32
909 )
910 ));
911 Ok((
912 Writer(Arc::clone(&self.finished)),
913 Reader(self.frames.clone().into()),
914 ))
915 }
916 }
917
918 #[tokio::test]
919 async fn preferred_ready_with_packs_is_direct_and_terminal_is_checked() {
920 for wrong_terminal in [false, true] {
921 let (mut open, ready, endpoint, artifacts) = super::super::tests::fixture();
922 open.delivery = fetch_open::Delivery::ProviderPreferred as i32;
923 open.routes = vec![ProviderDialRoute {
924 provider: Some(EndpointRef {
925 public_key: vec![9; 32],
926 kind: EndpointKind::Provider as i32,
927 }),
928 address: Some(provider_dial_route::Address::RelayUrl(
929 "https://relay.example/".to_string(),
930 )),
931 }];
932 let mut frames = vec![
933 FetchServerFrame {
934 body: Some(fetch_server_frame::Body::Ready(ready.clone())),
935 }
936 .encode_to_vec(),
937 ];
938 for (index, bytes) in artifacts.iter().enumerate() {
939 frames.push(
940 FetchServerFrame {
941 body: Some(fetch_server_frame::Body::Pack(PackChunk {
942 extent: Some(ready.packs[index].clone()),
943 data: bytes.clone(),
944 })),
945 }
946 .encode_to_vec(),
947 );
948 }
949 let mut checkpoint = ready.checkpoint.clone().expect("fixture checkpoint");
950 checkpoint.committed_bytes = if wrong_terminal {
951 0
952 } else {
953 artifacts.iter().map(|bytes| bytes.len() as u64).sum()
954 };
955 frames.push(
956 FetchServerFrame {
957 body: Some(fetch_server_frame::Body::Complete(FetchComplete {
958 revision: ready.current.clone(),
959 checkpoint: Some(checkpoint),
960 closure: Coverage::Complete as i32,
961 missing: vec![],
962 })),
963 }
964 .encode_to_vec(),
965 );
966 let finished = Arc::new(AtomicBool::new(false));
967 let remote = Remote {
968 api: Client::new(
969 Peer {
970 frames,
971 finished: Arc::clone(&finished),
972 },
973 [rpc::SyncServiceFetch::METHOD.path.into()],
974 ),
975 description: DescribeEndpointResponse {
976 endpoint: Some(endpoint),
977 ..Default::default()
978 },
979 };
980 let ProviderFetch::Direct(mut download) = remote
981 .begin_provider_fetch(open, Limits::default())
982 .await
983 .expect("direct fallback")
984 else {
985 panic!("expected direct branch")
986 };
987 assert!(
988 finished.load(Ordering::Acquire),
989 "direct path closes request half"
990 );
991 for _ in 0..2 {
992 assert!(matches!(
993 download.next().await.expect("pack"),
994 Some(Item::Pack(_))
995 ));
996 }
997 if wrong_terminal {
998 assert!(matches!(
999 download.next().await,
1000 Err(Error::Invalid(
1001 "download does not match its exact declared source coverage"
1002 ))
1003 ));
1004 } else {
1005 assert!(matches!(
1006 download.next().await.expect("Complete"),
1007 Some(Item::Complete(_))
1008 ));
1009 }
1010 }
1011 }
1012
1013 #[test]
1014 fn provider_ready_has_no_direct_pack_or_partial_fallback() {
1015 let (mut open, mut ready, endpoint, _) = super::super::tests::fixture();
1016 open.delivery = fetch_open::Delivery::ProviderPreferred as i32;
1017 assert!(
1018 Validation::new(
1019 open.clone(),
1020 ready.clone(),
1021 Some(&endpoint),
1022 Limits::default()
1023 )
1024 .is_err(),
1025 "provider mode cannot silently receive direct source packs"
1026 );
1027 ready.packs.clear();
1028 Validation::new(
1029 open.clone(),
1030 ready.clone(),
1031 Some(&endpoint),
1032 Limits::default(),
1033 )
1034 .expect("complete provider Ready with no direct artifact");
1035 ready.full_closure_available = false;
1036 assert!(
1037 Validation::new(open, ready, Some(&endpoint), Limits::default()).is_err(),
1038 "provider offer cannot downgrade whole-source disclosure"
1039 );
1040 }
1041
1042 #[test]
1043 fn encoded_record_digest_is_checked_before_completion() {
1044 let scratch = tempfile::tempdir().expect("test scratch");
1045 let body = b"one encoded record";
1046 let mut header = [0_u8; 16];
1047 header[..4].copy_from_slice(b"LMPK");
1048 header[4..8].copy_from_slice(&4_u32.to_be_bytes());
1049 header[8..].copy_from_slice(&1_u64.to_be_bytes());
1050 let spool = ProviderPackSpool::new_in(
1051 scratch.path(),
1052 ProviderPackManifest {
1053 header,
1054 output_pack_length: 16 + body.len() as u64 + 32,
1055 extents: vec![ProviderPackExtent {
1056 output_offset: 16,
1057 length: body.len() as u64,
1058 digest: *blake3::hash(body).as_bytes(),
1059 objects: vec![ProviderPackIndexEntry {
1060 id: PackObjectId::Hash(ContentHash::from_bytes([7; 32])),
1061 output_offset: 16,
1062 }],
1063 }],
1064 },
1065 )
1066 .expect("bounded spool fixture");
1067 let writer = spool.writer();
1068 writer
1069 .write_extent_chunk(0, 0, body)
1070 .expect("fixture bytes");
1071 let mut record = ProviderAssemblyRecord {
1072 encoded_length: body.len() as u64,
1073 encoded_digest: Some(ObjectAddress {
1074 algorithm: "blake3".into(),
1075 digest: vec![0; 32],
1076 }),
1077 ..Default::default()
1078 };
1079 assert!(
1080 verify_record(&writer, 0, &record).is_err(),
1081 "a fully received but wrong record cannot become verified"
1082 );
1083 record
1084 .encoded_digest
1085 .as_mut()
1086 .expect("fixture digest")
1087 .digest = blake3::hash(body).as_bytes().to_vec();
1088 verify_record(&writer, 0, &record).expect("exact record becomes verified");
1089 }
1090
1091 struct ConsentWriter;
1092 impl MessageWriter for ConsentWriter {
1093 type Error = transport::Error;
1094 async fn send(&mut self, _: Vec<u8>) -> Result<(), Self::Error> {
1095 Ok(())
1096 }
1097 async fn finish(&mut self) -> Result<(), Self::Error> {
1098 Ok(())
1099 }
1100 fn abort(&mut self) {}
1101 }
1102
1103 struct IssuerPeer {
1104 frames: Vec<Vec<u8>>,
1105 }
1106 impl RpcTransport for IssuerPeer {
1107 type Error = transport::Error;
1108 type Reader = Reader;
1109 type Writer = ConsentWriter;
1110 async fn unary(
1111 &self,
1112 _: &'static MethodDescriptor,
1113 _: Vec<u8>,
1114 ) -> Result<Vec<u8>, Self::Error> {
1115 Err(transport::Error::Protocol("unused"))
1116 }
1117 async fn observe(
1118 &self,
1119 _: &'static MethodDescriptor,
1120 _: Vec<u8>,
1121 ) -> Result<Reader, Self::Error> {
1122 Err(transport::Error::Protocol("unused"))
1123 }
1124 async fn exchange(
1125 &self,
1126 _: &'static MethodDescriptor,
1127 opening: Vec<u8>,
1128 ) -> Result<(ConsentWriter, Reader), Self::Error> {
1129 let frame = FetchClientFrame::decode(opening.as_slice())
1130 .map_err(|_| transport::Error::Protocol("bad opening"))?;
1131 assert!(
1132 matches!(
1133 frame.body,
1134 Some(fetch_client_frame::Body::Open(FetchOpen {
1135 delivery,
1136 ref routes,
1137 ..
1138 })) if delivery == fetch_open::Delivery::ProviderPreferred as i32
1139 && !routes.is_empty()
1140 ),
1141 "provider branch must open Fetch with ProviderPreferred routes"
1142 );
1143 Ok((ConsentWriter, Reader(self.frames.clone().into())))
1144 }
1145 }
1146
1147 struct ProviderPeer {
1148 frames: Vec<Vec<u8>>,
1149 range_reads: Arc<AtomicUsize>,
1150 }
1151 impl RpcTransport for ProviderPeer {
1152 type Error = transport::Error;
1153 type Reader = Reader;
1154 type Writer = ConsentWriter;
1155 async fn unary(
1156 &self,
1157 _: &'static MethodDescriptor,
1158 _: Vec<u8>,
1159 ) -> Result<Vec<u8>, Self::Error> {
1160 Err(transport::Error::Protocol("unused"))
1161 }
1162 async fn observe(
1163 &self,
1164 method: &'static MethodDescriptor,
1165 opening: Vec<u8>,
1166 ) -> Result<Reader, Self::Error> {
1167 assert_eq!(method.path, rpc::SyncServiceReadProviderExtent::METHOD.path);
1168 let _ = ReadProviderExtentRequest::decode(opening.as_slice())
1169 .map_err(|_| transport::Error::Protocol("bad provider extent request"))?;
1170 self.range_reads.fetch_add(1, Ordering::SeqCst);
1171 Ok(Reader(self.frames.clone().into()))
1172 }
1173 async fn exchange(
1174 &self,
1175 _: &'static MethodDescriptor,
1176 _: Vec<u8>,
1177 ) -> Result<(ConsentWriter, Reader), Self::Error> {
1178 Err(transport::Error::Protocol("unused"))
1179 }
1180 }
1181
1182 struct TestConsent {
1183 signer: Ed25519Signer,
1184 subject: String,
1185 client: EndpointRef,
1186 }
1187 impl ProviderConsentSigner for TestConsent {
1188 fn verified_subject(&self) -> Result<String, Error> {
1189 Ok(self.subject.clone())
1190 }
1191 fn client_endpoint(&self) -> Result<EndpointRef, Error> {
1192 Ok(self.client.clone())
1193 }
1194 fn public_key(&self) -> &[u8] {
1195 self.signer.public_key()
1196 }
1197 fn sign(&self, canonical: &[u8]) -> Result<Vec<u8>, Error> {
1198 self.signer
1199 .sign(canonical)
1200 .map_err(|error| Error::Preparation(error.to_string()))
1201 }
1202 }
1203
1204 fn offer_from(plan: &ProviderPlan) -> ProviderOffer {
1205 ProviderOffer {
1206 extent_set_digest: plan.extent_set_digest.clone(),
1207 extents: plan
1208 .extents
1209 .iter()
1210 .map(|extent| {
1211 let ticket = extent.ticket.as_ref().expect("fixture ticket");
1212 ProviderOfferExtent {
1213 provider: extent.provider.clone(),
1214 range: extent.range.clone(),
1215 spool: ticket.spool.clone(),
1216 facet: ticket.facet,
1217 audience: ticket.audience.clone(),
1218 content_root: ticket.content_root.clone(),
1219 }
1220 })
1221 .collect(),
1222 challenge: plan.challenge.clone(),
1223 assembly_digest: plan.assembly_digest.clone(),
1224 pack_header: plan.pack_header.clone(),
1225 output_pack_length: plan.output_pack_length,
1226 records: plan.records.clone(),
1227 }
1228 }
1229
1230 fn mixed_plan(
1231 thread: ThreadRef,
1232 revision: api::heddle::api::v1alpha2::RevisionRef,
1233 issuer: EndpointRef,
1234 client: EndpointRef,
1235 provider: EndpointRef,
1236 provider_bytes: &[u8],
1237 inline_bytes: &[u8],
1238 ) -> ProviderPlan {
1239 let spool = thread.spool.clone().expect("thread spool");
1240 let expiry = prost_types::Timestamp {
1241 seconds: i64::MAX / 2,
1242 nanos: 0,
1243 };
1244 let object = |digest: [u8; 32], size: u64| TransferObject {
1245 address: Some(ObjectAddress {
1246 algorithm: "blake3".into(),
1247 digest: digest.to_vec(),
1248 }),
1249 kind: "blob".into(),
1250 facet: SharedFacet::Source as i32,
1251 size,
1252 availability: Coverage::Complete as i32,
1253 };
1254 let records = vec![
1255 ProviderAssemblyRecord {
1256 object: Some(object([8; 32], provider_bytes.len() as u64)),
1257 encoded_length: provider_bytes.len() as u64,
1258 encoded_digest: Some(ObjectAddress {
1259 algorithm: "blake3".into(),
1260 digest: blake3::hash(provider_bytes).as_bytes().to_vec(),
1261 }),
1262 output_offset: 16,
1263 source: Some(provider_assembly_record::Source::Provider(
1264 ProviderRangeSource {
1265 extent_index: 0,
1266 source_offset: 0,
1267 },
1268 )),
1269 },
1270 ProviderAssemblyRecord {
1271 object: Some(object([10; 32], inline_bytes.len() as u64)),
1272 encoded_length: inline_bytes.len() as u64,
1273 encoded_digest: Some(ObjectAddress {
1274 algorithm: "blake3".into(),
1275 digest: blake3::hash(inline_bytes).as_bytes().to_vec(),
1276 }),
1277 output_offset: 16 + provider_bytes.len() as u64,
1278 source: Some(provider_assembly_record::Source::Inline(
1279 ProviderInlineSource {},
1280 )),
1281 },
1282 ];
1283 let mut range = ProviderPhysicalRange {
1284 pack_id: vec![5; 32],
1285 object_etag: "etag-1".into(),
1286 offset: 128,
1287 length: provider_bytes.len() as u64,
1288 record_set_commitment: vec![],
1289 };
1290 range.record_set_commitment = provider_record_set_commitment(&range, &records, 0)
1291 .expect("tiled provider range")
1292 .to_vec();
1293 let ticket = ProviderReadTicket {
1294 attenuated_capability: vec![99],
1295 extent_set_digest: vec![],
1296 spool: Some(spool.clone()),
1297 facet: SharedFacet::Source as i32,
1298 audience: "Public".into(),
1299 content_root: vec![6; 32],
1300 pack_id: range.pack_id.clone(),
1301 object_etag: range.object_etag.clone(),
1302 offset: range.offset,
1303 length: range.length,
1304 provider: Some(provider.clone()),
1305 client: Some(client.clone()),
1306 assembly_digest: vec![],
1307 expires_at: Some(expiry),
1308 record_set_commitment: range.record_set_commitment.clone(),
1309 };
1310 let mut header = b"LMPK".to_vec();
1311 header.extend_from_slice(&4_u32.to_be_bytes());
1312 header.extend_from_slice(&2_u64.to_be_bytes());
1313 let mut plan = ProviderPlan {
1314 extent_set_digest: vec![],
1315 extents: vec![ProviderExtent {
1316 provider: Some(provider),
1317 ticket: Some(ticket),
1318 range: Some(range),
1319 }],
1320 challenge: Some(ProviderPlanChallenge {
1321 nonce: vec![7; 16],
1322 thread: Some(thread),
1323 revision: Some(revision),
1324 issuer: Some(issuer),
1325 client: Some(client),
1326 expires_at: Some(expiry),
1327 extent_set_digest: vec![],
1328 assembly_digest: vec![],
1329 }),
1330 assembly_digest: vec![],
1331 pack_header: header,
1332 output_pack_length: 16 + provider_bytes.len() as u64 + inline_bytes.len() as u64 + 32,
1333 records,
1334 };
1335 let set = provider_extent_set_digest(&plan).expect("extent set");
1336 plan.extent_set_digest = set.to_vec();
1337 plan.challenge
1338 .as_mut()
1339 .expect("challenge")
1340 .extent_set_digest = set.to_vec();
1341 plan.extents[0]
1342 .ticket
1343 .as_mut()
1344 .expect("ticket")
1345 .extent_set_digest = set.to_vec();
1346 let assembly = provider_assembly_digest(&plan).expect("assembly");
1347 plan.assembly_digest = assembly.to_vec();
1348 plan.challenge.as_mut().expect("challenge").assembly_digest = assembly.to_vec();
1349 plan.extents[0]
1350 .ticket
1351 .as_mut()
1352 .expect("ticket")
1353 .assembly_digest = assembly.to_vec();
1354 plan
1355 }
1356
1357 #[tokio::test]
1358 async fn preferred_provider_fetch_negotiates_consent_and_rehashes_provider_ranges() {
1359 let (mut open, mut ready, issuer, _) = super::super::tests::fixture();
1360 let provider = EndpointRef {
1361 public_key: vec![9; 32],
1362 kind: EndpointKind::Provider as i32,
1363 };
1364 let client = EndpointRef {
1365 public_key: vec![3; 32],
1366 kind: EndpointKind::Device as i32,
1367 };
1368 let provider_bytes = b"provrec!".as_slice();
1369 let inline_bytes = b"inlinrec".as_slice();
1370 open.delivery = fetch_open::Delivery::ProviderPreferred as i32;
1371 open.routes = vec![ProviderDialRoute {
1372 provider: Some(provider.clone()),
1373 address: Some(provider_dial_route::Address::RelayUrl(
1374 "https://relay.example/".to_string(),
1375 )),
1376 }];
1377 let plan = mixed_plan(
1378 open.thread.clone().expect("thread"),
1379 open.revision.clone().expect("revision"),
1380 issuer.clone(),
1381 client.clone(),
1382 provider.clone(),
1383 provider_bytes,
1384 inline_bytes,
1385 );
1386 let offer = offer_from(&plan);
1387 ready.packs.clear();
1388 ready.full_closure_available = true;
1389 let mut checkpoint = ready.checkpoint.clone().expect("fixture checkpoint");
1390 checkpoint.plan_digest = offer.assembly_digest.clone();
1391 ready.checkpoint = Some(checkpoint.clone());
1392
1393 let issuer_frames = vec![
1394 FetchServerFrame {
1395 body: Some(fetch_server_frame::Body::Ready(ready)),
1396 }
1397 .encode_to_vec(),
1398 FetchServerFrame {
1399 body: Some(fetch_server_frame::Body::ProviderOffer(offer)),
1400 }
1401 .encode_to_vec(),
1402 FetchServerFrame {
1403 body: Some(fetch_server_frame::Body::ProviderPlan(plan.clone())),
1404 }
1405 .encode_to_vec(),
1406 FetchServerFrame {
1407 body: Some(fetch_server_frame::Body::ProviderInline(
1408 ProviderInlineChunk {
1409 assembly_digest: plan.assembly_digest.clone(),
1410 record_index: 1,
1411 offset: 0,
1412 data: inline_bytes.to_vec(),
1413 },
1414 )),
1415 }
1416 .encode_to_vec(),
1417 ];
1418 let range = plan.extents[0].range.clone().expect("range");
1419 let provider_frames = vec![
1420 ProviderExtentEvent {
1421 body: Some(provider_extent_event::Body::Ready(range.clone())),
1422 }
1423 .encode_to_vec(),
1424 ProviderExtentEvent {
1425 body: Some(provider_extent_event::Body::Chunk(ProviderRangeChunk {
1426 offset: 0,
1427 data: provider_bytes.to_vec(),
1428 })),
1429 }
1430 .encode_to_vec(),
1431 ProviderExtentEvent {
1432 body: Some(provider_extent_event::Body::Complete(TransferCheckpoint {
1433 committed_bytes: range.length,
1434 ..Default::default()
1435 })),
1436 }
1437 .encode_to_vec(),
1438 ];
1439 let range_reads = Arc::new(AtomicUsize::new(0));
1440 let issuer_remote = Remote {
1441 api: Client::new(
1442 IssuerPeer {
1443 frames: issuer_frames,
1444 },
1445 [rpc::SyncServiceFetch::METHOD.path.into()],
1446 ),
1447 description: DescribeEndpointResponse {
1448 endpoint: Some(issuer),
1449 ..Default::default()
1450 },
1451 };
1452 let provider_remote = Remote {
1453 api: Client::new(
1454 ProviderPeer {
1455 frames: provider_frames,
1456 range_reads: Arc::clone(&range_reads),
1457 },
1458 [rpc::SyncServiceReadProviderExtent::METHOD.path.into()],
1459 ),
1460 description: DescribeEndpointResponse {
1461 endpoint: Some(provider),
1462 ..Default::default()
1463 },
1464 };
1465 let signer = TestConsent {
1466 signer: Ed25519Signer::from_seed(&[61; 32]).expect("consent key"),
1467 subject: "provider-reader".into(),
1468 client,
1469 };
1470
1471 let ProviderFetch::Provider(download) = issuer_remote
1472 .begin_provider_fetch(open, Limits::default())
1473 .await
1474 .expect("provider Ready without direct packs")
1475 else {
1476 panic!("expected provider branch, not direct fallback");
1477 };
1478 let mut session = download
1479 .negotiate(&signer)
1480 .await
1481 .expect("negotiate offer, consent, ProviderPlan");
1482 assert!(
1483 session.plan().records.iter().any(|record| matches!(
1484 record.source,
1485 Some(provider_assembly_record::Source::Provider(_))
1486 )),
1487 "admitted plan must contain provider-sourced records"
1488 );
1489 let scratch = tempfile::tempdir().expect("provider scratch");
1490 session
1491 .receive_inline(scratch.path())
1492 .await
1493 .expect("inline records");
1494 session
1495 .receive_provider_ranges(std::slice::from_ref(&provider_remote))
1496 .await
1497 .expect("provider ranges");
1498 let reads = range_reads.load(Ordering::SeqCst);
1499 assert_eq!(
1500 reads, 1,
1501 "records came via receive_provider_ranges (ReadProviderExtent observe count={reads})"
1502 );
1503 }
1504}