Skip to main content

heddle_thread_api/fetch/
provider.rs

1//! Capability-free candidate admission followed by exact signed consent.
2//! The caller's verified credential implementation owns its terminal PoP key;
3//! this module never accepts an arbitrary identity string or private key per
4//! Fetch invocation.
5
6use 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
28/// Implemented by the same credential that signed the Fetch opening. The
29/// server verifies this identity against its authenticated Biscuit subject
30/// and requires the signature key to equal the credential's terminal cnf key.
31pub 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
38/// An authenticated Fetch exchange whose request half remains open for exact
39/// consent and the final verified result. Dropping it aborts both halves.
40pub 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
51/// The issuer may explicitly select ordinary direct source delivery when a
52/// preferred provider transfer cannot be offered. Only the admitted Ready
53/// decides the branch; an arbitrary stream error never triggers a retry.
54pub 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
62/// Only an issued plan matching the signed candidate can reach this stage.
63pub 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    /// The exact ticketed layout admitted after client consent.
80    pub fn plan(&self) -> &ProviderPlan {
81        &self.plan
82    }
83
84    /// The exact issuer admission that preceded the unsigned offer.
85    pub fn ready(&self) -> &TransferReady {
86        &self.ready
87    }
88
89    /// Reserve a bounded virtual pack and receive only the issuer's inline
90    /// records. Every chunk must name its exact record and next offset; the
91    /// wire framing cannot allocate or write outside the canonical layout.
92    /// The session retains the spool for separately authorized provider ranges.
93    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    /// Read each ticketed physical range from its exact provider endpoint.
184    /// The provider cannot change virtual placement: every byte is written
185    /// only into its precommitted record, then independently rehashed.
186    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    /// Finalize exact pack integrity and signed source closure before sending
316    /// the provider result. A terminal Complete must acknowledge the same
317    /// transfer, plan and verified output length before staging is returned.
318    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
610/// Convert only a canonical, ticketed plan to the positional spool layout.
611/// Each encoded source record is verified separately before the spool may
612/// finalize, including records that arrive inline rather than from a provider.
613fn 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
670/// An unsigned offer bound to the authenticated issuer, client, and exact
671/// selected source. It grants no provider read until final ticket admission.
672pub 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}