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;