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