use super::*;
use crate::completion::CompletionRequest;
use crate::observe::ObservationLog;
#[test]
fn native_http_errors_preserve_distinct_boundaries_before_report_erasure() {
use crate::error::ProviderError;
use crate::http_client::Error;
for (native, boundary) in [
(Error::NoHeaders, AdapterErrorBoundary::Request),
(
Error::InvalidContentType(http::HeaderValue::from_static("text/plain")),
AdapterErrorBoundary::Decode,
),
(Error::StreamEnded, AdapterErrorBoundary::Transport),
(
Error::instance(std::io::Error::other("opaque client failure")),
AdapterErrorBoundary::Unknown,
),
] {
let error = ProviderError::Http(native.into());
let report = crate::error::ErrorReport::from(&error);
assert_eq!(report.kind, crate::error::ErrorKind::Http);
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "call");
let request = http::Request::new(());
let slot = AdapterSlot::default();
slot.install(context.attempt_for(&request, "/completion"));
slot.fail(&error);
drop(slot);
let trace = log.trace();
assert_eq!(trace.observations.len(), 2);
assert!(
matches!(&trace.observations[1].action, Action::Adapter { observation }
if observation.event == AdapterEvent::Finished { ending: AdapterEnding::Error {
boundary, kind: "http".into(), status: None, retryable: report.is_retryable(),
}})
);
assert_eq!(crate::error::ErrorReport::from(&error), report);
}
}
#[test]
fn host_attempts_share_send_ordinals_without_relabeling_in_flight_facts() {
let log = Arc::new(ObservationLog::default());
let operation = AdapterContext::new(log.clone(), Subject::default(), "logical-call");
let first_subject = Subject::scoped("run/1");
let retry_subject = Subject::scoped("run/2");
let first = operation.for_host_attempt(first_subject.clone(), 1.try_into().unwrap());
let retry = operation.for_host_attempt(retry_subject.clone(), 2.try_into().unwrap());
let mut send_one = first.begin(&http::Method::POST, "/completion").unwrap();
let mut send_two = retry.begin(&http::Method::POST, "/completion").unwrap();
let send_three = first
.clone()
.begin(&http::Method::POST, "/completion")
.unwrap();
send_two.finish(AdapterEnding::Decoded);
send_one.finish(AdapterEnding::Error {
boundary: AdapterErrorBoundary::ProviderResponse,
kind: "http".into(),
status: Some(503),
retryable: true,
});
drop(send_three);
let trace = log.trace();
assert_eq!(trace.observations.len(), 6);
let expected = [(1, 1), (2, 2), (3, 1), (2, 2), (1, 1), (3, 1)];
for (o, (send, host)) in trace.observations.iter().zip(expected) {
let Action::Adapter { observation } = &o.action else {
panic!("adapter fact");
};
assert_eq!(observation.operation, "logical-call");
assert_eq!(observation.attempt, Some(send));
assert_eq!(
observation.host_attempt.map(std::num::NonZeroU64::get),
Some(host)
);
assert_eq!(
o.subject,
if host == 1 {
first_subject.clone()
} else {
retry_subject.clone()
}
);
}
let decoded: crate::observe::ObservationTrace =
serde_json::from_str(&serde_json::to_string(&trace).unwrap()).unwrap();
assert_eq!(decoded, trace);
let separate = AdapterContext::new(log.clone(), Subject::default(), "another-call");
drop(separate.begin(&http::Method::POST, "/completion"));
let last = log.trace().observations.pop().unwrap();
assert!(matches!(last.action, Action::Adapter { observation }
if observation.operation == "another-call" && observation.attempt == Some(1) && observation.host_attempt.is_none()));
}
#[test]
fn cloned_context_numbers_attempts_and_closes_once_without_payloads() {
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "operation-1");
let mut first = context
.begin(&http::Method::POST, "/models/{model}")
.unwrap();
first.response_with_headers(http::StatusCode::OK, None);
first.response_with_headers(http::StatusCode::OK, None);
first.finish(AdapterEnding::Decoded);
drop(first);
drop(
context
.clone()
.begin(&http::Method::POST, "/models/{model}"),
);
let trace = log.trace();
assert_eq!(trace.observations.len(), 5);
let facts: Vec<_> = trace
.observations
.iter()
.map(|o| {
let Action::Adapter { observation } = &o.action else {
panic!("adapter fact")
};
observation
})
.collect();
assert_eq!(
facts.iter().map(|f| f.attempt).collect::<Vec<_>>(),
[Some(1), Some(1), Some(1), Some(2), Some(2)]
);
assert_eq!(
facts[4].event,
AdapterEvent::Finished {
ending: AdapterEnding::Dropped
}
);
let encoded = serde_json::to_string(&trace).unwrap();
assert_eq!(
serde_json::from_str::<crate::observe::ObservationTrace>(&encoded).unwrap(),
trace
);
}
#[test]
fn context_is_not_part_of_serialized_completion_requests() {
let request = crate::completion::CompletionRequest::new("hello");
let encoded = serde_json::to_string(&request).unwrap();
assert!(!encoded.contains("local-only-secret"));
assert!(!encoded.contains("observation"));
let decoded: crate::completion::CompletionRequest = serde_json::from_str(&encoded).unwrap();
assert_eq!(
serde_json::to_value(decoded).unwrap(),
serde_json::to_value(request).unwrap()
);
}
#[tokio::test]
async fn gemini_unary_emits_the_actual_http_boundary_without_changing_the_request() {
use crate::test_utils::RecordingHttpClient;
let body = r#"{"candidates":[{"content":{"parts":[{"text":"pong"}],"role":"model"},"finishReason":"STOP"}]}"#;
let http = RecordingHttpClient::new(body);
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new("synthetic-secret-key")
.completion("gemini-test"),
http.clone(),
);
let plain = CompletionRequest::new("hello");
model.call(plain.clone()).await.unwrap();
let log = Arc::new(ObservationLog::default());
let observed = plain;
let context = AdapterContext::new(log.clone(), Subject::default(), "call-1");
model.call_observed(observed, context).await.unwrap();
let requests = http.requests();
assert_eq!(requests.len(), 2);
assert_eq!(requests[0], requests[1]);
let trace = log.trace();
let events: Vec<_> = trace
.observations
.iter()
.map(|o| {
let Action::Adapter { observation } = &o.action else {
panic!("adapter fact")
};
observation.event.clone()
})
.collect();
assert_eq!(
events,
[
AdapterEvent::Started {
method: "POST".into(),
route: "/v1beta/models/gemini-test:generateContent".into()
},
AdapterEvent::Response { status: 200 },
AdapterEvent::Provider {
verdict: AdapterVerdict {
finish_reason: Some("STOP".into()),
..AdapterVerdict::default()
}
},
AdapterEvent::TransportEof {
after: 1,
partial_bytes: 0
},
AdapterEvent::Finished {
ending: AdapterEnding::Terminal
},
]
);
assert!(
!serde_json::to_string(&trace)
.unwrap()
.contains("synthetic-secret-key")
);
}
#[test]
fn exhausted_identity_never_wraps_or_reuses_an_attempt() {
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "last-call");
*context.inner.next.lock().unwrap() = Some(u64::MAX);
drop(context.begin(&http::Method::POST, "/completion"));
assert!(context.begin(&http::Method::POST, "/completion").is_none());
let trace = log.trace();
assert!(trace.observations.iter().any(|o| matches!(
&o.action, Action::Adapter { observation } if observation.event == AdapterEvent::IdentityExhausted
)));
}
#[tokio::test]
async fn unary_failure_facts_preserve_retryability_without_copying_error_bodies() {
use crate::test_utils::RecordingHttpClient;
let http = RecordingHttpClient::with_error(
http::StatusCode::TOO_MANY_REQUESTS,
r#"{"error":{"message":"synthetic-sensitive-body"},"usageMetadata":{"promptTokenCount":3}}"#,
);
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new("synthetic-sensitive-body")
.completion("gemini-test"),
http.clone(),
);
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "retry-operation");
for _ in 0..2 {
let error = model
.call_observed(CompletionRequest::new("hello"), context.clone())
.await
.unwrap_err();
assert!(error.is_retryable());
}
assert_eq!(http.requests().len(), 2);
let trace = log.trace();
assert_eq!(trace.observations.len(), 10);
for (index, observation) in trace.observations.iter().enumerate() {
let Action::Adapter { observation } = &observation.action else {
panic!("adapter fact")
};
assert_eq!(observation.attempt, Some((index / 5 + 1) as u64));
if index % 5 == 2 {
assert_eq!(
observation.event,
AdapterEvent::Usage {
usage: AdapterUsage {
input_tokens: Some(3),
..AdapterUsage::default()
}
}
);
}
if index % 5 == 3 {
assert_eq!(
observation.event,
AdapterEvent::ErrorEnvelope {
error: AdapterErrorEnvelope {
message: Some("[redacted]".into()),
..AdapterErrorEnvelope::default()
}
}
);
}
if index % 5 == 4 {
assert_eq!(
observation.event,
AdapterEvent::Finished {
ending: AdapterEnding::Error {
boundary: AdapterErrorBoundary::ProviderResponse,
kind: "provider_response".into(),
status: Some(429),
retryable: true
}
}
);
}
}
assert!(
!serde_json::to_string(&trace)
.unwrap()
.contains("synthetic-sensitive-body")
);
}
async fn observed_stream(bytes: &str, stop_after_first: bool) -> crate::observe::ObservationTrace {
use crate::test_utils::MockStreamingClient;
use futures::StreamExt;
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new("test-key").completion("gemini-test"),
MockStreamingClient {
sse_bytes: bytes::Bytes::copy_from_slice(bytes.as_bytes()),
},
);
let log = Arc::new(ObservationLog::default());
let request = CompletionRequest::new("hello");
let mut plain_stream = model.stream(request.clone()).unwrap();
let mut plain_items = Vec::new();
while let Some(item) = plain_stream.next().await {
plain_items.push(
item.map(|event| serde_json::to_value(event).unwrap())
.map_err(|e| e.to_string()),
);
if stop_after_first {
break;
}
}
drop(plain_stream);
let context = AdapterContext::new(log.clone(), Subject::default(), "stream-call");
let mut stream = model.stream_observed(request, context).unwrap();
assert!(
log.is_empty(),
"an unpolled lazy stream has not sent a request"
);
let mut observed_items = Vec::new();
while let Some(item) = stream.next().await {
observed_items.push(
item.map(|event| serde_json::to_value(event).unwrap())
.map_err(|e| e.to_string()),
);
if stop_after_first {
break;
}
}
drop(stream);
assert_eq!(
observed_items, plain_items,
"observation preserves stream semantics"
);
log.trace()
}
#[tokio::test]
async fn stream_terminal_eof_error_and_drop_have_distinct_closures() {
let content = "data: {\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"hi\"}],\"role\":\"model\"},\"index\":0}]}\n\n";
let terminal = "data: {\"candidates\":[{\"finishReason\":\"STOP\",\"index\":0}]}\n\n";
let error = "data: {\"error\":{\"code\":503,\"message\":\"unavailable\",\"status\":\"UNAVAILABLE\"}}\n\n";
for (bytes, stop, ending) in [
(
format!("{content}{terminal}"),
false,
AdapterEnding::Terminal,
),
(content.to_owned(), false, AdapterEnding::Eof { after: 1 }),
(
format!("{content}data: {{"),
false,
AdapterEnding::PartialFrame {
byte_count: 7,
after: 1,
},
),
(
error.to_owned(),
false,
AdapterEnding::Error {
boundary: AdapterErrorBoundary::ProviderResponse,
kind: "provider_response".into(),
status: Some(503),
retryable: true,
},
),
(format!("{content}{terminal}"), true, AdapterEnding::Dropped),
] {
let trace = observed_stream(&bytes, stop).await;
let events: Vec<_> = trace
.observations
.iter()
.map(|o| {
let Action::Adapter { observation } = &o.action else {
panic!("adapter fact")
};
assert_eq!(observation.attempt, Some(1));
&observation.event
})
.collect();
assert_eq!(
events[0],
&AdapterEvent::Started {
method: "POST".into(),
route: "/v1beta/models/gemini-test:streamGenerateContent".into()
}
);
assert_eq!(events[1], &AdapterEvent::Response { status: 200 });
assert_eq!(events.last().unwrap(), &&AdapterEvent::Finished { ending });
assert_eq!(
events
.iter()
.filter(|e| matches!(e, AdapterEvent::Finished { .. }))
.count(),
1
);
}
}
#[tokio::test]
async fn corrupt_frame_is_evidence_separate_from_recovery_or_consumer_drop() {
let corrupt = "data: {\n\n";
let terminal = "data: {\"candidates\":[{\"finishReason\":\"STOP\",\"index\":0}]}\n\n";
let ending = AdapterEnding::Error {
boundary: AdapterErrorBoundary::Decode,
kind: "json".into(),
status: None,
retryable: false,
};
for stop in [false, true] {
let ending = ending.clone();
let trace = observed_stream(&format!("{corrupt}{terminal}"), stop).await;
let events: Vec<_> = trace
.observations
.iter()
.filter_map(|o| {
let Action::Adapter { observation } = &o.action else {
return None;
};
Some(&observation.event)
})
.collect();
assert_eq!(events[2], &AdapterEvent::Corrupt { frame: 1 });
assert_eq!(events.last().unwrap(), &&AdapterEvent::Finished { ending });
assert_eq!(events.len(), 4, "{events:?}");
assert!(!serde_json::to_string(&trace).unwrap().contains("data:"));
}
}
#[tokio::test]
async fn empty_unary_rejection_preserves_optional_usage_before_failure() {
use crate::test_utils::RecordingHttpClient;
for (metadata, expected) in [
(
serde_json::json!({"promptTokenCount": 7, "totalTokenCount": 9}),
Some(AdapterUsage {
input_tokens: Some(7),
total_tokens: Some(9),
..AdapterUsage::default()
}),
),
(
serde_json::json!({"candidatesTokenCount": 0, "promptTokenCount": -1}),
Some(AdapterUsage {
output_tokens: Some(0),
..AdapterUsage::default()
}),
),
(serde_json::Value::Null, None),
] {
let body = serde_json::json!({"candidates": [], "usageMetadata": metadata}).to_string();
let http = RecordingHttpClient::new(body);
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new("test-key").completion("gemini-test"),
http.clone(),
);
let request = CompletionRequest::new("hello");
let plain_error = model.call(request.clone()).await.unwrap_err();
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "empty-call");
let error = model.call_observed(request, context).await.unwrap_err();
assert_eq!(error.to_string(), plain_error.to_string());
assert_eq!(http.requests()[0], http.requests()[1]);
let trace = log.trace();
let usage: Vec<_> = trace
.observations
.iter()
.filter_map(|o| match &o.action {
Action::Adapter {
observation:
AdapterObservation {
event: AdapterEvent::Usage { usage },
..
},
} => Some(usage.clone()),
_ => None,
})
.collect();
assert_eq!(usage, expected.into_iter().collect::<Vec<_>>());
assert!(matches!(&trace.observations.last().unwrap().action,
Action::Adapter { observation } if observation.event == AdapterEvent::Finished { ending: AdapterEnding::Eof { after: 1 } }
));
}
}
#[tokio::test]
async fn streaming_http_rejection_preserves_usage_and_the_original_error() {
use crate::test_utils::HttpErrorStreamingClient;
use futures::StreamExt;
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new("synthetic-sensitive-body")
.completion("gemini-test"),
HttpErrorStreamingClient::new(
http::StatusCode::TOO_MANY_REQUESTS,
r#"{"error":{"message":"synthetic-sensitive-body"},"usageMetadata":{"promptTokenCount":3}}"#,
),
);
let request = CompletionRequest::new("hello");
let mut plain = model.stream(request.clone()).unwrap();
let plain_error = plain.next().await.unwrap().unwrap_err();
assert!(plain.next().await.is_none());
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "rejected-stream");
let mut stream = model.stream_observed(request, context).unwrap();
let error = stream.next().await.unwrap().unwrap_err();
assert_eq!(error.to_string(), plain_error.to_string());
assert!(stream.next().await.is_none());
drop(stream);
let trace = log.trace();
assert_eq!(trace.observations.len(), 5);
assert!(matches!(&trace.observations[2].action,
Action::Adapter { observation } if observation.event == AdapterEvent::Usage { usage: AdapterUsage { input_tokens: Some(3), ..AdapterUsage::default() } }
));
assert!(matches!(&trace.observations[4].action,
Action::Adapter { observation } if matches!(observation.event, AdapterEvent::Finished { ending: AdapterEnding::Error { status: Some(429), retryable: true, .. } })
));
assert!(
!serde_json::to_string(&trace)
.unwrap()
.contains("synthetic-sensitive-body")
);
}
#[tokio::test]
async fn provider_metadata_and_headers_are_scrubbed_before_observation() {
use crate::test_utils::RecordingHttpClient;
let secret = "synthetic-credential-12345";
let body = serde_json::json!({
"candidates": [{"finishReason": "MAX_TOKENS", "finishMessage": secret}],
"modelVersion": "gemini-test", "responseId": secret,
"error": {"code": "FUTURE_STATUS", "status": "RESOURCE_EXHAUSTED", "message": secret}
})
.to_string();
let mut headers = http::HeaderMap::new();
headers.insert("retry-after", "2".parse().unwrap());
headers.insert("x-request-id", secret.parse().unwrap());
headers.insert(
"x-ratelimit-reset-tokens",
"x".repeat(4096).parse().unwrap(),
);
headers.insert("set-cookie", secret.parse().unwrap());
headers.insert("authorization", secret.parse().unwrap());
headers.insert("x-private-detail", secret.parse().unwrap());
for returned_response in [false, true] {
let http = if returned_response {
RecordingHttpClient::with_error_response_headers(
http::StatusCode::TOO_MANY_REQUESTS,
body.clone(),
headers.clone(),
)
} else {
RecordingHttpClient::with_error_headers(
http::StatusCode::TOO_MANY_REQUESTS,
body.clone(),
headers.clone(),
)
};
let model = crate::driver::Model::new(
crate::providers::gemini::GeminiConfig::new(secret).completion("gemini-test"),
http,
);
let log = Arc::new(ObservationLog::default());
let request = CompletionRequest::new("hello");
let context = AdapterContext::new(log.clone(), Subject::default(), "metadata-call");
assert!(
model
.call_observed(request, context)
.await
.unwrap_err()
.is_retryable()
);
let trace = log.trace();
let facts: Vec<_> = trace
.observations
.iter()
.filter_map(|o| {
let Action::Adapter { observation } = &o.action else {
return None;
};
Some(observation)
})
.collect();
assert_eq!(facts.len(), 5);
let captured = facts[1]
.analysis
.as_ref()
.unwrap()
.headers
.as_ref()
.unwrap();
assert_eq!(captured.len(), 3);
assert_eq!(captured["retry-after"], "2");
assert_eq!(captured["x-request-id"], "[redacted]");
assert_eq!(captured["x-ratelimit-reset-tokens"], "[truncated]");
assert_eq!(
facts[2].event,
AdapterEvent::Provider {
verdict: AdapterVerdict {
finish_reason: Some("MAX_TOKENS".into()),
detail: Some("[redacted]".into()),
model: Some("gemini-test".into()),
..AdapterVerdict::default()
}
}
);
assert_eq!(
facts[2].analysis.as_ref().unwrap().response_id.as_deref(),
Some("[redacted]")
);
assert_eq!(
facts[3].event,
AdapterEvent::ErrorEnvelope {
error: AdapterErrorEnvelope {
code: Some("FUTURE_STATUS".into()),
status: Some("RESOURCE_EXHAUSTED".into()),
message: Some("[redacted]".into())
}
}
);
let serialized = serde_json::to_string(&trace).unwrap();
assert!(!serialized.contains(secret));
assert!(!serialized.contains("set-cookie"));
let mut changed = trace.clone();
for o in &mut changed.observations {
if let Action::Adapter { observation } = &mut o.action {
observation.analysis = None;
}
}
assert!(matches!(
crate::test_utils::observations::compare(&trace, &changed),
crate::test_utils::observations::Comparison::Equal
));
if let Action::Adapter { observation } = &mut changed.observations[1].action {
observation.event = AdapterEvent::Response { status: 503 };
}
assert!(matches!(
crate::test_utils::observations::compare(&trace, &changed),
crate::test_utils::observations::Comparison::Diverged(_)
));
}
}
#[test]
fn scrubbing_controls_cannot_reconstitute_a_request_credential() {
let secrets = vec!["credential-value".into()];
assert_eq!(
super::scrub::text("credential-\nvalue", &secrets),
"[redacted]"
);
assert_eq!(
super::scrub::text("https://example.invalid/?key=escaped%20value", &[]),
"[redacted]"
);
assert_eq!(super::scrub::text(&"🦀".repeat(200), &[]), "[truncated]");
}
#[test]
fn url_credentials_are_scrubbed_in_messages_and_allowlisted_header_echoes() {
for (uri, echoes) in [
(
"https://name%2Dpiece:opaque%2Dvalue@example.invalid/v1",
vec![
"name-piece",
"name%2dpiece",
"opaque-value",
"opaque%2Dvalue",
],
),
(
"https://alice:opaque%252Dvalue@example.invalid/v1",
vec!["opaque%2Dvalue", "opaque%252Dvalue"],
),
(
"https://example.invalid/v1?API_KEY=opaque%2Dvalue",
vec!["opaque-value", "opaque%2dvalue"],
),
(
"/v1?%61ccess_token=opaque%2Dvalue",
vec!["opaque-value", "opaque%2Dvalue"],
),
(
"/v1?key=opaque+value",
vec!["opaque value", "opaque+value", "opaque%20value"],
),
(
"https://opaque+name:opaque%2Bvalue@example.invalid/v1",
vec!["opaque+name", "opaque+value"],
),
] {
let request = http::Request::builder().uri(uri).body(()).unwrap();
let secrets = super::scrub::request_secrets(&request);
for echo in echoes {
assert_eq!(
super::scrub::text(echo, &secrets),
"[redacted]",
"{uri}: {echo}"
);
let mut headers = http::HeaderMap::new();
headers.insert("request-id", echo.parse().unwrap());
assert_eq!(
super::scrub::headers(&headers, &secrets)["request-id"],
"[redacted]"
);
}
}
let secrets = super::scrub::url_secrets("/v1?search=ordinary-value&key=");
assert!(secrets.is_empty());
assert_eq!(
super::scrub::text("ordinary%20value", &secrets),
"ordinary%20value"
);
}
#[test]
fn diagnostic_comparison_normalizes_both_sides_without_persisting_decoded_text() {
for (secret, value) in [
("opaque-\nvalue", "opaque-value"),
("opaque-value", "opaque%2dvalue"),
("opaque-value", "opaque-%0Avalue"),
("opaque-%0Avalue", "opaque-value"),
("opaque%2Dvalue", "opaque%252Dvalue"),
("opaque%\n2Dvalue", "opaque-value"),
] {
assert_eq!(super::scrub::text(value, &[secret.into()]), "[redacted]");
}
assert_eq!(super::scrub::text("ordinary", &["\n".into()]), "ordinary");
assert_eq!(super::scrub::text("ordinary%2", &[]), "ordinary%2");
assert_eq!(super::scrub::text("%61pi_key=opaque", &[]), "[redacted]");
assert_eq!(super::scrub::text(&"%41".repeat(200), &[]), "[truncated]");
}
#[tokio::test]
async fn optional_response_ids_do_not_create_semantic_stream_events() {
let content = serde_json::json!({
"candidates": [{"content": {"parts": [{"text": "hi"}], "role": "model"}, "index": 0}]
});
let terminal = "data: {\"candidates\":[{\"finishReason\":\"STOP\",\"index\":0}]}\n\n";
for (suffix, stop) in [("", false), (terminal, false), (terminal, true)] {
let baseline = observed_stream(&format!("data: {content}\n\n{suffix}"), stop).await;
for id in ["response-one", "response-two", "test-key"] {
let mut identified = content.clone();
identified["responseId"] = id.into();
let actual = observed_stream(&format!("data: {identified}\n\n{suffix}"), stop).await;
assert_eq!(
crate::test_utils::observations::compare(&baseline, &actual),
crate::test_utils::observations::Comparison::Equal
);
let ids: Vec<_> = actual
.observations
.iter()
.filter_map(|o| {
let Action::Adapter { observation } = &o.action else {
return None;
};
observation.analysis.as_ref()?.response_id.as_deref()
})
.collect();
assert_eq!(ids, [if id == "test-key" { "[redacted]" } else { id }]);
}
}
let baseline = observed_stream("", false).await;
let actual = observed_stream("data: {\"responseId\":\"response-only\"}\n\n", false).await;
assert_eq!(
crate::test_utils::observations::compare(&baseline, &actual),
crate::test_utils::observations::Comparison::Equal
);
assert!(actual.observations.iter().any(|o| matches!(
&o.action, Action::Adapter { observation }
if observation.analysis.as_ref().and_then(|a| a.response_id.as_deref()) == Some("response-only")
)));
for suffix in ["", "data: {broken\n\n", "data: partial"] {
let baseline = observed_stream(&format!("data: {content}\n\n{suffix}"), false).await;
let actual = observed_stream(
&format!("data: {{\"responseId\":\"first\"}}\n\ndata: {content}\n\ndata: {{\"responseId\":\"last\"}}\n\n{suffix}"),
false,
).await;
assert_eq!(
crate::test_utils::observations::compare(&baseline, &actual),
crate::test_utils::observations::Comparison::Equal
);
}
for extra in [
"{}",
"{\"unknown\":1}",
"{\"responseId\":42}",
"{\"responseId\":\"a\",\"responseId\":\"b\"}",
"{\"usageMetadata\":{\"totalTokenCount\":1}}",
] {
let actual = observed_stream(&format!("data: {extra}\n\n"), false).await;
assert_ne!(
crate::test_utils::observations::compare(&baseline, &actual),
crate::test_utils::observations::Comparison::Equal
);
assert!(actual.observations.iter().all(|fact| !matches!(
&fact.action, Action::Adapter { observation }
if observation.event == AdapterEvent::Corrupt { frame: 0 }
)));
}
}
#[tokio::test]
async fn a_transport_that_records_nothing_writes_no_facts() {
use crate::test_utils::{MockCompletionModel, MockStreamEvent, MockTurn};
use futures::StreamExt;
let log = Arc::new(ObservationLog::default());
let context = AdapterContext::new(log.clone(), Subject::default(), "ordinary-call");
let request = CompletionRequest::new("hello");
let response = MockCompletionModel::from_turns([MockTurn::text("ordinary")])
.call_observed(request.clone(), context.clone())
.await
.unwrap();
assert!(
serde_json::to_string(&response.choice)
.unwrap()
.contains("ordinary")
);
let mut stream = MockCompletionModel::from_stream_turns([[
MockStreamEvent::text("ordinary stream"),
MockStreamEvent::final_response_with_default_usage(),
]])
.stream_observed(request, context)
.unwrap();
let mut events = Vec::new();
while let Some(event) = stream.next().await {
events.push(event.unwrap());
}
assert!(
serde_json::to_string(&events)
.unwrap()
.contains("ordinary stream")
);
assert!(log.trace().observations.is_empty());
}
#[test]
fn a_host_scrubs_its_own_diagnostics_with_the_adapter_rules() {
let secrets = crate::observe::diagnostic_url_secrets(
"https://user:hunter2@example.invalid/v1?key=abc123&search=fine",
);
assert!(secrets.iter().any(|s| s == "hunter2"), "{secrets:?}");
assert!(secrets.iter().any(|s| s == "abc123"), "{secrets:?}");
assert!(!secrets.iter().any(|s| s == "fine"), "{secrets:?}");
let scrubbed = crate::observe::scrub_diagnostic("request to key=abc123 failed", &secrets);
assert!(!scrubbed.contains("abc123"), "{scrubbed}");
assert_eq!(
crate::observe::scrub_diagnostic("plain failure", &secrets),
"plain failure"
);
}