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