Skip to main content

rig_core/observe/
adapter.rs

1//! Execution-local correlation and facts from provider request boundaries.
2//!
3//! ```
4//! use rig_core::observe::scrub_diagnostic;
5//!
6//! assert_eq!(scrub_diagnostic("Bearer secret", &[]), "[redacted]");
7//! ```
8
9use std::sync::{Arc, Mutex};
10
11use serde::{Deserialize, Serialize};
12
13use super::{Action, Emitter, Observation, Stage, Subject, Witness};
14
15/// Bounds diagnostic text and removes control characters. Known credentials
16/// or credential markers redact the whole message; oversized messages are
17/// replaced rather than truncated to a potentially sensitive prefix.
18pub fn scrub_diagnostic(value: &str, secrets: &[String]) -> String {
19    scrub::text(value, secrets)
20}
21
22/// The secrets a URL carries for diagnostic redaction: userinfo and known
23/// credential query parameters, origin-form request URIs included. The
24/// returned values are secrets: keep them runtime-only and never write
25/// them into a diagnostic or an artifact.
26pub fn diagnostic_url_secrets(url: &str) -> Vec<String> {
27    scrub::url_secrets(url)
28}
29
30/// One fact about a provider operation, independently of its effect record.
31#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct AdapterObservation {
33    /// Caller-owned logical operation identity within this execution.
34    pub operation: String,
35    /// One-based HTTP send ordinal; absent only when correlation is exhausted.
36    pub attempt: Option<u64>,
37    /// One-based host dispatch attempt, when explicitly supplied by the host.
38    /// Independent of the HTTP send ordinal: a dispatch may send more than once.
39    #[serde(default, skip_serializing_if = "Option::is_none")]
40    pub host_attempt: Option<std::num::NonZeroU64>,
41    /// What the owning boundary observed.
42    pub event: AdapterEvent,
43    /// Scrubbed volatile diagnostics, excluded from semantic comparison.
44    #[serde(default, skip_serializing_if = "Option::is_none")]
45    pub analysis: Option<AdapterAnalysis>,
46}
47
48/// Metadata from the provider boundary. Request/response bodies are not included.
49#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
50#[serde(tag = "event", rename_all = "snake_case")]
51pub enum AdapterEvent {
52    /// Fields actually present in a provider payload; absent fields are not updates.
53    Provider {
54        /// The provider's verdict and selected model, independent of transport closure.
55        verdict: AdapterVerdict,
56    },
57    /// A provider error envelope, including envelopes carried under HTTP 200.
58    ErrorEnvelope {
59        /// Scrubbed original envelope fields, not an inferred HTTP status.
60        error: AdapterErrorEnvelope,
61    },
62    /// Provider-reported usage for this attempt, including rejected responses.
63    Usage {
64        /// A cumulative snapshot: replace earlier snapshots, never sum them.
65        usage: AdapterUsage,
66    },
67    /// A request is about to be sent.
68    Started {
69        /// HTTP method.
70        method: String,
71        /// Provider-declared route template, never a credential-bearing URI.
72        route: String,
73    },
74    /// The transport returned response headers.
75    Response {
76        /// Actual HTTP response status, including successful responses.
77        status: u16,
78    },
79    /// The body reached EOF, independently of any provider terminal verdict.
80    TransportEof {
81        /// Number of complete frames, excluding provider-recognized analysis-only frames.
82        after: usize,
83        /// Undelimited SSE bytes remaining at EOF; zero means no partial frame.
84        partial_bytes: usize,
85    },
86    /// The owning request boundary closed this attempt.
87    Finished {
88        /// Whether the attempt completed, failed, or was dropped.
89        ending: AdapterEnding,
90    },
91    /// A complete transport frame could not be decoded; the driver may continue.
92    Corrupt {
93        /// One-based ordinal counting known, unknown and corrupt frames,
94        /// excluding provider-recognized analysis-only frames.
95        frame: usize,
96    },
97    /// The context cannot assign another unique send ordinal.
98    IdentityExhausted,
99}
100
101/// Sparse provider metadata: a present field supersedes its previous value.
102#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
103pub struct AdapterVerdict {
104    /// Provider finish code, including unknown codes after scrubbing.
105    pub finish_reason: Option<String>,
106    /// Provider prompt-block/refusal reason, when reported.
107    pub block_reason: Option<String>,
108    /// Provider detail accompanying a finish/block reason.
109    pub detail: Option<String>,
110    /// Actual provider model/version, not the requested alias.
111    pub model: Option<String>,
112}
113
114/// Analysis-only response metadata. Never use identifiers as operation join keys.
115#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
116pub struct AdapterAnalysis {
117    /// Provider response identifier, scrubbed and bounded.
118    pub response_id: Option<String>,
119    /// Allowlisted response headers. None means the transport supplied no map;
120    /// an empty map means it supplied no allowlisted values.
121    pub headers: Option<std::collections::BTreeMap<String, String>>,
122}
123
124/// Original provider error-envelope fields, separate from Rig retry classification.
125#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
126pub struct AdapterErrorEnvelope {
127    /// Provider code in its scalar wire spelling; unknown codes stay observable.
128    pub code: Option<String>,
129    /// Provider status name, when present.
130    pub status: Option<String>,
131    /// Bounded, scrubbed provider message, when present.
132    pub message: Option<String>,
133}
134
135/// Error-envelope projection with optional code, message, and `type`/`status`.
136/// Codes remain JSON values until emission validates their scalar shape.
137#[derive(Default, Deserialize)]
138pub struct ObservedError {
139    pub code: Option<serde_json::Value>,
140    #[serde(rename = "type", alias = "status")]
141    pub kind: Option<String>,
142    pub message: Option<String>,
143}
144
145impl ObservedError {
146    /// Emit this envelope as the attempt's error-envelope fact, scrubbed.
147    pub fn emit(self, sink: &mut ObservationSink<'_>) {
148        let code = self.code.map(|code| match code {
149            serde_json::Value::String(code) => sink.scrub(&code),
150            serde_json::Value::Number(code) => code.to_string(),
151            _ => "[invalid]".to_owned(),
152        });
153        sink.emit(AdapterEvent::ErrorEnvelope {
154            error: AdapterErrorEnvelope {
155                code,
156                status: self.kind.map(|value| sink.scrub(&value)),
157                message: self.message.map(|value| sink.scrub(&value)),
158            },
159        });
160    }
161}
162
163/// A provider's cumulative usage snapshot for one HTTP attempt.
164///
165/// Missing, invalid or negative counts remain unknown. A present zero is a
166/// reported zero. These fields may overlap (for example cached tokens are
167/// included in input tokens); never sum fields to invent a total. A later
168/// snapshot replaces the earlier snapshot, including its unknown fields.
169#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
170pub struct AdapterUsage {
171    /// Input tokens, including cached input where the provider includes it.
172    pub input_tokens: Option<u64>,
173    /// Output candidate tokens, as reported by the provider.
174    pub output_tokens: Option<u64>,
175    /// Provider-reported total, not a sum of the other fields.
176    pub total_tokens: Option<u64>,
177    /// Cached input tokens, when reported.
178    pub cached_input_tokens: Option<u64>,
179    /// Reasoning tokens, when reported separately.
180    pub reasoning_tokens: Option<u64>,
181    /// Provider tool-use input tokens, when reported separately.
182    pub tool_input_tokens: Option<u64>,
183}
184
185/// Boundary identified by the adapter's original typed error, not its message.
186#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
187#[serde(rename_all = "snake_case")]
188pub enum AdapterErrorBoundary {
189    /// Request construction or validation failed.
190    Request,
191    /// The provider returned a failure response or error envelope.
192    ProviderResponse,
193    /// Response decoding or validation failed.
194    Decode,
195    /// A typed transport termination was reported.
196    Transport,
197    /// An erased client error does not establish a more specific boundary.
198    #[default]
199    Unknown,
200}
201
202impl AdapterErrorBoundary {
203    pub(crate) fn from_http(error: &crate::http_client::Error) -> Self {
204        use crate::http_client::Error as H;
205        match error {
206            H::Protocol(_) | H::InvalidHeaderValue(_) | H::NoHeaders => Self::Request,
207            H::InvalidContentType(_) => Self::Decode,
208            H::StreamEnded => Self::Transport,
209            H::InvalidStatusCodeWithDetails { .. } => Self::ProviderResponse,
210            H::Instance(_) => Self::Unknown,
211        }
212    }
213}
214
215/// How a provider attempt closed, distinct from the run's ending.
216#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
217#[serde(tag = "ending", rename_all = "snake_case")]
218pub enum AdapterEnding {
219    /// The request produced a decoded response.
220    Decoded,
221    /// The request failed; provider payloads remain in the existing error report.
222    Error {
223        /// Known boundary before error conversion erases transport subtypes.
224        #[serde(default)]
225        boundary: AdapterErrorBoundary,
226        /// Stable Rig error classification.
227        kind: String,
228        /// HTTP status reported by the failure, when available.
229        status: Option<u16>,
230        /// Retryability under the existing policy, not a decision to retry.
231        retryable: bool,
232    },
233    /// A provider terminal record was decoded from the stream.
234    Terminal,
235    /// Transport EOF without a provider terminal record.
236    Eof {
237        /// Number of complete frames, excluding provider-recognized analysis-only frames.
238        after: usize,
239    },
240    /// EOF left an undelimited SSE event; bytes are never retained in the fact.
241    PartialFrame {
242        /// Raw bytes since the last blank event delimiter, saturated at usize::MAX.
243        byte_count: usize,
244        /// Complete frames previously delivered, excluding analysis-only frames.
245        after: usize,
246    },
247    /// The future or stream was dropped before its boundary closed.
248    Dropped,
249}
250
251/// A caller-supplied witness and logical operation identity.
252///
253/// Clone this context for retries of the same operation. Create a new context
254/// with a distinct identity for a different call, including identical parallel
255/// requests. Identity is execution-local; it does not assert replay equivalence.
256/// This handle is runtime-only and must never be serialized into provider data.
257#[derive(Clone)]
258pub struct AdapterContext {
259    inner: Arc<AdapterContextInner>,
260}
261
262struct AdapterContextInner {
263    sink: Arc<dyn Witness + Send + Sync>,
264    subject: Subject,
265    operation: String,
266    next: Arc<Mutex<Option<u64>>>,
267    host_attempt: Option<std::num::NonZeroU64>,
268}
269
270impl std::fmt::Debug for AdapterContext {
271    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
272        // Caller-supplied identity can contain sensitive data; do not log it.
273        f.debug_struct("AdapterContext").finish_non_exhaustive()
274    }
275}
276
277impl AdapterContext {
278    /// Caller-owned identity of this logical operation within the execution.
279    pub fn operation(&self) -> &str {
280        &self.inner.operation
281    }
282
283    /// Bind a logical operation to its witness and existing subject.
284    /// The caller must use a non-sensitive, execution-unique operation identity.
285    pub fn new(
286        sink: Arc<dyn Witness + Send + Sync>,
287        subject: Subject,
288        operation: impl Into<String>,
289    ) -> Self {
290        Self {
291            inner: Arc::new(AdapterContextInner {
292                sink,
293                subject,
294                operation: operation.into(),
295                next: Arc::new(Mutex::new(Some(1))),
296                host_attempt: None,
297            }),
298        }
299    }
300
301    /// Bind an explicit host attempt to its current dispatch subject.
302    ///
303    /// The logical operation, witness and HTTP send counter remain shared.
304    /// Existing clones retain their own subject and host ordinal, so an older
305    /// in-flight attempt cannot be relabeled by a concurrent retry. The host
306    /// owns ordinal allocation and must start a new context for a changed
307    /// logical request rather than reusing this operation identity.
308    pub fn for_host_attempt(&self, subject: Subject, attempt: std::num::NonZeroU64) -> Self {
309        Self {
310            inner: Arc::new(AdapterContextInner {
311                sink: self.inner.sink.clone(),
312                subject,
313                operation: self.inner.operation.clone(),
314                next: self.inner.next.clone(),
315                host_attempt: Some(attempt),
316            }),
317        }
318    }
319
320    /// Begins a send and captures request credentials for diagnostic scrubbing.
321    /// `route` must be a credential-free provider template without base-URL
322    /// prefixes or query data. Returns `None` when attempt IDs are exhausted.
323    pub(crate) fn attempt_for<B>(
324        &self,
325        request: &http::Request<B>,
326        route: &str,
327    ) -> Option<AdapterAttempt> {
328        let mut attempt = self.begin(request.method(), route)?;
329        attempt.secrets = scrub::request_secrets(request);
330        Some(attempt)
331    }
332
333    fn emit(&self, attempt: Option<u64>, event: AdapterEvent) {
334        self.emit_with_analysis(attempt, event, None);
335    }
336
337    fn emit_with_analysis(
338        &self,
339        attempt: Option<u64>,
340        event: AdapterEvent,
341        analysis: Option<AdapterAnalysis>,
342    ) {
343        self.inner.sink.observe(Observation::new(
344            self.inner.subject.clone(),
345            Stage::Handler,
346            Emitter::named("rig-core/adapter"),
347            Action::Adapter {
348                observation: AdapterObservation {
349                    operation: self.inner.operation.clone(),
350                    attempt,
351                    host_attempt: self.inner.host_attempt,
352                    event,
353                    analysis,
354                },
355            },
356        ));
357    }
358
359    /// Allocates a send ordinal and emits its start under a credential-free
360    /// route template. Emits exhaustion at the last ordinal; later calls return
361    /// `None` without reusing identities or affecting provider execution.
362    pub(crate) fn begin(&self, method: &http::Method, route: &str) -> Option<AdapterAttempt> {
363        let number = {
364            let mut next = self
365                .inner
366                .next
367                .lock()
368                .unwrap_or_else(std::sync::PoisonError::into_inner);
369            let number = (*next)?;
370            *next = number.checked_add(1);
371            number
372        };
373        self.emit(
374            Some(number),
375            AdapterEvent::Started {
376                method: method.to_string(),
377                route: route.to_owned(),
378            },
379        );
380        if number == u64::MAX {
381            self.emit(None, AdapterEvent::IdentityExhausted);
382        }
383        Some(AdapterAttempt {
384            context: self.clone(),
385            number,
386            closed: false,
387            response_seen: false,
388            sse_tail: super::sse_tail::SseTail::default(),
389            secrets: Vec::new(),
390            pending_response_id: None,
391            error_boundary: None,
392        })
393    }
394}
395
396/// The request owns this guard until completion or cancellation.
397pub(crate) struct AdapterAttempt {
398    context: AdapterContext,
399    number: u64,
400    closed: bool,
401    response_seen: bool,
402    sse_tail: super::sse_tail::SseTail,
403    secrets: Vec<String>,
404    pending_response_id: Option<String>,
405    error_boundary: Option<AdapterErrorBoundary>,
406}
407
408impl AdapterAttempt {
409    pub(crate) fn text(&self, text: &str) -> String {
410        scrub::text(text, &self.secrets)
411    }
412
413    pub(crate) fn emit_with_analysis(&self, event: AdapterEvent, analysis: AdapterAnalysis) {
414        let analysis = (analysis != AdapterAnalysis::default()).then_some(analysis);
415        self.context
416            .emit_with_analysis(Some(self.number), event, analysis);
417    }
418
419    pub(crate) fn emit(&self, event: AdapterEvent) {
420        self.context.emit(Some(self.number), event);
421    }
422
423    /// Project a payload's facts through the projector that understands it.
424    pub(crate) fn project(&mut self, project: impl FnOnce(&mut ObservationSink<'_>)) {
425        project(&mut ObservationSink { attempt: self });
426    }
427
428    pub(crate) fn provider(&mut self, verdict: AdapterVerdict, response_id: Option<String>) {
429        if response_id.is_some() {
430            self.pending_response_id = response_id;
431        }
432        // An ID-only payload must not manufacture a semantic provider event.
433        // Retain at most one ID until the next verdict or the attempt closure.
434        if verdict != AdapterVerdict::default() {
435            let response_id = self.pending_response_id.take();
436            self.emit_with_analysis(
437                AdapterEvent::Provider { verdict },
438                AdapterAnalysis {
439                    response_id,
440                    ..AdapterAnalysis::default()
441                },
442            );
443        }
444    }
445
446    pub(crate) fn response_with_headers(
447        &mut self,
448        status: http::StatusCode,
449        headers: Option<&http::HeaderMap>,
450    ) {
451        if !self.response_seen {
452            self.response_seen = true;
453            self.emit_with_analysis(
454                AdapterEvent::Response {
455                    status: status.as_u16(),
456                },
457                AdapterAnalysis {
458                    headers: headers.map(|h| scrub::headers(h, &self.secrets)),
459                    ..AdapterAnalysis::default()
460                },
461            );
462        }
463    }
464
465    pub(crate) fn finish(&mut self, ending: AdapterEnding) {
466        if !self.closed {
467            self.closed = true;
468            let response_id = self.pending_response_id.take();
469            self.emit_with_analysis(
470                AdapterEvent::Finished { ending },
471                AdapterAnalysis {
472                    response_id,
473                    ..AdapterAnalysis::default()
474                },
475            );
476        }
477    }
478}
479
480impl Drop for AdapterAttempt {
481    fn drop(&mut self) {
482        self.finish(AdapterEnding::Dropped);
483    }
484}
485
486/// Where a payload projector writes its observation facts.
487///
488/// A projector reads verdicts, usage, ids and error envelopes off a raw
489/// payload before normalization discards them. Text it forwards must go
490/// through [`Self::scrub`]: a payload can echo credentials.
491pub struct ObservationSink<'a> {
492    attempt: &'a mut AdapterAttempt,
493}
494
495impl ObservationSink<'_> {
496    /// Record one boundary fact.
497    pub fn emit(&mut self, event: AdapterEvent) {
498        self.attempt.emit(event);
499    }
500
501    /// Record the provider's verdict, and the response id it named.
502    pub fn provider(&mut self, verdict: AdapterVerdict, response_id: Option<String>) {
503        self.attempt.provider(verdict, response_id);
504    }
505
506    /// Bound and redact diagnostic text from the payload.
507    pub fn scrub(&self, value: &str) -> String {
508        self.attempt.text(value)
509    }
510}
511
512/// Shared by the SSE transport and frame driver, which own different boundaries.
513#[derive(Clone, Default)]
514pub(crate) struct AdapterSlot(Arc<Mutex<Option<AdapterAttempt>>>);
515
516impl AdapterSlot {
517    pub(crate) fn transport_eof(&self, after: usize) {
518        if let Some(attempt) = self
519            .0
520            .lock()
521            .unwrap_or_else(std::sync::PoisonError::into_inner)
522            .as_ref()
523        {
524            attempt.emit(AdapterEvent::TransportEof {
525                after,
526                partial_bytes: attempt.sse_tail.pending(),
527            });
528        }
529    }
530
531    /// Project a reply payload's facts through the projector that
532    /// understands it.
533    pub(crate) fn project(&self, project: impl FnOnce(&mut ObservationSink<'_>)) {
534        if let Some(attempt) = self
535            .0
536            .lock()
537            .unwrap_or_else(std::sync::PoisonError::into_inner)
538            .as_mut()
539        {
540            attempt.project(project);
541        }
542    }
543
544    /// Install the attempt this send's facts belong to.
545    pub(crate) fn install(&self, attempt: Option<AdapterAttempt>) {
546        *self
547            .0
548            .lock()
549            .unwrap_or_else(std::sync::PoisonError::into_inner) = attempt;
550    }
551
552    pub(crate) fn response(&self, status: http::StatusCode) {
553        self.response_with_headers(status, None);
554    }
555
556    pub(crate) fn response_with_headers(
557        &self,
558        status: http::StatusCode,
559        headers: Option<&http::HeaderMap>,
560    ) {
561        if let Some(attempt) = self
562            .0
563            .lock()
564            .unwrap_or_else(std::sync::PoisonError::into_inner)
565            .as_mut()
566        {
567            attempt.response_with_headers(status, headers);
568        }
569    }
570
571    pub(crate) fn finish(&self, ending: AdapterEnding) {
572        if let Some(attempt) = self
573            .0
574            .lock()
575            .unwrap_or_else(std::sync::PoisonError::into_inner)
576            .as_mut()
577        {
578            attempt.finish(ending);
579        }
580    }
581
582    /// Preserve information before the SSE transport converts the owned error.
583    pub(crate) fn error_boundary(&self, boundary: AdapterErrorBoundary) {
584        if let Some(attempt) = self
585            .0
586            .lock()
587            .unwrap_or_else(std::sync::PoisonError::into_inner)
588            .as_mut()
589        {
590            attempt.error_boundary = Some(boundary);
591        }
592    }
593
594    pub(crate) fn fail(&self, error: &crate::error::ProviderError) {
595        if let Some(status) = error.provider_response_status() {
596            self.response(status);
597        }
598        let report = error.report();
599        let boundary = self
600            .0
601            .lock()
602            .unwrap_or_else(std::sync::PoisonError::into_inner)
603            .as_ref()
604            .and_then(|attempt| attempt.error_boundary)
605            .unwrap_or_else(|| error.boundary());
606        self.finish(AdapterEnding::Error {
607            boundary,
608            kind: report.kind.code().to_owned(),
609            status: report.http_status,
610            retryable: report.is_retryable(),
611        });
612    }
613
614    pub(crate) fn bytes(&self, bytes: &[u8]) {
615        if let Some(attempt) = self
616            .0
617            .lock()
618            .unwrap_or_else(std::sync::PoisonError::into_inner)
619            .as_mut()
620        {
621            attempt.sse_tail.feed(bytes);
622        }
623    }
624
625    pub(crate) fn eof(&self, after: usize) {
626        if let Some(attempt) = self
627            .0
628            .lock()
629            .unwrap_or_else(std::sync::PoisonError::into_inner)
630            .as_mut()
631        {
632            let byte_count = attempt.sse_tail.pending();
633            let ending = if byte_count == 0 {
634                AdapterEnding::Eof { after }
635            } else {
636                AdapterEnding::PartialFrame { byte_count, after }
637            };
638            attempt.finish(ending);
639        }
640    }
641
642    pub(crate) fn corrupt(&self, frame: usize) {
643        if let Some(attempt) = self
644            .0
645            .lock()
646            .unwrap_or_else(std::sync::PoisonError::into_inner)
647            .as_ref()
648        {
649            attempt
650                .context
651                .emit(Some(attempt.number), AdapterEvent::Corrupt { frame });
652        }
653    }
654}
655
656mod scrub;
657#[cfg(test)]
658mod tests;