use super::*;
fn write_log(root: &Path, token: &str, records: &[Value]) {
let directory = root.join(token);
std::fs::create_dir_all(&directory).expect("create token directory");
let contents = records
.iter()
.map(|record| serde_json::to_string(record).expect("serialize"))
.collect::<Vec<_>>()
.join("\n");
std::fs::write(directory.join("requests.jsonl"), contents + "\n").expect("write log");
}
#[test]
fn a_truncated_stream_is_found_without_searching_for_error_text() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "a", "phase": "client_request", "uri": "/v1/messages"}),
json!({"correlation_id": "a", "phase": "client_response", "status": 200}),
json!({"correlation_id": "a", "phase": "upstream_response_body", "body": "data: x"}),
json!({
"correlation_id": "a",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"complete": false,
"frames": 444
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(unparsable, 0);
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.incomplete_streams, 1, "{summary:?}");
assert_eq!(summary.statuses.get(&200).copied(), Some(1));
let found = anomalies(&exchanges);
let cut = found
.iter()
.find(|anomaly| anomaly.kind == "stream_ended_without_terminator")
.expect("the truncated stream is named");
assert_eq!(cut.correlation_ids, vec!["a".to_string()]);
}
#[test]
fn a_compressed_body_is_reported_as_undecodable_not_as_a_failure() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "b", "phase": "client_request", "uri": "/v1/messages"}),
json!({"correlation_id": "b", "phase": "client_response", "status": 200}),
json!({
"correlation_id": "b",
"phase": "upstream_response_body",
"body": {"base64": "H4sIAAAA", "bytes": 6}
}),
json!({
"correlation_id": "b",
"phase": "stream_end",
"outcome": "completed",
"complete": true,
"frames": 1
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.incomplete_streams, 0, "{summary:?}");
assert_eq!(summary.undecodable_bodies, 1);
let found = anomalies(&exchanges);
assert!(
found
.iter()
.any(|anomaly| anomaly.kind == "undecodable_bodies"),
"undecodable bodies must be stated, not silently ignored"
);
assert!(
!found
.iter()
.any(|anomaly| anomaly.kind == "stream_ended_without_terminator"),
"a completed stream must not be reported as truncated"
);
}
#[test]
fn unparsable_records_are_reported() {
let root = tempfile::tempdir().expect("temporary log root");
let directory = root.path().join("tokenhash");
std::fs::create_dir_all(&directory).expect("create");
std::fs::write(
directory.join("requests.jsonl"),
"{\"correlation_id\":\"c\",\"phase\":\"client_response\",\"status\":200}\nnot json\n",
)
.expect("write");
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read");
assert_eq!(unparsable, 1);
let rendered = summarise(&exchanges, unparsable, 0).render();
assert!(rendered.contains("1 unparsable"), "{rendered}");
}
#[test]
fn repeated_authentication_failures_are_named() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "d", "phase": "client_response", "status": 401}),
json!({"correlation_id": "e", "phase": "client_response", "status": 401}),
],
);
let (exchanges, _, _) = read_exchanges(root.path(), None).expect("read");
let found = anomalies(&exchanges);
let refused = found
.iter()
.find(|anomaly| anomaly.kind == "repeated_authentication_failure")
.expect("named");
assert_eq!(refused.correlation_ids.len(), 2);
}
#[test]
fn a_healthy_log_reports_nothing() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "f", "phase": "client_response", "status": 200}),
json!({
"correlation_id": "f",
"phase": "stream_end",
"outcome": "completed",
"complete": true,
"frames": 3
}),
],
);
let (exchanges, _, _) = read_exchanges(root.path(), None).expect("read");
assert!(anomalies(&exchanges).is_empty());
}
#[test]
fn a_single_exchange_can_be_shown() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "g", "phase": "client_request", "uri": "/v1/messages"}),
json!({"correlation_id": "h", "phase": "client_request", "uri": "/other"}),
],
);
let rendered = show(root.path(), None, "g").expect("show");
assert!(rendered.contains("/v1/messages"), "{rendered}");
assert!(!rendered.contains("/other"), "{rendered}");
let missing = show(root.path(), None, "nope").expect("show");
assert!(missing.contains("no records"), "{missing}");
}
#[test]
fn a_complete_non_streamed_exchange_is_not_an_anomaly() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "j", "phase": "client_request", "uri": "/v1/messages"}),
json!({"correlation_id": "j", "phase": "upstream_request"}),
json!({
"correlation_id": "j",
"phase": "upstream_response",
"status": 200,
"headers": {"content-type": "application/json"}
}),
json!({"correlation_id": "j", "phase": "upstream_response_body", "body": "{\"ok\":true}"}),
json!({
"correlation_id": "j",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "application/json"}
}),
json!({"correlation_id": "j", "phase": "client_response_body", "body": "{\"ok\":true}"}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(unparsable, 0);
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(
summary.streamed, 0,
"a JSON response is not a stream: {summary:?}"
);
assert_eq!(summary.unterminated_streams, 0, "{summary:?}");
assert_eq!(summary.incomplete_streams, 0, "{summary:?}");
let found = anomalies(&exchanges);
assert!(
found.is_empty(),
"a complete non-streamed exchange must raise nothing: {found:?}"
);
}
#[test]
fn a_gzip_compressed_json_reply_is_not_a_truncated_stream() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "g", "phase": "client_request", "uri": "/v1/messages"}),
json!({
"correlation_id": "g",
"phase": "upstream_response",
"status": 200,
"headers": {"content-type": "application/json", "content-encoding": "gzip"}
}),
json!({"correlation_id": "g", "phase": "upstream_response_body", "body": {"base64": "H4sIAAAAAAAA"}}),
json!({"correlation_id": "g", "phase": "upstream_response_body", "body": {"base64": "A/3VRy07DMBD8"}}),
json!({
"correlation_id": "g",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "application/json", "content-encoding": "gzip"}
}),
json!({
"correlation_id": "g",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"complete": false,
"frames": 2
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(
summary.streamed, 0,
"a compressed JSON reply is not a stream: {summary:?}"
);
assert_eq!(summary.incomplete_streams, 0, "{summary:?}");
assert_eq!(summary.unterminated_streams, 0, "{summary:?}");
let found = anomalies(&exchanges);
assert!(
!found
.iter()
.any(|anomaly| anomaly.kind.starts_with("stream_")
|| anomaly.kind == "no_terminal_record"),
"a successful compressed reply must raise no stream verdict: {found:?}"
);
}
#[test]
fn an_sse_response_is_still_a_stream() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "s", "phase": "client_request", "uri": "/v1/messages"}),
json!({
"correlation_id": "s",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream"}
}),
json!({"correlation_id": "s", "phase": "upstream_response_body", "body": "data: x"}),
json!({
"correlation_id": "s",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"complete": false,
"frames": 12
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.streamed, 1, "{summary:?}");
assert_eq!(
summary.incomplete_streams, 1,
"a truncated SSE stream must still be reported: {summary:?}"
);
let found = anomalies(&exchanges);
assert!(
found
.iter()
.any(|anomaly| anomaly.kind == "stream_ended_without_terminator"),
"{found:?}"
);
}
#[test]
fn a_parameterised_media_type_is_recognised() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "p", "phase": "client_request"}),
json!({
"correlation_id": "p",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream; charset=utf-8"}
}),
json!({"correlation_id": "p", "phase": "upstream_response_body", "body": "data: x"}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(summarise(&exchanges, unparsable, 0).streamed, 1);
}
#[test]
fn a_requested_stream_counts_when_the_response_type_is_unknown() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({
"correlation_id": "r",
"phase": "client_request",
"body": {"json": {"stream": true}}
}),
json!({"correlation_id": "r", "phase": "client_response", "status": 200}),
json!({"correlation_id": "r", "phase": "upstream_response_body", "body": "data: x"}),
json!({
"correlation_id": "r",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"complete": false,
"frames": 3
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.streamed, 1, "{summary:?}");
assert_eq!(summary.incomplete_streams, 1, "{summary:?}");
}
#[test]
fn a_json_answer_to_a_streaming_request_is_not_a_stream() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({
"correlation_id": "m",
"phase": "client_request",
"body": {"json": {"stream": true}}
}),
json!({
"correlation_id": "m",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "application/json"}
}),
json!({"correlation_id": "m", "phase": "upstream_response_body", "body": "{}"}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(
summarise(&exchanges, unparsable, 0).streamed,
0,
"the response decides, not the request"
);
}
#[test]
fn a_compressed_sse_stream_is_not_reported_as_truncated() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({
"correlation_id": "z",
"phase": "client_request",
"uri": "/v1/messages?beta=true",
"body": {"json": {"stream": true}}
}),
json!({
"correlation_id": "z",
"phase": "client_response",
"status": 200,
"headers": {
"content-type": "text/event-stream; charset=utf-8",
"content-encoding": "gzip"
}
}),
json!({"correlation_id": "z", "phase": "upstream_response_body", "body": {"base64": "H4sIAAAAAAAA"}}),
json!({
"correlation_id": "z",
"phase": "stream_end",
"outcome": "encoded_not_verifiable",
"streamed": true,
"inspectable": false,
"complete": false,
"frames": 11
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.streamed, 1, "it is still a stream: {summary:?}");
assert_eq!(
summary.incomplete_streams, 0,
"an unreadable stream is not a demonstrated truncation: {summary:?}"
);
assert_eq!(summary.unterminated_streams, 0, "{summary:?}");
assert_eq!(
summary.unverifiable_streams, 1,
"it is reported as its own class: {summary:?}"
);
let found = anomalies(&exchanges);
assert!(
!found
.iter()
.any(|anomaly| anomaly.kind == "stream_ended_without_terminator"),
"a healthy compressed stream must not be called truncated: {found:?}"
);
let named = found
.iter()
.find(|anomaly| anomaly.kind == "stream_not_verifiable")
.expect("the unverifiable stream is named");
assert_eq!(named.correlation_ids, vec!["z".to_string()]);
}
#[test]
fn a_compressed_stream_is_recognised_from_its_headers() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "h", "phase": "client_request"}),
json!({
"correlation_id": "h",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream", "content-encoding": "gzip"}
}),
json!({"correlation_id": "h", "phase": "upstream_response_body", "body": {"base64": "H4sI"}}),
json!({
"correlation_id": "h",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"complete": false,
"frames": 11
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.unverifiable_streams, 1, "{summary:?}");
assert_eq!(summary.incomplete_streams, 0, "{summary:?}");
}
#[test]
fn an_uncompressed_truncated_stream_is_still_reported() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "t", "phase": "client_request"}),
json!({
"correlation_id": "t",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream"}
}),
json!({"correlation_id": "t", "phase": "upstream_response_body", "body": "data: x"}),
json!({
"correlation_id": "t",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"streamed": true,
"inspectable": true,
"complete": false,
"frames": 444
}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.incomplete_streams, 1, "{summary:?}");
assert_eq!(summary.unverifiable_streams, 0, "{summary:?}");
assert!(
anomalies(&exchanges)
.iter()
.any(|anomaly| anomaly.kind == "stream_ended_without_terminator"),
"a real truncation must still be named"
);
}
fn stream_without_a_terminal_record(id: &str, uri: &str, body: &str) -> Vec<Value> {
vec![
json!({"correlation_id": id, "phase": "client_request", "uri": uri}),
json!({
"correlation_id": id,
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream"}
}),
json!({"correlation_id": id, "phase": "upstream_response_body", "body": body}),
]
}
#[test]
fn a_terminator_in_the_recorded_body_settles_the_ending() {
let root = tempfile::tempdir().expect("temporary log root");
let mut records = Vec::new();
records.extend(stream_without_a_terminal_record(
"openai",
"/v1/chat/completions",
"data: {\"choices\":[{\"delta\":{}}]}\n\ndata: [DONE]\n\n",
));
records.extend(stream_without_a_terminal_record(
"gemini",
"/api/gemini/v1beta/models/gemini-2.0:streamGenerateContent",
"data: {\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"hi\"}]},\
\"finishReason\":\"STOP\",\"index\":0}]}\n\n",
));
records.extend(stream_without_a_terminal_record(
"anthropic",
"/v1/messages",
"event: message_stop\ndata: {\"type\":\"message_stop\"}\n\n",
));
records.extend(stream_without_a_terminal_record(
"responses",
"/v1/responses",
"event: response.completed\ndata: {\"type\":\"response.completed\"}\n\n",
));
write_log(root.path(), "tokenhash", &records);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.streamed, 4, "{summary:?}");
assert_eq!(
summary.unterminated_streams, 0,
"every one of these carries its own dialect's terminator: {summary:?}"
);
let found = anomalies(&exchanges);
assert!(
!found
.iter()
.any(|anomaly| anomaly.kind == "no_terminal_record"),
"a completed stream must not be reported as unknown: {found:?}"
);
}
#[test]
fn a_stream_with_no_terminator_is_still_reported() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&stream_without_a_terminal_record("empty", "/v1/chat/completions", ""),
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(
summary.unterminated_streams, 1,
"an empty body settles nothing: {summary:?}"
);
assert!(
anomalies(&exchanges)
.iter()
.any(|anomaly| anomaly.kind == "no_terminal_record"),
"the genuinely unaccounted-for stream must still be named"
);
}
#[test]
fn a_terminal_record_outranks_the_recorded_body() {
let root = tempfile::tempdir().expect("temporary log root");
let mut records = stream_without_a_terminal_record(
"cut",
"/v1/chat/completions",
"data: {\"choices\":[{\"delta\":{}}]}\n\n",
);
records.push(json!({
"correlation_id": "cut",
"phase": "stream_end",
"outcome": "ended_without_terminator",
"streamed": true,
"inspectable": true,
"complete": false,
"frames": 12
}));
write_log(root.path(), "tokenhash", &records);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(
summary.incomplete_streams, 1,
"the relay saw the stream stop early: {summary:?}"
);
assert_eq!(summary.unterminated_streams, 0, "{summary:?}");
}
#[test]
fn a_client_side_body_can_settle_the_ending() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "c", "phase": "client_request", "uri": "/v1/chat/completions"}),
json!({
"correlation_id": "c",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream"}
}),
json!({"correlation_id": "c", "phase": "client_response_body", "body": "data: [DONE]\n\n"}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(summarise(&exchanges, unparsable, 0).unterminated_streams, 0);
}
#[test]
fn a_compressed_body_is_not_settled_by_its_encoded_text() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"tokenhash",
&[
json!({"correlation_id": "gz", "phase": "client_request", "uri": "/v1/messages"}),
json!({
"correlation_id": "gz",
"phase": "client_response",
"status": 200,
"headers": {"content-type": "text/event-stream", "content-encoding": "gzip"}
}),
json!({"correlation_id": "gz", "phase": "upstream_response_body", "body": {"base64": "bWVzc2FnZV9zdG9w"}}),
],
);
let (exchanges, unparsable, _) = read_exchanges(root.path(), None).expect("read the log");
let summary = summarise(&exchanges, unparsable, 0);
assert_eq!(summary.unverifiable_streams, 1, "{summary:?}");
assert_eq!(summary.unterminated_streams, 0, "{summary:?}");
}
#[test]
fn a_damaged_line_does_not_discard_the_records_around_it() {
let root = tempfile::tempdir().expect("temporary log root");
let directory = root.path().join("tokenhash");
std::fs::create_dir_all(&directory).expect("create token directory");
std::fs::write(
directory.join("requests.jsonl"),
"{\"correlation_id\":\"ok\",\"phase\":\"client_response\",\"status\":200}\n\
not json at all\n\
\n\
{\"correlation_id\":\"ok\",\"phase\":\"upstream_response_body\",\"body\":\"data: [DONE]\"}\n",
)
.expect("write log");
let (exchanges, unparsable, bytes) = read_exchanges(root.path(), None).expect("read the log");
assert_eq!(unparsable, 1, "the damaged line is counted");
assert!(bytes > 0, "the bytes read are reported");
assert_eq!(exchanges.len(), 1, "the readable records still assemble");
assert_eq!(exchanges[0].correlation_id, "ok");
}
#[test]
fn a_token_filter_selects_one_directory() {
let root = tempfile::tempdir().expect("temporary log root");
write_log(
root.path(),
"aaaa",
&[json!({"correlation_id": "a", "phase": "client_response", "status": 200})],
);
write_log(
root.path(),
"bbbb",
&[json!({"correlation_id": "b", "phase": "client_response", "status": 500})],
);
let (all, _, _) = read_exchanges(root.path(), None).expect("read every token");
assert_eq!(all.len(), 2);
let (one, _, _) = read_exchanges(root.path(), Some("aaaa")).expect("read one token");
assert_eq!(one.len(), 1, "only the named token is read");
assert_eq!(one[0].correlation_id, "a");
}
#[test]
fn a_missing_root_yields_nothing() {
let (exchanges, unparsable, bytes) =
read_exchanges(std::path::Path::new("/nonexistent-log-root-258"), None)
.expect("a missing root is not an error");
assert!(exchanges.is_empty());
assert_eq!(unparsable, 0);
assert_eq!(bytes, 0);
}