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