1use std::sync::{Arc, Mutex};
10
11use serde::{Deserialize, Serialize};
12
13use super::{Action, Emitter, Observation, Stage, Subject, Witness};
14
15pub fn scrub_diagnostic(value: &str, secrets: &[String]) -> String {
19 scrub::text(value, secrets)
20}
21
22pub fn diagnostic_url_secrets(url: &str) -> Vec<String> {
27 scrub::url_secrets(url)
28}
29
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct AdapterObservation {
33 pub operation: String,
35 pub attempt: Option<u64>,
37 #[serde(default, skip_serializing_if = "Option::is_none")]
40 pub host_attempt: Option<std::num::NonZeroU64>,
41 pub event: AdapterEvent,
43 #[serde(default, skip_serializing_if = "Option::is_none")]
45 pub analysis: Option<AdapterAnalysis>,
46}
47
48#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
50#[serde(tag = "event", rename_all = "snake_case")]
51pub enum AdapterEvent {
52 Provider {
54 verdict: AdapterVerdict,
56 },
57 ErrorEnvelope {
59 error: AdapterErrorEnvelope,
61 },
62 Usage {
64 usage: AdapterUsage,
66 },
67 Started {
69 method: String,
71 route: String,
73 },
74 Response {
76 status: u16,
78 },
79 TransportEof {
81 after: usize,
83 partial_bytes: usize,
85 },
86 Finished {
88 ending: AdapterEnding,
90 },
91 Corrupt {
93 frame: usize,
96 },
97 IdentityExhausted,
99}
100
101#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
103pub struct AdapterVerdict {
104 pub finish_reason: Option<String>,
106 pub block_reason: Option<String>,
108 pub detail: Option<String>,
110 pub model: Option<String>,
112}
113
114#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
116pub struct AdapterAnalysis {
117 pub response_id: Option<String>,
119 pub headers: Option<std::collections::BTreeMap<String, String>>,
122}
123
124#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
126pub struct AdapterErrorEnvelope {
127 pub code: Option<String>,
129 pub status: Option<String>,
131 pub message: Option<String>,
133}
134
135#[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 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#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
170pub struct AdapterUsage {
171 pub input_tokens: Option<u64>,
173 pub output_tokens: Option<u64>,
175 pub total_tokens: Option<u64>,
177 pub cached_input_tokens: Option<u64>,
179 pub reasoning_tokens: Option<u64>,
181 pub tool_input_tokens: Option<u64>,
183}
184
185#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
187#[serde(rename_all = "snake_case")]
188pub enum AdapterErrorBoundary {
189 Request,
191 ProviderResponse,
193 Decode,
195 Transport,
197 #[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
217#[serde(tag = "ending", rename_all = "snake_case")]
218pub enum AdapterEnding {
219 Decoded,
221 Error {
223 #[serde(default)]
225 boundary: AdapterErrorBoundary,
226 kind: String,
228 status: Option<u16>,
230 retryable: bool,
232 },
233 Terminal,
235 Eof {
237 after: usize,
239 },
240 PartialFrame {
242 byte_count: usize,
244 after: usize,
246 },
247 Dropped,
249}
250
251#[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 f.debug_struct("AdapterContext").finish_non_exhaustive()
274 }
275}
276
277impl AdapterContext {
278 pub fn operation(&self) -> &str {
280 &self.inner.operation
281 }
282
283 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 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 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 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
396pub(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 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 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
486pub struct ObservationSink<'a> {
492 attempt: &'a mut AdapterAttempt,
493}
494
495impl ObservationSink<'_> {
496 pub fn emit(&mut self, event: AdapterEvent) {
498 self.attempt.emit(event);
499 }
500
501 pub fn provider(&mut self, verdict: AdapterVerdict, response_id: Option<String>) {
503 self.attempt.provider(verdict, response_id);
504 }
505
506 pub fn scrub(&self, value: &str) -> String {
508 self.attempt.text(value)
509 }
510}
511
512#[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 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 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 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;