1use serde_json::{Map, Value};
2use std::fs::{self, File, create_dir_all};
3use std::io::{self, Write};
4use std::path::{Path, PathBuf};
5use std::sync::{Mutex, MutexGuard};
6
7use crate::config::AliasProvider;
8use crate::logging::REDACT_KEYS;
9use crate::paths;
10
11#[derive(Debug)]
12pub struct TrafficCapture {
13 root: PathBuf,
14 artifact_counter: Mutex<usize>,
15 event_counter: Mutex<usize>,
16}
17
18pub const MAX_SSE_CAPTURE_BYTES: usize = 8 * 1024 * 1024;
19pub const MAX_STREAM_CAPTURE_EVENT_BYTES: usize = 8 * 1024 * 1024;
20pub const MAX_STREAM_CAPTURE_EVENTS: usize = 1_024;
21pub const MAX_STREAM_CAPTURE_FRAME_BYTES: usize = 64 * 1024;
22
23#[derive(Debug)]
24pub struct TrafficCaptureOptions {
25 pub req_id: String,
26 pub session_id: Option<String>,
27 pub session_seq: Option<u64>,
28 pub provider: Option<String>,
29 pub state_dir_override: Option<PathBuf>,
30}
31
32pub fn traffic_capture_enabled() -> bool {
33 traffic_capture_enabled_for_env(&std::env::vars().collect())
34}
35
36pub fn traffic_capture_enabled_for_env(env: &std::collections::HashMap<String, String>) -> bool {
37 match env.get("CCP_TRAFFIC_LOG").map(String::as_str) {
38 Some(v) => matches!(v, "1" | "true" | "yes"),
39 None => false,
40 }
41}
42
43pub fn create_traffic_capture(opts: TrafficCaptureOptions) -> Option<TrafficCapture> {
44 if !traffic_capture_enabled() {
45 return None;
46 }
47 let state_root = opts
48 .state_dir_override
49 .unwrap_or_else(paths::state_dir)
50 .join("traffic")
51 .join(sanitize_path_part(
52 opts.session_id.as_deref().unwrap_or("no-session"),
53 ))
54 .join(format!(
55 "{:06}-{}-{}",
56 opts.session_seq.unwrap_or(0),
57 sanitize_path_part(opts.provider.as_deref().unwrap_or("unknown-provider")),
58 sanitize_path_part(&opts.req_id),
59 ));
60
61 Some(TrafficCapture {
62 root: state_root,
63 artifact_counter: Mutex::new(0),
64 event_counter: Mutex::new(0),
65 })
66}
67
68impl TrafficCapture {
69 pub fn root(&self) -> &Path {
70 &self.root
71 }
72
73 pub fn write_json(&self, name: &str, value: &Value) {
74 let value = redact_traffic(value);
75 let payload = serde_json::to_string_pretty(&value)
76 .unwrap_or_else(|_| "{}".to_string())
77 .into_bytes();
78 let path = self.next_artifact_path(name, true);
79 let _ = write_bytes(path, &payload);
80 }
81
82 pub fn write_text(&self, name: &str, text: &str) {
83 let file = if name.ends_with(".txt") {
84 name.to_string()
85 } else {
86 format!("{name}.txt")
87 };
88 let path = self.next_artifact_path(&file, false);
89 let _ = write_bytes(path, text.as_bytes());
90 }
91
92 pub fn write_bytes(&self, name: &str, value: &[u8]) {
93 let path = self.next_artifact_path(name, false);
94 let _ = write_bytes(path, value);
95 }
96
97 pub fn write_json_event(&self, name: &str, value: &Value) {
98 let value = redact_traffic(value);
99 let payload = serde_json::to_string_pretty(&value)
100 .unwrap_or_else(|_| "{}".to_string())
101 .into_bytes();
102 let path = self.next_event_path(name, true);
103 let _ = write_bytes(path, &payload);
104 }
105
106 pub fn stream_capture(&self) -> StreamTrafficCapture {
107 StreamTrafficCapture::default()
108 }
109
110 fn next_artifact_path(&self, name: &str, ensure_ext_json: bool) -> PathBuf {
111 let mut counter: MutexGuard<'_, usize> = self
112 .artifact_counter
113 .lock()
114 .unwrap_or_else(|_| self.artifact_counter.lock().unwrap());
115 *counter += 1;
116 let file = if ensure_ext_json && !name.ends_with(".json") {
117 format!("{name}.json")
118 } else {
119 name.to_string()
120 };
121 self.root
122 .join(format!("{:03}-{}", *counter, sanitize_path_part(&file)))
123 }
124
125 fn next_event_path(&self, name: &str, ensure_ext_json: bool) -> PathBuf {
126 let mut counter: MutexGuard<'_, usize> = self
127 .event_counter
128 .lock()
129 .unwrap_or_else(|_| self.event_counter.lock().unwrap());
130 *counter += 1;
131 let file = if ensure_ext_json && !name.ends_with(".json") {
132 format!("{name}.json")
133 } else {
134 name.to_string()
135 };
136 self.root
137 .join("events")
138 .join(format!("{:06}-{}", *counter, sanitize_path_part(&file)))
139 }
140}
141
142#[cfg(test)]
143pub(crate) fn test_capture(root: PathBuf) -> TrafficCapture {
144 TrafficCapture {
145 root,
146 artifact_counter: Mutex::new(0),
147 event_counter: Mutex::new(0),
148 }
149}
150
151fn write_bytes(path: PathBuf, value: &[u8]) -> io::Result<()> {
152 if let Some(parent) = path.parent() {
153 create_dir_all(parent)?;
154 if let Ok(meta) = fs::metadata(parent) {
155 set_mode(parent, 0o700);
156 if meta.is_dir() {
157 #[cfg(unix)]
158 {
159 use std::os::unix::fs::PermissionsExt;
160 let mut perm = meta.permissions();
161 perm.set_mode(0o700);
162 let _ = fs::set_permissions(parent, perm);
163 }
164 }
165 }
166 }
167 let mut out = File::create(&path)?;
168 out.write_all(value)?;
169 #[cfg(unix)]
170 {
171 use std::os::unix::fs::PermissionsExt;
172 let mut perm = out.metadata()?.permissions();
173 perm.set_mode(0o600);
174 let _ = fs::set_permissions(&path, perm);
175 }
176 Ok(())
177}
178
179pub struct StreamTrafficCapture {
180 upstream_sse: Vec<u8>,
181 upstream_events: Vec<Value>,
182 downstream_events: Vec<Value>,
183 malformed: Vec<Value>,
184 upstream_event_bytes: usize,
185 downstream_event_bytes: usize,
186 upstream_sse_truncated: u64,
187 upstream_events_truncated: u64,
188 downstream_events_truncated: u64,
189 malformed_truncated: u64,
190 upstream_frames_truncated: u64,
191}
192
193impl Default for StreamTrafficCapture {
194 fn default() -> Self {
195 Self {
196 upstream_sse: Vec::with_capacity(MAX_SSE_CAPTURE_BYTES.min(64 * 1024)),
197 upstream_events: Vec::new(),
198 downstream_events: Vec::new(),
199 malformed: Vec::new(),
200 upstream_event_bytes: 0,
201 downstream_event_bytes: 0,
202 upstream_sse_truncated: 0,
203 upstream_events_truncated: 0,
204 downstream_events_truncated: 0,
205 malformed_truncated: 0,
206 upstream_frames_truncated: 0,
207 }
208 }
209}
210
211impl StreamTrafficCapture {
212 pub fn upstream_event(&mut self, event: Option<&str>, value: &Value) {
213 let value = redact_traffic(value);
214 let frame = serde_json::to_vec(&value).unwrap_or_default();
215 let event = event.unwrap_or("message");
216 let frame_len = event.len().saturating_add(frame.len()).saturating_add(16);
217 if frame_len > MAX_STREAM_CAPTURE_FRAME_BYTES {
218 self.upstream_frames_truncated = self.upstream_frames_truncated.saturating_add(1);
219 } else if self.upstream_sse.len().saturating_add(frame_len) <= MAX_SSE_CAPTURE_BYTES {
220 self.upstream_sse.extend_from_slice(b"event: ");
221 self.upstream_sse.extend_from_slice(event.as_bytes());
222 self.upstream_sse.extend_from_slice(b"\ndata: ");
223 self.upstream_sse.extend_from_slice(&frame);
224 self.upstream_sse.extend_from_slice(b"\n\n");
225 } else {
226 self.upstream_sse_truncated = self.upstream_sse_truncated.saturating_add(1);
227 }
228 self.push_event(true, serde_json::json!({"event":event,"data":value}));
229 }
230
231 pub fn malformed(&mut self, stage: &str, kind: &str) {
232 if self.malformed.len() < MAX_STREAM_CAPTURE_EVENTS {
233 self.malformed
234 .push(serde_json::json!({"stage":stage,"kind":kind}));
235 } else {
236 self.malformed_truncated = self.malformed_truncated.saturating_add(1);
237 }
238 }
239
240 pub fn downstream_event(&mut self, event: &str, data: Value) {
241 self.push_event(
242 false,
243 serde_json::json!({"event":event,"data":redact_traffic(&data)}),
244 );
245 }
246
247 fn push_event(&mut self, upstream: bool, value: Value) {
248 let bytes = serde_json::to_vec(&value).map_or(0, |value| value.len());
249 let (events, total, truncated) = if upstream {
250 (
251 &mut self.upstream_events,
252 &mut self.upstream_event_bytes,
253 &mut self.upstream_events_truncated,
254 )
255 } else {
256 (
257 &mut self.downstream_events,
258 &mut self.downstream_event_bytes,
259 &mut self.downstream_events_truncated,
260 )
261 };
262 if events.len() < MAX_STREAM_CAPTURE_EVENTS
263 && total.saturating_add(bytes) <= MAX_STREAM_CAPTURE_EVENT_BYTES
264 {
265 *total += bytes;
266 events.push(value);
267 } else {
268 *truncated = truncated.saturating_add(1);
269 }
270 }
271
272 pub fn finish(self, traffic: &TrafficCapture, completion: Value) {
273 self.finish_named(traffic, completion, "061-grok-stream-summary");
274 }
275
276 pub fn finish_named(self, traffic: &TrafficCapture, completion: Value, summary_name: &str) {
277 let upstream_event_count = self.upstream_events.len();
278 let downstream_event_count = self.downstream_events.len();
279 if !self.upstream_sse.is_empty() {
280 traffic.write_bytes("032-upstream-response-body.sse", &self.upstream_sse);
281 }
282 traffic.write_json(
283 "033-upstream-response-capture",
284 &serde_json::json!({
285 "truncated": self.upstream_sse_truncated > 0 || self.upstream_frames_truncated > 0 || self.upstream_events_truncated > 0,
286 "captured_bytes": self.upstream_sse.len(),
287 "truncated_frames": self.upstream_sse_truncated,
288 "oversized_frames": self.upstream_frames_truncated,
289 "captured_events": self.upstream_events.len(),
290 "captured_event_bytes": self.upstream_event_bytes,
291 "truncated_events": self.upstream_events_truncated,
292 "malformed": self.malformed,
293 "truncated_malformed": self.malformed_truncated,
294 }),
295 );
296 for value in self.upstream_events {
297 traffic.write_json_event("040-upstream-event", &value);
298 }
299 for value in self.downstream_events {
300 traffic.write_json_event("050-downstream-event", &value);
301 }
302 traffic.write_json(
303 summary_name,
304 &serde_json::json!({
305 "completion": completion,
306 "upstream_sse": {
307 "captured_bytes": self.upstream_sse.len(),
308 "truncated_frames": self.upstream_sse_truncated,
309 "oversized_frames": self.upstream_frames_truncated,
310 },
311 "upstream_events": {
312 "captured": upstream_event_count,
313 "captured_bytes": self.upstream_event_bytes,
314 "truncated": self.upstream_events_truncated,
315 },
316 "downstream_events": {
317 "captured": downstream_event_count,
318 "captured_bytes": self.downstream_event_bytes,
319 "truncated": self.downstream_events_truncated,
320 },
321 }),
322 );
323 }
324}
325
326pub fn sanitize_path_part(input: &str) -> String {
327 let cleaned: String = input
328 .chars()
329 .map(|ch| {
330 if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' || ch == '.' {
331 ch
332 } else {
333 '_'
334 }
335 })
336 .collect();
337
338 let truncated = if cleaned.len() > 160 {
339 &cleaned[..160]
340 } else {
341 &cleaned
342 };
343 if truncated.is_empty() {
344 "unknown".to_string()
345 } else {
346 truncated.to_string()
347 }
348}
349
350pub fn redact_traffic(value: &Value) -> Value {
351 redact_traffic_with_depth(value, 0)
352}
353
354fn redact_traffic_with_depth(value: &Value, depth: u16) -> Value {
355 if depth > 100 {
356 return Value::String("[depth-limit]".to_string());
357 }
358
359 match value {
360 Value::Object(map) => {
361 let mut out = Map::new();
362 for (key, value) in map {
363 let normalized = key.to_lowercase();
364 if normalized == "image_url" {
365 out.insert(key.clone(), redact_traffic_value(value));
368 } else if normalized == "data" && looks_like_image_source(map) {
369 out.insert(key.clone(), redact_traffic_value(value));
372 } else if REDACT_KEYS.contains(&normalized.as_str())
373 || matches!(
374 normalized.as_str(),
375 "token"
376 | "bearer_token"
377 | "oauth_token"
378 | "oauth_access_token"
379 | "oauth_refresh_token"
380 | "client_secret"
381 | "secret"
382 | "password"
383 | "email"
384 | "user_id"
385 | "account_id"
386 | "identity"
387 | "identity_id"
388 | "subject"
389 | "sub"
390 )
391 {
392 out.insert(key.clone(), redact_traffic_value(value));
393 } else {
394 out.insert(key.clone(), redact_traffic_with_depth(value, depth + 1));
395 }
396 }
397 Value::Object(out)
398 }
399 Value::Array(values) => Value::Array(
400 values
401 .iter()
402 .map(|value| redact_traffic_with_depth(value, depth + 1))
403 .collect(),
404 ),
405 Value::String(text) if is_data_url(text) => {
406 Value::String(format!("[redacted data-url len={}]", text.len()))
407 }
408 _ => value.clone(),
409 }
410}
411
412fn is_data_url(text: &str) -> bool {
413 let lower = text.to_ascii_lowercase();
414 lower.starts_with("data:image/") && lower.contains(";base64,")
415}
416
417fn looks_like_image_source(map: &Map<String, Value>) -> bool {
423 if !map
424 .get("type")
425 .and_then(Value::as_str)
426 .is_some_and(|source_type| source_type.eq_ignore_ascii_case("base64"))
427 {
428 return false;
429 }
430 match map.get("media_type").and_then(Value::as_str) {
431 Some(media_type) => media_type.to_ascii_lowercase().starts_with("image/"),
432 None => true,
433 }
434}
435
436fn redact_traffic_value(value: &Value) -> Value {
437 match value {
438 Value::String(s) => Value::String(format!("[redacted len={}]", s.len())),
439 Value::Object(_) | Value::Array(_) => redact_traffic(value),
442 _ => Value::String("[redacted]".to_string()),
443 }
444}
445
446fn set_mode(path: &Path, mode: u32) {
447 #[cfg(unix)]
448 {
449 use std::os::unix::fs::PermissionsExt;
450 if let Ok(meta) = fs::metadata(path) {
451 let mut perm = meta.permissions();
452 perm.set_mode(mode);
453 let _ = fs::set_permissions(path, perm);
454 }
455 }
456}
457
458#[allow(dead_code)]
459fn _provider_alias(_provider: &str) -> Option<AliasProvider> {
460 None
461}
462
463#[cfg(test)]
464mod tests {
465 use super::*;
466
467 #[test]
468 fn redact_traffic_strips_input_image_data_urls() {
469 let value = serde_json::json!({
470 "input": [
471 {"type": "message", "role": "user", "content": [
472 {"type": "input_text", "text": "what color?"},
473 {"type": "input_image", "image_url": "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nGP4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC"}
474 ]}
475 ]
476 });
477 let redacted = redact_traffic(&value);
478 let rendered = redacted.to_string();
479 assert!(
480 !rendered.contains("iVBORw0KGgo"),
481 "base64 payload leaked: {rendered}"
482 );
483 assert!(
484 !rendered.contains("data:image/png;base64,iVBOR"),
485 "data URL payload leaked: {rendered}"
486 );
487 assert!(rendered.contains("input_image"));
489 }
490
491 #[test]
492 fn redact_traffic_strips_inline_base64_image_values() {
493 let value = serde_json::json!({
494 "output": "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nGP4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC"
495 });
496 let redacted = redact_traffic(&value);
497 let rendered = redacted.to_string();
498 assert!(!rendered.contains("iVBORw0KGgo"));
499 }
500
501 #[test]
502 fn redact_traffic_strips_anthropic_image_source_data() {
503 let value = serde_json::json!({
506 "messages": [{
507 "role": "user",
508 "content": [
509 {"type": "text", "text": "look"},
510 {"type": "image", "source": {
511 "type": "base64",
512 "media_type": "image/png",
513 "data": "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nGP4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC"
514 }}
515 ]
516 }]
517 });
518 let redacted = redact_traffic(&value);
519 let rendered = redacted.to_string();
520 assert!(
521 !rendered.contains("iVBORw0KGgo"),
522 "source.data base64 leaked: {rendered}"
523 );
524 assert!(rendered.contains("redacted"));
525 assert!(rendered.contains("image/png"));
527 }
528
529 #[test]
530 fn redact_traffic_keeps_unrelated_data_keys() {
531 let value = serde_json::json!({"data": "some-non-image-payload", "count": 3});
532 let redacted = redact_traffic(&value);
533 assert_eq!(redacted["data"], "some-non-image-payload");
534 }
535
536 #[test]
537 fn redact_traffic_strips_data_urls_in_array_shaped_tool_output() {
538 let value = serde_json::json!({
541 "input": [
542 {"type": "function_call_output", "call_id": "call_1", "output": [
543 {"type": "input_text", "text": "screenshot"},
544 {"type": "input_image", "image_url": "data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nGP4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC"}
545 ]}
546 ]
547 });
548 let redacted = redact_traffic(&value);
549 let rendered = redacted.to_string();
550 assert!(
551 !rendered.contains("iVBORw0KGgo"),
552 "base64 payload leaked: {rendered}"
553 );
554 assert!(
555 !rendered.contains("data:image/png;base64,iVBOR"),
556 "data URL payload leaked: {rendered}"
557 );
558 assert!(rendered.contains("input_image"));
560 assert!(rendered.contains("screenshot"));
561 assert!(rendered.contains("redacted"));
562 }
563
564 #[test]
565 fn redact_traffic_strips_uppercase_media_type_image_source() {
566 let value = serde_json::json!({
569 "source": {"type": "base64", "media_type": "IMAGE/PNG", "data": "QUJDREVGRw=="}
570 });
571 let redacted = redact_traffic(&value);
572 let rendered = redacted.to_string();
573 assert!(
574 !rendered.contains("QUJDREVGRw"),
575 "payload leaked: {rendered}"
576 );
577 assert!(
578 redacted["source"]["data"]
579 .as_str()
580 .unwrap()
581 .contains("redacted")
582 );
583 }
584
585 #[test]
586 fn redact_traffic_strips_uppercase_image_source_type() {
587 let value = serde_json::json!({
588 "source": {"type": "BASE64", "media_type": "image/png", "data": "QUJDREVGRw=="}
589 });
590 let redacted = redact_traffic(&value);
591 let rendered = redacted.to_string();
592 assert!(
593 !rendered.contains("QUJDREVGRw"),
594 "payload leaked: {rendered}"
595 );
596 }
597
598 #[test]
599 fn redact_traffic_strips_image_source_without_media_type() {
600 let value = serde_json::json!({
603 "source": {"type": "base64", "data": "QUJDREVGRw=="}
604 });
605 let redacted = redact_traffic(&value);
606 let rendered = redacted.to_string();
607 assert!(
608 !rendered.contains("QUJDREVGRw"),
609 "payload leaked: {rendered}"
610 );
611 }
612
613 #[test]
614 fn redact_traffic_strips_uppercase_data_url_scheme() {
615 let value = serde_json::json!({
616 "image_url": "DATA:IMAGE/PNG;base64,iVBORw0KGgoAAAANSUhEUg"
617 });
618 let redacted = redact_traffic(&value);
619 let rendered = redacted.to_string();
620 assert!(
621 !rendered.contains("iVBORw0KGgo"),
622 "data URL payload leaked: {rendered}"
623 );
624 }
625
626 #[test]
627 fn redact_traffic_keeps_structured_image_url_shape() {
628 let value = serde_json::json!({
631 "image_url": {"url": "DATA:IMAGE/PNG;BASE64,iVBORw0KGgoAAAANSUhEUg"}
632 });
633 let redacted = redact_traffic(&value);
634 assert!(redacted["image_url"].is_object(), "shape lost: {redacted}");
635 let rendered = redacted.to_string();
636 assert!(
637 !rendered.contains("iVBORw0KGgo"),
638 "nested payload leaked: {rendered}"
639 );
640 }
641
642 #[test]
643 fn traffic_redacts_proxy_authorization() {
644 let redacted = redact_traffic(&serde_json::json!({
645 "headers": {
646 "proxy-authorization": "Basic dXNlcjpwYXNz",
647 "x-safe": "kept"
648 }
649 }));
650
651 assert_eq!(redacted["headers"]["x-safe"], "kept");
652 assert_eq!(
653 redacted["headers"]["proxy-authorization"],
654 "[redacted len=18]"
655 );
656 }
657}