Skip to main content

rig_core/
wire.rs

1//! Operations, provider wires and their decoders. An [`Operation`] names what
2//! goes in, what comes out event by event and what ends a reply, and its
3//! [`Fold`] turns one reply's events and end into the response. A [`Wire`]
4//! describes a provider endpoint as plain data, encodes requests without
5//! transport access, and names a fresh [`Decoder`] for each reply. A decoder
6//! is told nothing about how the reply arrives, so buffered and streamed
7//! replies go through the same decoder and fold.
8//!
9//! ```
10//! use rig_core::wire::WireFrame;
11//!
12//! let frame = WireFrame::Bytes(b"response".to_vec());
13//! assert_eq!(frame.as_str(), "response");
14//! ```
15
16use crate::error::{EncodeError, ProviderError};
17use crate::http_client::MultipartForm;
18use crate::wasm_compat::{WasmCompatSend, WasmCompatSync};
19
20pub use crate::http_client::framing::{Framing, WireFrame};
21pub use crate::observe::{
22    AdapterErrorEnvelope, AdapterEvent, AdapterUsage, AdapterVerdict, ObservationSink,
23};
24
25mod citation;
26pub mod document;
27pub(crate) mod secret;
28
29pub use citation::{SpanUnit, WireCitation, WireSpan};
30pub use secret::Secret;
31
32/// The request a wire sends, and how its reply is framed.
33///
34/// Data only: built by [`Wire::encode`] from the wire and the request, and
35/// read by the driver. A wire never touches a socket or a request extension.
36///
37/// `Debug` shows request methods and URI paths, not schemes, authorities,
38/// queries, header values or bodies. Paths are not scrubbed: callers must
39/// still avoid placing sensitive data in them.
40pub struct Encoded {
41    /// The HTTP request.
42    pub request: http::Request<Body>,
43    /// How the reply's bytes split into frames.
44    pub framing: Framing,
45    /// The reply header carrying the provider's transport request id
46    /// (Anthropic `request-id`, OpenAI `x-request-id`), when the provider
47    /// reports one. `None` is "does not report one", never an error.
48    pub request_id_header: Option<&'static str>,
49    /// Whether a streamed reply may omit `Content-Type` (one gateway
50    /// replays Responses bodies without it). A *wrong* content type is
51    /// still rejected.
52    pub relaxed_content_type: bool,
53    /// Stable endpoint template for observation grouping, without base-URL
54    /// prefixes or interpolated values. `None` uses the concrete request path.
55    pub route: Option<&'static str>,
56    /// Reads observation facts off each raw reply payload before
57    /// normalization discards them. `None` projects nothing.
58    pub project: Option<Projector>,
59    /// Whether a frame carries only analysis metadata: it still decodes, but
60    /// does not advance observation's EOF and corruption positions.
61    pub analysis_only: Option<fn(&WireFrame) -> bool>,
62}
63
64/// Reads one raw payload's observation facts (verdicts, usage, ids, error
65/// envelopes) through the sink. Text it forwards goes through
66/// [`ObservationSink::scrub`]: a payload can echo credentials.
67pub type Projector = fn(&[u8], &mut ObservationSink<'_>);
68
69impl Encoded {
70    /// One request whose provider reports no transport request id.
71    pub fn new(request: http::Request<Body>, framing: Framing) -> Self {
72        Self {
73            request,
74            framing,
75            request_id_header: None,
76            relaxed_content_type: false,
77            route: None,
78            project: None,
79            analysis_only: None,
80        }
81    }
82
83    /// Name the reply header carrying the provider's transport request id.
84    pub fn with_request_id_header(mut self, header: Option<&'static str>) -> Self {
85        self.request_id_header = header;
86        self
87    }
88
89    /// Accept a streamed reply that names no content type.
90    pub fn with_relaxed_content_type(mut self) -> Self {
91        self.relaxed_content_type = true;
92        self
93    }
94
95    /// Name the endpoint template observation groups attempts under.
96    pub fn with_route(mut self, route: Option<&'static str>) -> Self {
97        self.route = route;
98        self
99    }
100
101    /// Read observation facts off each reply payload through `project`.
102    pub fn with_projection(mut self, project: Projector) -> Self {
103        self.project = Some(project);
104        self
105    }
106
107    /// Exempt frames `analysis_only` accepts from observation's frame
108    /// positions.
109    pub fn with_analysis_only(mut self, analysis_only: fn(&WireFrame) -> bool) -> Self {
110        self.analysis_only = Some(analysis_only);
111        self
112    }
113}
114
115/// A request body: bytes, or a multipart form for the upload endpoints.
116///
117/// `Debug` prints the body's shape and size, never its bytes: a request
118/// body carries prompts, documents and uploaded files.
119pub enum Body {
120    /// A serialized body (JSON for every wire in this crate, or empty).
121    Bytes(Vec<u8>),
122    /// A multipart form (audio transcription, image edits).
123    Multipart(MultipartForm),
124}
125
126impl Body {
127    /// An empty body, for a `GET`.
128    pub fn empty() -> Self {
129        Self::Bytes(Vec::new())
130    }
131}
132
133impl std::fmt::Debug for Body {
134    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
135        match self {
136            Self::Bytes(bytes) => write!(f, "Bytes({} bytes)", bytes.len()),
137            Self::Multipart(_) => f.write_str("Multipart"),
138        }
139    }
140}
141
142impl std::fmt::Debug for Encoded {
143    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
144        // Diagnostics omit query, authority, headers and bodies. A caller's
145        // path can still contain sensitive data; it is not scrubbed here.
146        f.debug_struct("Encoded")
147            .field(
148                "request",
149                &(self.request.method(), self.request.uri().path()),
150            )
151            .field("framing", &self.framing)
152            .field("request_id_header", &self.request_id_header)
153            .field("relaxed_content_type", &self.relaxed_content_type)
154            .field("route", &self.route)
155            .field("project", &self.project.is_some())
156            .finish()
157    }
158}
159
160/// Reply mode used by the wire to select request encoding and framing.
161#[derive(Debug, Clone, Copy, PartialEq, Eq)]
162pub enum Mode {
163    /// One whole reply ([`Model::call`](crate::driver::Model::call)).
164    Unary,
165    /// A streamed reply ([`Model::stream`](crate::driver::Model::stream)).
166    Streaming,
167}
168
169/// An operation: what goes in, what comes out event by event, what ends a
170/// reply, and how those fold into one response.
171///
172/// Implemented once per operation, never per provider, so every provider of
173/// an operation is interchangeable behind a [`DynModel`](crate::DynModel).
174pub trait Operation: Sized + 'static {
175    /// The normalized request this operation accepts.
176    type Request: WasmCompatSend + 'static;
177    /// One decoded step of a reply.
178    type Event: WasmCompatSend + 'static;
179    /// What the provider sends when it ends the reply. A fold needs one to
180    /// produce a response, so a reply that stopped early has none.
181    type End: WasmCompatSend + 'static;
182    /// The normalized response the events and the end fold into.
183    type Response: WasmCompatSend + 'static;
184    /// One reply's state: it absorbs the events the consumer sees, then
185    /// folds them with the end into the response.
186    type Fold: Fold<Self> + WasmCompatSend + 'static;
187    /// Who builds the events: [`Free`] for events a decoder writes itself,
188    /// [`Assembled`] for the completion events only the writer's part
189    /// handles build.
190    type Emit: Emit<Self>;
191
192    /// The state for one reply to `request`, before its first frame. The
193    /// call says who answers and in which mode; an operation that records
194    /// telemetry opens its span here and hands it to
195    /// [`Call::instrument`].
196    fn fold(request: &Self::Request, call: &mut Call<'_>) -> Self::Fold;
197
198    /// Shape `request` for the wire `wire` describes, or reject a request
199    /// no provider of this operation can answer. The driver calls it before
200    /// the request is encoded, so a rejected request reaches no wire or
201    /// transport. Passes every request through by default.
202    fn prepare(
203        request: Self::Request,
204        wire: &Descriptor<'_>,
205    ) -> Result<Self::Request, ProviderError> {
206        let _ = wire;
207        Ok(request)
208    }
209}
210
211/// Events a decoder builds and writes with [`Out::event`].
212#[derive(Debug, Clone, Copy, PartialEq, Eq)]
213pub enum Free {}
214
215/// Events only the completion writer's part handles build.
216#[derive(Debug, Clone, Copy, PartialEq, Eq)]
217pub enum Assembled {}
218
219/// Who builds an operation's events: [`Free`] or [`Assembled`]. Sealed.
220pub trait Emit<Op: Operation>: reply::Closing<Op> {}
221
222impl<Op: Operation> reply::Closing<Op> for Free {
223    fn close(_shared: &mut Shared<Op>) {}
224}
225
226impl<Op: Operation> Emit<Op> for Free {}
227
228/// What the driver knows about a call before its first frame.
229pub struct Call<'a> {
230    /// What the wire says about itself.
231    pub wire: &'a Descriptor<'a>,
232    /// How the reply arrives.
233    pub mode: Mode,
234    pub(crate) span: tracing::Span,
235}
236
237impl<'a> Call<'a> {
238    /// A call to the wire `wire` describes, in `mode`, under no span.
239    pub(crate) fn new(wire: &'a Descriptor<'a>, mode: Mode) -> Self {
240        Self {
241            wire,
242            mode,
243            span: tracing::Span::none(),
244        }
245    }
246
247    /// Run the call under `span`: a unary call sends under it and a stream
248    /// decodes under it.
249    pub fn instrument(&mut self, span: tracing::Span) {
250        self.span = span;
251    }
252}
253
254/// One reply's state for an operation.
255///
256/// Each event is [`absorb`](Self::absorb)ed as the consumer takes it, so
257/// the fold never runs ahead of what was delivered. The reply's end is what
258/// [`finish`](Self::finish) needs: without the provider's end there is no
259/// response.
260pub trait Fold<Op: Operation> {
261    /// Absorb one event the consumer is about to see. An error fails the
262    /// reply.
263    fn absorb(&mut self, event: &Op::Event) -> Result<(), ProviderError>;
264
265    /// The response, from the events absorbed, the provider's end and what
266    /// the driver learned about the reply.
267    fn finish(self, end: Op::End, reply: Reply) -> Result<Op::Response, ProviderError>;
268}
269
270/// What the driver learned about a reply beyond its events.
271#[derive(Debug, Clone, PartialEq)]
272pub struct Reply {
273    /// The provider descriptor name, for the response's `provider` field.
274    pub provider: String,
275    /// The reply's provider document (a completion's `raw`): the whole body
276    /// when it is one JSON document, else what the wire's reassembler
277    /// rebuilt or a free-event decoder recorded, else `Null`.
278    pub raw: serde_json::Value,
279    /// The provider's transport request id from the reply headers.
280    pub provider_request_id: Option<String>,
281}
282
283/// Whether a decoder step left the reply open. Only [`Out::end`] makes an
284/// [`Ended`].
285#[derive(Debug)]
286#[must_use]
287pub enum Flow {
288    /// The reply continues.
289    More,
290    /// The provider ended the reply; nothing is read after it.
291    Ended(Ended),
292}
293
294/// Proof that a decoder ended its reply with [`Out::end`].
295#[derive(Debug)]
296pub struct Ended(());
297
298pub(crate) use reply::Shared;
299
300pub(crate) mod reply {
301    use std::collections::VecDeque;
302
303    use super::{Fold, Operation, ProviderError, Reply};
304    use crate::streaming::Item;
305
306    /// One reply's items and end, shared by the driver that writes them and
307    /// the stream that takes them.
308    pub struct Shared<Op: Operation> {
309        pub(crate) fold: Op::Fold,
310        pub(crate) items: VecDeque<Result<Item<Op::Event>, ProviderError>>,
311        pub(crate) end: Option<Op::End>,
312        pub(crate) raw: Option<serde_json::Value>,
313        /// The response a relayed reply's origin already folded.
314        pub(crate) response: Option<Op::Response>,
315        /// What the transport reported: the reply's request id, its whole
316        /// document and its request path.
317        pub(crate) request_id: Option<String>,
318        pub(crate) document: Option<serde_json::Value>,
319        pub(crate) route: String,
320    }
321
322    impl<Op: Operation> Shared<Op> {
323        pub(crate) fn new(fold: Op::Fold) -> Self {
324            Self {
325                fold,
326                items: VecDeque::new(),
327                end: None,
328                raw: None,
329                response: None,
330                request_id: None,
331                document: None,
332                route: String::new(),
333            }
334        }
335
336        /// Take the next item: an event is absorbed by the fold before it
337        /// leaves, and an error is stamped with the reply's request id.
338        pub(crate) fn take(&mut self) -> Option<Result<Item<Op::Event>, ProviderError>> {
339            Some(match self.items.pop_front()? {
340                Ok(Item::Event(event)) => self.fold.absorb(&event).map(|()| Item::Event(event)),
341                Ok(unknown) => Ok(unknown),
342                // An id an upstream constructor already attached wins: it
343                // saw the reply.
344                Err(error) => Err(error.with_provider_request_id(self.request_id.clone())),
345            })
346        }
347
348        /// What the driver learned about the reply so far, from `provider`.
349        pub(crate) fn reply(&self, provider: &str) -> Reply {
350            Reply {
351                provider: provider.to_owned(),
352                // A whole body is the reply's document; a stream's is what
353                // the wire's reassembler rebuilt, or what the decoder of an
354                // operation with free events recorded.
355                raw: self
356                    .document
357                    .clone()
358                    .or_else(|| self.raw.clone())
359                    .unwrap_or(serde_json::Value::Null),
360                provider_request_id: self.request_id.clone(),
361            }
362        }
363
364        /// Fold the absorbed reply with the provider's end into the
365        /// response, or return the response a relay's origin folded. A
366        /// reply the provider did not end is [`ProviderError::Truncated`].
367        pub(crate) fn conclude(self, provider: &str) -> Result<Op::Response, ProviderError> {
368            if let Some(response) = self.response {
369                return Ok(response);
370            }
371            let reply = self.reply(provider);
372            let end = self.end.ok_or(ProviderError::Truncated)?;
373            self.fold.finish(end, reply)
374        }
375    }
376
377    /// What the reply's end closes before it is recorded.
378    pub trait Closing<Op: Operation> {
379        fn close(shared: &mut Shared<Op>);
380    }
381}
382
383/// What a runtime accounts for about a model. Each field matters to the
384/// operations that name it and is left at its default by the others.
385#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
386pub struct Capabilities {
387    /// What a completion provider composes.
388    pub completion: crate::completion::ProviderCapabilities,
389    /// The most documents one embedding or reranking request takes.
390    pub max_documents: usize,
391    /// The dimensionality of the returned vectors. Zero is unknown.
392    pub ndims: usize,
393    /// The embedding width the caller asked for, when they named one. `None`
394    /// and zero disable the check that replies honour it.
395    pub declared: Option<usize>,
396}
397
398impl Capabilities {
399    /// A completion provider's capabilities.
400    pub const fn completion(completion: crate::completion::ProviderCapabilities) -> Self {
401        Self {
402            completion,
403            ..Self::embedding(0, 0)
404        }
405    }
406
407    /// An embedding model's batch limit and width, with no width declared.
408    pub const fn embedding(max_documents: usize, ndims: usize) -> Self {
409        Self {
410            completion: crate::completion::ProviderCapabilities::new(),
411            max_documents,
412            ndims,
413            declared: None,
414        }
415    }
416
417    /// A reranking model's batch limit.
418    pub const fn rerank(max_documents: usize) -> Self {
419        Self::embedding(max_documents, 0)
420    }
421
422    /// Record the width the caller named, when they named one.
423    pub const fn declaring(mut self, declared: Option<usize>) -> Self {
424        self.declared = declared;
425        self
426    }
427
428    /// Returns [`ProviderError::MismatchedDimensions`] for the first width
429    /// differing from a positive caller declaration. Otherwise succeeds.
430    pub(crate) fn honour_declaration(
431        &self,
432        provider: &str,
433        widths: impl IntoIterator<Item = usize>,
434    ) -> Result<(), ProviderError> {
435        // Zero is rig's sentinel for an unknown width, never a claim about
436        // one: a model absent from every table this build knows resolves to
437        // it, and treating that as a declaration would fail every reply.
438        let Some(requested) = self.declared.filter(|declared| *declared > 0) else {
439            return Ok(());
440        };
441        let Some(returned) = widths.into_iter().find(|width| *width != requested) else {
442            return Ok(());
443        };
444        Err(ProviderError::MismatchedDimensions {
445            provider: provider.to_owned(),
446            requested,
447            returned,
448        })
449    }
450}
451
452/// What a wire says about itself: plain data, read before every call.
453#[derive(Debug, Clone)]
454pub struct Descriptor<'a> {
455    /// The model replay shapes a completion history for. Every completion
456    /// wire names one; other operations leave it `None`.
457    pub replay: Option<&'a dyn crate::completion::ReplayTarget>,
458    /// The provider descriptor name (`"anthropic"`), as records and
459    /// telemetry name it.
460    pub name: &'a str,
461    /// The model id the wire addresses, for telemetry. `None` for
462    /// operations that address no model.
463    pub model: Option<&'a str>,
464    /// What a runtime accounts for.
465    pub capabilities: Capabilities,
466    /// The telemetry operation for a call in each mode, when the endpoint
467    /// has a canonical name of its own (Gemini `generate_content`). The mode
468    /// itself is recorded separately, as `gen_ai.request.stream`.
469    pub telemetry: Option<fn(Mode) -> crate::telemetry::GenAiOperation>,
470}
471
472impl<'a> Descriptor<'a> {
473    /// A wire named `name`, addressing no model, with default capabilities.
474    pub fn new(name: &'a str) -> Self {
475        Self {
476            replay: None,
477            name,
478            model: None,
479            capabilities: Capabilities::default(),
480            telemetry: None,
481        }
482    }
483
484    /// The model replay shapes a completion history for.
485    pub fn replay(mut self, target: &'a dyn crate::completion::ReplayTarget) -> Self {
486        self.replay = Some(target);
487        self
488    }
489
490    /// The model id the wire addresses, when it addresses one.
491    pub fn model(mut self, model: impl Into<Option<&'a str>>) -> Self {
492        self.model = model.into();
493        self
494    }
495
496    /// What a runtime accounts for.
497    pub fn capabilities(mut self, capabilities: Capabilities) -> Self {
498        self.capabilities = capabilities;
499        self
500    }
501
502    /// The endpoint's own telemetry operation for each mode.
503    pub fn telemetry(mut self, telemetry: fn(Mode) -> crate::telemetry::GenAiOperation) -> Self {
504        self.telemetry = Some(telemetry);
505        self
506    }
507}
508
509/// One classified wire frame.
510#[derive(Debug)]
511pub enum WireEvent<T> {
512    /// The frame carries a discriminator this client models and its payload
513    /// decoded fully.
514    Known(T),
515    /// Valid JSON not recognized by this classifier.
516    /// Drivers log structural metadata only and skip interpretation.
517    Unknown {
518        /// The unmodeled discriminator value.
519        event_type: String,
520        /// Full payload for raw passthrough, never warning logs. Debug is redacted.
521        value: crate::streaming::UnknownPayload,
522    },
523    /// Invalid JSON or a recognized frame that failed typed decoding.
524    /// Must not be demoted to `Unknown`.
525    Corrupt(serde_json::Error),
526}
527
528impl<T> WireEvent<T> {
529    /// An event a typed transport (an SDK event stream, a gRPC stream)
530    /// reported that this client does not model. The detail is kept for
531    /// raw passthrough, never for warning logs.
532    pub fn unrecognized(event_type: impl Into<String>, detail: impl Into<String>) -> Self {
533        Self::Unknown {
534            event_type: event_type.into(),
535            value: serde_json::Value::String(detail.into()).into(),
536        }
537    }
538
539    /// A modeled event a typed transport failed to decode.
540    pub fn malformed(message: impl std::fmt::Display) -> Self {
541        Self::Corrupt(<serde_json::Error as serde::de::Error>::custom(message))
542    }
543
544    /// Map the `Known` payload, preserving the classification.
545    ///
546    /// This is how an adapter layers a pure event-shape mapping on top of a
547    /// classifier without restating the triage: `Unknown` and `Corrupt` pass
548    /// through untouched, so policy stays with the driver.
549    pub fn map<U>(self, f: impl FnOnce(T) -> U) -> WireEvent<U> {
550        match self {
551            Self::Known(event) => WireEvent::Known(f(event)),
552            Self::Unknown { event_type, value } => WireEvent::Unknown { event_type, value },
553            Self::Corrupt(error) => WireEvent::Corrupt(error),
554        }
555    }
556}
557
558/// Where a decoder writes one reply. The `'id` brand ties it to that one
559/// reply: a writer cannot be kept past the decode step that received it.
560///
561/// ```compile_fail
562/// use rig_core::operation::Completion;
563/// use rig_core::wire::Out;
564///
565/// // A writer cannot be stashed to write into the reply later.
566/// struct Stash(Option<Out<'static, Completion>>);
567///
568/// fn keep<'id>(stash: &mut Stash, out: Out<'id, Completion>) {
569///     stash.0 = Some(out);
570/// }
571/// ```
572pub struct Out<'id, Op: Operation> {
573    shared: &'id std::sync::Mutex<Shared<Op>>,
574    brand: std::marker::PhantomData<fn(&'id ()) -> &'id ()>,
575}
576
577impl<'id, Op: Operation> Out<'id, Op> {
578    pub(crate) fn new(shared: &'id std::sync::Mutex<Shared<Op>>) -> Self {
579        Self {
580            shared,
581            brand: std::marker::PhantomData,
582        }
583    }
584
585    pub(crate) fn lock(&self) -> std::sync::MutexGuard<'id, Shared<Op>> {
586        self.shared
587            .lock()
588            .unwrap_or_else(std::sync::PoisonError::into_inner)
589    }
590
591    /// End the reply with what the provider sent at its end. It consumes the
592    /// writer: nothing is written after the end.
593    pub fn end(self, end: Op::End) -> Flow {
594        let mut shared = self.lock();
595        <Op::Emit as reply::Closing<Op>>::close(&mut shared);
596        shared.end = Some(end);
597        Flow::Ended(Ended(()))
598    }
599
600    /// A payload the provider sent that this decoder does not model. It
601    /// reaches the consumer as [`Item::Unknown`](crate::streaming::Item).
602    pub fn unknown(&mut self, payload: crate::streaming::UnknownPayload) {
603        self.lock()
604            .items
605            .push_back(Ok(crate::streaming::Item::Unknown(payload)));
606    }
607}
608
609impl<Op: Operation<Emit = Free>> Out<'_, Op> {
610    /// Record the reply's provider document, the response's `raw`. A whole
611    /// JSON body the transport reported outranks it: a decoder whose reply
612    /// is not one JSON document records it here.
613    ///
614    /// Only an operation whose decoders build their own events has this. A
615    /// completion's `raw` is what its wire's
616    /// [`Reassembler`](Wire::Reassembler) rebuilds, never what a decoder
617    /// writes:
618    ///
619    /// ```compile_fail,E0599
620    /// use rig_core::operation::Completion;
621    /// use rig_core::wire::Out;
622    ///
623    /// fn record(out: &mut Out<'_, Completion>) {
624    ///     out.raw(serde_json::json!({ "response_id": "x" }));
625    /// }
626    /// ```
627    pub fn raw(&mut self, raw: serde_json::Value) {
628        self.lock().raw = Some(raw);
629    }
630
631    /// One event of the reply.
632    ///
633    /// Only an operation whose decoders build their own events has this; a
634    /// completion's events come from its part handles:
635    ///
636    /// ```compile_fail,E0599
637    /// use rig_core::operation::Completion;
638    /// use rig_core::streaming::StreamEvent;
639    /// use rig_core::wire::Out;
640    ///
641    /// fn reinject(out: &mut Out<'_, Completion>, seen: &StreamEvent) {
642    ///     out.event(seen.clone());
643    /// }
644    /// ```
645    ///
646    /// An event rebuilt from its serialized form is no different:
647    ///
648    /// ```compile_fail,E0599
649    /// use rig_core::operation::Completion;
650    /// use rig_core::streaming::Transcript;
651    /// use rig_core::wire::Out;
652    ///
653    /// fn rebuild(out: &mut Out<'_, Completion>, recorded: serde_json::Value) {
654    ///     let Ok(transcript) = Transcript::parse(recorded) else { return };
655    ///     for event in transcript.events() {
656    ///         out.event(event.clone());
657    ///     }
658    /// }
659    /// ```
660    pub fn event(&mut self, event: Op::Event) {
661        self.lock()
662            .items
663            .push_back(Ok(crate::streaming::Item::Event(event)));
664    }
665}
666
667/// Synchronous state machine for one reply. It classifies frames and
668/// decodes known events into the reply's writer, without transport access
669/// and without being told how the reply arrives. HTTP wires read
670/// [`WireFrame`]s; other transports name their own frame type.
671pub trait Decoder<'id, Op: Operation, Frame = WireFrame> {
672    /// The wire's typed event, produced by this decoder's classifier.
673    type Event;
674
675    /// Decode and classify one frame. A JSON wire MUST delegate to a
676    /// classifier in [`crate::providers::internal::wire`], so the
677    /// decode-then-validate policy is not re-derived per provider.
678    fn classify(&self, frame: Frame) -> WireEvent<Self::Event>;
679
680    /// Write one `Known` event into the reply. An error ends the reply with
681    /// it.
682    fn decode(&mut self, event: Self::Event, out: Out<'id, Op>) -> Result<Flow, ProviderError>;
683
684    /// The frames ran out without an end. A decoder that already saw the
685    /// provider's end, in a form the provider sends before its last frame,
686    /// ends the reply here; otherwise the reply is truncated.
687    fn eof(&mut self, out: Out<'id, Op>) -> Result<Flow, ProviderError> {
688        let _ = out;
689        Err(ProviderError::Truncated)
690    }
691}
692
693/// A provider endpoint: plain data that encodes requests and names a decoder
694/// for each reply.
695///
696/// No transport, no future, no type parameter. Implementations are plain
697/// data (`Clone + PartialEq + Debug + Serialize + Deserialize`, with
698/// credentials held in [`Secret`]), so a host can store one in a scene, a
699/// component, or a config file. `Clone` is a supertrait: the driver clones
700/// the wire into every call's `'static` stream.
701pub trait Wire: Clone + WasmCompatSend + WasmCompatSync + 'static {
702    /// The operation this wire performs.
703    type Op: Operation;
704    /// What [`Self::encode`] produces for the transport to send.
705    type Payload: WasmCompatSend + 'static;
706    /// One unit of a reply, as the transport delivers it.
707    type Frame: WasmCompatSend + 'static;
708    /// The decoder for one of its replies, branded with that reply.
709    type Decoder<'id>: Decoder<'id, Self::Op, Self::Frame> + WasmCompatSend;
710    /// What rebuilds a reply's provider document from its frames when the
711    /// transport reports no whole document, as for a stream. The driver
712    /// feeds it every frame and records what it finishes with as `raw`.
713    /// A wire whose operation's decoders record `raw` themselves names
714    /// [`document::Unreassembled`]; a completion wire cannot, because it
715    /// does not [`Serve`](document::Serves) completions.
716    type Reassembler: document::Reassemble<Self::Frame> + document::Serves<Self::Op>;
717
718    /// What the wire says about itself.
719    fn describe(&self) -> Descriptor<'_>;
720
721    /// The request to send. Pure: it may read `self`, `request` and `mode`,
722    /// and nothing else. A request that cannot be built is an
723    /// [`EncodeError`], which always reports as a request failure.
724    fn encode(&self, request: Request<Self>, mode: Mode) -> Result<Self::Payload, EncodeError>;
725
726    /// A fresh decoder for one reply.
727    fn decoder<'id>(&self) -> Self::Decoder<'id>;
728
729    /// A fresh reassembler for one reply. A wire that picks its API per
730    /// reply builds the matching one.
731    fn reassembler(&self) -> Self::Reassembler {
732        Self::Reassembler::default()
733    }
734}
735
736/// The decoder of a reply that is one JSON document: the document is the
737/// operation's end, and is kept as the reply's `raw`.
738#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
739pub struct Json;
740
741impl<'id, Op> Decoder<'id, Op, WireFrame> for Json
742where
743    Op: Operation<Emit = Free>,
744    Op::End: serde::de::DeserializeOwned,
745{
746    type Event = (Op::End, serde_json::Value);
747
748    fn classify(&self, frame: WireFrame) -> WireEvent<Self::Event> {
749        // Read from the text, so an end holding raw JSON can borrow it.
750        let text = frame.as_str();
751        let document = match serde_json::from_str::<serde_json::Value>(&text) {
752            Ok(document) => document,
753            Err(error) => return WireEvent::Corrupt(error),
754        };
755        match serde_json::from_str::<Op::End>(&text) {
756            Ok(end) => WireEvent::Known((end, document)),
757            Err(error) => WireEvent::Corrupt(error),
758        }
759    }
760
761    fn decode(
762        &mut self,
763        (end, document): Self::Event,
764        mut out: Out<'id, Op>,
765    ) -> Result<Flow, ProviderError> {
766        out.raw(document);
767        Ok(out.end(end))
768    }
769}
770
771/// A wire's request type.
772pub type Request<W> = <<W as Wire>::Op as Operation>::Request;
773/// A wire's response type.
774pub type Response<W> = <<W as Wire>::Op as Operation>::Response;
775/// A wire's event type.
776pub type Event<W> = <<W as Wire>::Op as Operation>::Event;