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