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