use axum::{
body::Body,
extract::State,
http::{Request, StatusCode},
response::Response,
};
use serde_json::Value;
use super::ProxyState;
use super::compress_shared::{self, ToolKind};
use super::forward;
use super::tool_kind::{self, ToolResultKind};
use super::{cache_safety, prose};
use crate::core::config::{HistoryMode, ProseRole};
pub async fn handler(
State(state): State<ProxyState>,
req: Request<Body>,
) -> Result<Response, StatusCode> {
let upstream = state.openai_upstream();
forward::forward_request(
State(state),
req,
&upstream,
"/v1/responses",
compress_request_body,
"OpenAI",
&[],
)
.await
}
pub async fn ws_handler(
State(state): State<ProxyState>,
headers: axum::http::HeaderMap,
ws: axum::extract::ws::WebSocketUpgrade,
) -> Response {
super::openai_responses_ws::upgrade(state, ws, &headers)
}
pub(super) fn compress_request_body(
parsed: Value,
original_size: usize,
) -> (Vec<u8>, usize, usize) {
let mut doc = parsed;
let cfg = crate::core::config::Config::load();
let system_aggr = cfg.proxy.resolved_role_aggressiveness(ProseRole::System);
let user_aggr = cfg.proxy.resolved_role_aggressiveness(ProseRole::User);
let live_compress = cfg.proxy.live_compresses();
let mode = cfg.proxy.resolved_history_mode();
let mut modified = false;
let arm = super::holdout::assign(
&super::holdout::openai_responses_key(&doc),
cfg.proxy.output_holdout_fraction(),
);
if cfg.proxy.ccr_inband_enabled() {
modified |= super::ccr::splice_inband_in_place(&mut doc);
}
if arm == super::holdout::Arm::Treatment {
if let Some(effort) = cfg.proxy.resolved_effort() {
modified |= super::effort::apply_openai_responses(&mut doc, effort);
}
if cfg.proxy.verbosity_steer_enabled() {
modified |= super::verbosity::apply_openai_responses(&mut doc);
}
}
if !live_compress
&& mode == HistoryMode::Off
&& system_aggr.is_none()
&& user_aggr.is_none()
&& !modified
{
let out = serde_json::to_vec(&doc).unwrap_or_default();
return (out, original_size, original_size);
}
let mut prose_segments: u64 = 0;
if let Some(a) = system_aggr {
prose_segments += u64::from(prose::compress_string_field(&mut doc, "instructions", a));
}
modified |= prune_responses_input(&mut doc);
modified |= compress_responses_input(&mut doc);
if let Some(a) = user_aggr {
prose_segments += u64::from(compress_responses_user_prose(&mut doc, mode, a));
}
if prose_segments > 0 {
modified = true;
}
cache_safety::record(prose_segments, true);
let out = serde_json::to_vec(&doc).unwrap_or_default();
let compressed_size = if modified { out.len() } else { original_size };
(out, original_size, compressed_size)
}
fn compress_responses_user_prose(doc: &mut Value, mode: HistoryMode, aggressiveness: f64) -> u32 {
let Some(input) = doc.get_mut("input").and_then(|i| i.as_array_mut()) else {
return 0;
};
let boundary = super::history_prune::prune_boundary(mode, input.len());
if boundary == 0 {
return 0;
}
let mut segments = 0;
for item in input.iter_mut().take(boundary) {
let item_type = item
.get("type")
.and_then(|t| t.as_str())
.unwrap_or("message");
let role = item.get("role").and_then(|r| r.as_str());
if item_type == "message" && role == Some("user") {
segments += prose::compress_message_content(item, aggressiveness);
}
}
segments
}
pub(super) fn prune_responses_input(doc: &mut Value) -> bool {
let mode = crate::core::config::Config::load()
.proxy
.resolved_history_mode();
let Some(input) = doc.get_mut("input").and_then(|i| i.as_array_mut()) else {
return false;
};
let boundary = super::history_prune::prune_boundary(mode, input.len());
if boundary == 0 {
return false;
}
let tool_names = tool_kind::responses_tool_names(input);
let mut modified = false;
for item in input.iter_mut().take(boundary) {
if compress_shared::classify_tool_kind(item) != ToolKind::ToolResult {
continue;
}
let tool_name = item
.get("call_id")
.and_then(|v| v.as_str())
.and_then(|id| tool_names.get(id))
.map(String::as_str);
let kind = compress_shared::tool_result_kind(tool_name);
if let Some(output) = item.get_mut("output") {
modified |= prune_output_field(output, kind);
}
}
modified
}
fn prune_output_field(output: &mut Value, kind: ToolResultKind) -> bool {
super::tool_output::prune_value(output, kind)
}
pub(super) fn compress_responses_input(doc: &mut Value) -> bool {
let cfg = crate::core::config::Config::load();
if !cfg.proxy.live_compresses() {
return false;
}
let mut modified = false;
if let Some(input) = doc.get_mut("input").and_then(|i| i.as_array_mut()) {
let tool_names = tool_kind::responses_tool_names(input);
for item in input.iter_mut() {
if compress_shared::classify_tool_kind(item) != ToolKind::ToolResult {
continue;
}
let name = item
.get("call_id")
.and_then(|v| v.as_str())
.and_then(|id| tool_names.get(id))
.map(String::as_str);
if name.is_some_and(|n| cfg.proxy.is_tool_live_compress_excluded(n)) {
continue;
}
let kind = compress_shared::tool_result_kind(name);
if let Some(output) = item.get_mut("output") {
modified |= compress_output_field(output, name, kind);
}
}
}
modified
}
fn compress_output_field(
output: &mut Value,
tool_name: Option<&str>,
kind: ToolResultKind,
) -> bool {
super::tool_output::compress_value(output, tool_name, kind)
}
#[cfg(test)]
mod tests {
use super::super::compress::compress_tool_result;
use super::*;
fn long_git_status() -> String {
let mut s = String::from(
"$ git status\nOn branch main\nYour branch is up to date with 'origin/main'.\n\nChanges not staged for commit:\n (use \"git add <file>...\" to update what will be committed)\n",
);
for i in 0..80 {
s.push_str(&format!("\tmodified: src/module_{i}/file_{i}.rs\n"));
}
s.push_str("\nno changes added to commit (use \"git add\" and/or \"git commit -a\")\n");
s
}
fn big_prose() -> String {
let p = "You are a careful, senior software engineer. You always explain your \
reasoning before making changes, you prefer small reviewable diffs, and \
you never introduce mock data or placeholders into production code. ";
[p; 6].join("\n")
}
fn long_tool_json() -> String {
let rows = (0..32)
.map(|i| {
serde_json::json!({
"path": format!("/Users/alex/work/app/src/module_{i}.rs"),
"regex": r"src/[a-z_]+\.rs:\d+",
"error": format!("error[E0{i:03}]: expected exact diagnostic text"),
})
})
.collect::<Vec<_>>();
serde_json::to_string(&serde_json::json!({ "results": rows })).unwrap()
}
#[test]
fn shell_json_envelope_text_is_compressed() {
let _lock = crate::core::data_dir::test_env_lock();
let raw = long_git_status();
let expected = compress_tool_result(&raw, Some("Bash"));
let envelope = serde_json::to_string(&serde_json::json!({
"content": [{"type": "text", "text": raw}],
"isError": false,
}))
.unwrap();
let body = serde_json::json!({
"model": "gpt-5",
"input": [
{"type": "function_call", "call_id": "call_1", "name": "Bash", "arguments": "{}"},
{"type": "function_call_output", "call_id": "call_1", "output": envelope}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig);
let parsed: Value = serde_json::from_slice(&out).unwrap();
let output = parsed["input"][1]["output"].as_str().unwrap();
let envelope: Value = serde_json::from_str(output).unwrap();
assert_eq!(envelope["content"][0]["text"].as_str().unwrap(), expected);
}
#[test]
fn shell_json_envelope_non_text_field_is_compressed() {
let _lock = crate::core::data_dir::test_env_lock();
let raw = long_git_status();
let expected = compress_tool_result(&raw, Some("Bash"));
let envelope = serde_json::to_string(&serde_json::json!({
"stdout": raw,
"exit_code": 0,
}))
.unwrap();
let body = serde_json::json!({
"model": "gpt-5",
"input": [
{"type": "function_call", "call_id": "call_1", "name": "Bash", "arguments": "{}"},
{"type": "function_call_output", "call_id": "call_1", "output": envelope}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "non-text shell field should be compressed");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let output = parsed["input"][1]["output"].as_str().unwrap();
let envelope: Value = serde_json::from_str(output).unwrap();
assert_eq!(envelope["stdout"].as_str().unwrap(), expected);
assert_eq!(envelope["exit_code"].as_i64().unwrap(), 0);
}
#[test]
fn old_shell_json_envelope_text_is_pruned() {
let _iso = crate::core::data_dir::isolated_data_dir();
let raw = long_git_status();
let envelope = serde_json::to_string(&serde_json::json!({
"content": [{"type": "text", "text": raw}],
"isError": false,
}))
.unwrap();
let mut output = Value::String(envelope);
assert!(prune_output_field(&mut output, ToolResultKind::Shell));
let envelope: Value = serde_json::from_str(output.as_str().unwrap()).unwrap();
assert!(
envelope["content"][0]["text"].as_str().unwrap().len() < raw.len(),
"nested text payload should be pruned"
);
}
#[test]
fn string_output_mirrors_engine_and_shrinks() {
let _lock = crate::core::data_dir::test_env_lock();
let raw = long_git_status();
let expected = compress_tool_result(&raw, None);
assert!(
expected.len() < raw.len(),
"fixture must be compressible by the shared engine"
);
let body = serde_json::json!({
"model": "gpt-5",
"input": [
{"type": "function_call_output", "call_id": "call_1", "output": raw}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "compressed body must be smaller");
let parsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(
parsed["input"][0]["output"].as_str().unwrap(),
expected,
"output must be exactly what the shared compressor produces"
);
}
#[test]
fn array_output_text_is_compressed() {
let _lock = crate::core::data_dir::test_env_lock();
let raw = long_git_status();
let expected = compress_tool_result(&raw, None);
let body = serde_json::json!({
"input": [
{
"type": "function_call_output",
"call_id": "call_1",
"output": [{"type": "input_text", "text": raw}]
}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig);
let parsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(
parsed["input"][0]["output"][0]["text"].as_str().unwrap(),
expected
);
}
#[test]
fn non_tool_output_items_are_untouched() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| c.proxy.proxy_mode = Some("cache".into()))
.unwrap();
let body = serde_json::json!({
"input": [
{"type": "message", "role": "user", "content": long_git_status()},
{"type": "function_call", "call_id": "c", "name": "x", "arguments": "{}"}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body.clone(), bytes.len());
assert_eq!(comp, orig, "no function_call_output → passthrough");
let reparsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(reparsed, body);
}
#[test]
fn plain_string_input_passthrough() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| c.proxy.proxy_mode = Some("cache".into()))
.unwrap();
let body = serde_json::json!({"model": "gpt-5", "input": "hello world"});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body.clone(), bytes.len());
assert_eq!(comp, orig);
let reparsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(reparsed, body);
}
#[test]
fn no_input_field_passthrough() {
let _iso = crate::core::data_dir::isolated_data_dir();
let body = serde_json::json!({"model": "gpt-5", "previous_response_id": "resp_abc"});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body.clone(), bytes.len());
assert_eq!(comp, orig);
let reparsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(reparsed, body);
}
#[test]
fn chatgpt_responses_eval_fixture_keeps_exact_payloads_and_pairing() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.user = Some(0.8);
c.proxy.proxy_mode = Some("cache".into());
})
.unwrap();
let command_input = "$ cargo test --lib proxy::openai_responses\nerror[E0425]: cannot find value `x` in this scope\nsrc/proxy/openai_responses.rs:12:9";
let body = serde_json::json!({"model": "gpt-5", "input": command_input});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body.clone(), bytes.len());
assert_eq!(comp, orig, "top-level exact command string must stay raw");
assert_eq!(serde_json::from_slice::<Value>(&out).unwrap(), body);
let old_json = long_tool_json();
let recent_json = long_tool_json();
let old_shell = long_git_status();
let input_text_block = "```rust\nfn main() { panic!(\"exact\"); }\n```\nRegex: src/[a-z_]+\\.rs:\\d+\nPath: /Users/alex/work/app/src/main.rs\nError: error[E0425]: cannot find value `x` in this scope";
let mut input = vec![
serde_json::json!({"type": "reasoning", "id": "rs_1", "summary": []}),
serde_json::json!({"type": "function_call", "call_id": "json_old", "name": "submit_tool_json", "arguments": "{\"strict\":true}"}),
serde_json::json!({"type": "function_call_output", "call_id": "json_old", "output": old_json}),
serde_json::json!({"type": "function_call", "call_id": "shell_old", "name": "Bash", "arguments": "{\"cmd\":\"git status\"}"}),
serde_json::json!({"type": "function_call_output", "call_id": "shell_old", "output": old_shell}),
serde_json::json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": input_text_block}]}),
];
while input.len() < 22 {
input.push(serde_json::json!({
"type": "message",
"role": "user",
"content": format!("filler {}", input.len()),
}));
}
input.push(serde_json::json!({"type": "function_call", "call_id": "json_recent", "name": "submit_tool_json", "arguments": "{\"strict\":true}"}));
input.push(serde_json::json!({"type": "function_call_output", "call_id": "json_recent", "output": recent_json}));
let body = serde_json::json!({"model": "gpt-5", "input": input});
let item_count = body["input"].as_array().unwrap().len();
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "old shell output should still provide savings");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let input = parsed["input"].as_array().unwrap();
assert_eq!(input.len(), item_count, "no Responses item may be dropped");
assert_eq!(input[0]["type"], "reasoning");
assert_eq!(input[1]["type"], "function_call");
assert_eq!(input[2]["type"], "function_call_output");
assert_eq!(input[2]["output"].as_str().unwrap(), old_json);
assert_ne!(input[4]["output"].as_str().unwrap(), old_shell);
assert_eq!(
input[5]["content"][0]["text"].as_str().unwrap(),
input_text_block
);
assert_eq!(input[22]["type"], "function_call");
assert_eq!(input[23]["type"], "function_call_output");
assert_eq!(input[23]["output"].as_str().unwrap(), recent_json);
assert_eq!(input[1]["call_id"], input[2]["call_id"]);
assert_eq!(input[22]["call_id"], input[23]["call_id"]);
}
#[test]
fn responses_instructions_prose_compressed_and_assistant_untouched() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.6);
})
.unwrap();
let prose = big_prose();
let body = serde_json::json!({
"model": "gpt-5",
"instructions": prose,
"input": [
{"type": "message", "role": "user", "content": "hi"},
{"type": "message", "role": "assistant", "content": prose},
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "enabled instructions prose must save bytes");
let parsed: Value = serde_json::from_slice(&out).unwrap();
assert!(
parsed["instructions"].as_str().unwrap().len() < prose.len(),
"Responses instructions must be compressed when enabled"
);
assert_eq!(
parsed["input"][1]["content"].as_str().unwrap(),
prose,
"assistant turns must pass through verbatim (#710)"
);
}
#[test]
fn responses_user_prose_compressed_only_in_frozen_region() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.user = Some(0.7);
c.proxy.proxy_mode = Some("cache".into());
})
.unwrap();
let prose = big_prose();
let mut input = Vec::new();
for i in 0..30 {
let role = if i % 2 == 0 { "user" } else { "assistant" };
input.push(serde_json::json!({
"type": "message",
"role": role,
"content": prose,
}));
}
let body = serde_json::json!({"model": "gpt-5", "input": input});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "old user prose must save bytes");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let frozen_user = parsed["input"][0]["content"].as_str().unwrap();
assert!(
frozen_user.len() < prose.len(),
"old user prose should compress"
);
assert_eq!(
parsed["input"][1]["content"].as_str().unwrap(),
prose,
"assistant prose must stay verbatim"
);
assert_eq!(
parsed["input"][16]["content"].as_str().unwrap(),
prose,
"live-tail user prose must stay verbatim"
);
}
#[test]
fn responses_prose_compression_is_deterministic() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.6);
})
.unwrap();
let prose = big_prose();
let mk = || {
serde_json::json!({
"model": "gpt-5",
"instructions": prose,
"input": [{"type": "message", "role": "user", "content": "hi"}],
})
};
let (a, b) = (mk(), mk());
let (la, lb) = (
serde_json::to_vec(&a).unwrap().len(),
serde_json::to_vec(&b).unwrap().len(),
);
assert_eq!(
compress_request_body(a, la).0,
compress_request_body(b, lb).0,
"identical input must yield byte-identical output (#498)"
);
}
#[test]
fn short_output_unchanged() {
let _iso = crate::core::data_dir::isolated_data_dir();
let body = serde_json::json!({
"input": [
{"type": "function_call_output", "call_id": "c", "output": "ok"}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body.clone(), bytes.len());
assert_eq!(comp, orig);
let reparsed: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(reparsed, body);
}
fn responses_read_turns(pairs: usize) -> Vec<Value> {
let code = (0..40)
.map(|i| format!(" let v{i} = compute_{i}(ctx, opts);"))
.collect::<Vec<_>>()
.join("\n");
let mut input = Vec::new();
for t in 0..pairs {
input.push(serde_json::json!({
"type": "function_call", "call_id": format!("c{t}"),
"name": "read_file", "arguments": "{}"
}));
input.push(serde_json::json!({
"type": "function_call_output", "call_id": format!("c{t}"),
"output": format!("{code}\n// turn {t}")
}));
}
input
}
#[test]
fn cache_aware_prune_stubs_old_reads_keeps_recent_and_pairing() {
let _iso = crate::core::data_dir::isolated_data_dir();
let body = serde_json::json!({"model": "gpt-5", "input": responses_read_turns(14)});
let item_count = body["input"].as_array().unwrap().len();
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "old reads must be pruned for savings");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let input = parsed["input"].as_array().unwrap();
assert_eq!(input.len(), item_count, "no items may be removed (pairing)");
for (i, item) in input.iter().enumerate() {
let expect = if i.is_multiple_of(2) {
"function_call"
} else {
"function_call_output"
};
assert_eq!(item["type"], expect, "item {i} type/order changed");
}
let old = input[1]["output"].as_str().unwrap();
assert!(
old.contains("Re-read the file"),
"old read should be stubbed, got: {old}"
);
let recent = input[27]["output"].as_str().unwrap();
assert!(
recent.contains("v39"),
"recent read must be protected, got: {recent}"
);
}
#[test]
fn responses_compression_is_deterministic() {
let _iso = crate::core::data_dir::isolated_data_dir();
let mk = || serde_json::json!({"model": "gpt-5", "input": responses_read_turns(14)});
let (a, b) = (mk(), mk());
let (la, lb) = (
serde_json::to_vec(&a).unwrap().len(),
serde_json::to_vec(&b).unwrap().len(),
);
let (out_a, _, _) = compress_request_body(a, la);
let (out_b, _, _) = compress_request_body(b, lb);
assert_eq!(out_a, out_b, "identical input must yield identical bytes");
}
#[test]
fn cache_aware_responses_prefix_is_byte_stable_across_turns() {
let _iso = crate::core::data_dir::isolated_data_dir();
let mut prev: Vec<String> = Vec::new();
let mut prev_boundary = 0;
for pairs in 1..=20 {
let input = responses_read_turns(pairs);
let len = input.len();
let body = serde_json::json!({"model": "gpt-5", "input": input});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _, _) = compress_request_body(body, bytes.len());
let parsed: Value = serde_json::from_slice(&out).unwrap();
let items: Vec<String> = parsed["input"]
.as_array()
.unwrap()
.iter()
.map(Value::to_string)
.collect();
for i in 0..prev_boundary {
assert_eq!(
prev[i], items[i],
"Responses item {i} changed at turn {pairs} — prompt cache prefix broken"
);
}
prev = items;
prev_boundary = crate::proxy::history_prune::prune_boundary(
crate::core::config::HistoryMode::CacheAware,
len,
);
}
}
#[test]
fn effort_control_sets_nested_reasoning_effort() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_EFFORT");
crate::core::config::Config::update_global(|c| {
c.proxy.effort = Some("low".into());
})
.unwrap();
let body = serde_json::json!({"model": "gpt-5.5", "input": []});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
serde_json::from_slice::<Value>(&out).unwrap()["reasoning"]["effort"],
"low"
);
}
}