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