Skip to main content

claude_codex/
traffic.rs

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                    // Grok/OpenAI `input_image` parts carry data URLs; never
366                    // persist the payload in traffic captures.
367                    out.insert(key.clone(), redact_traffic_value(value));
368                } else if normalized == "data" && looks_like_image_source(map) {
369                    // Anthropic image blocks carry raw base64 under
370                    // `source.data`; redact it for the same reason.
371                    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
417/// An object shaped like an Anthropic image source: `type: "base64"` and
418/// either an image `media_type` or none at all. Used to redact only image
419/// payloads, not arbitrary `data` keys. Missing `media_type` is treated as
420/// an image: the capture is written before translation, so a malformed block
421/// that the translator would reject must still not persist its payload.
422fn 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        // Structured values (e.g. Kimi's `{url: ...}` image_url) keep their
440        // shape so captures stay parseable; only the payload leaves.
441        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        // The image part is still structurally visible (redacted marker).
488        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        // The incoming Anthropic request body carries raw base64 under
504        // `source.data` (no data: URL wrapper) — it must not persist either.
505        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        // Structure is preserved.
526        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        // L2b inline mode: function_call_output.output is an array of content
539        // parts — the image_url inside must be redacted just like message parts.
540        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        // Structure and text survive.
559        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        // Captures are written before translation validates anything, so a
567        // non-standard casing must still not persist its payload.
568        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        // A malformed block the translator would reject must still be redacted
601        // in the pre-translation capture.
602        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        // Kimi-style object image_url: redact the payload but keep the object
629        // structure so captures stay parseable.
630        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}