use axum::{
body::Body,
extract::State,
http::{Request, StatusCode},
response::Response,
};
use serde_json::Value;
use super::ProxyState;
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/chat/completions",
compress_request_body,
"OpenAI",
&[],
)
.await
}
fn compress_request_body(parsed: Value, original_size: usize) -> (Vec<u8>, usize, usize) {
let mut doc = parsed;
let mut modified = false;
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 arm = super::holdout::assign(
&super::holdout::openai_chat_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_chat(&mut doc, effort);
}
if cfg.proxy.verbosity_steer_enabled() {
modified |= super::verbosity::apply_openai_chat(&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(messages) = doc.get_mut("messages").and_then(|m| m.as_array_mut()) {
let tool_names = tool_kind::openai_tool_names(messages);
let boundary = super::history_prune::prune_boundary(mode, messages.len());
let cached = super::history_prune::cached_prefix_len(messages);
modified |=
super::history_prune::prune_history_range(messages, cached, boundary, &tool_names);
for msg in messages.iter_mut() {
let role = msg.get("role").and_then(|r| r.as_str()).unwrap_or("");
if role != "tool" {
continue;
}
let name = msg
.get("tool_call_id")
.and_then(|v| v.as_str())
.and_then(|id| tool_names.get(id))
.map(String::as_str);
let kind = name.map_or(ToolResultKind::Other, tool_kind::classify_tool_name);
if !live_compress || name.is_some_and(|n| cfg.proxy.is_tool_live_compress_excluded(n)) {
continue;
}
if let Some(mut content) = msg
.get_mut("content")
.and_then(|c| c.as_str().map(String::from))
&& super::tool_output::compress_text(&mut content, name, kind)
{
msg["content"] = Value::String(content);
modified = true;
}
}
if system_aggr.is_some() || user_aggr.is_some() {
for (i, msg) in messages.iter_mut().enumerate() {
if i < cached {
continue;
}
let role = msg.get("role").and_then(|r| r.as_str()).unwrap_or("");
let aggr = match role {
"system" | "developer" => system_aggr,
"user" if i < boundary => user_aggr,
_ => None,
};
if let Some(a) = aggr {
prose_segments += u64::from(prose::compress_message_content(msg, a));
}
}
}
}
if prose_segments > 0 {
modified = true;
}
cache_safety::record(prose_segments, true);
maybe_inject_usage_reporting(&mut doc);
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 maybe_inject_usage_reporting(doc: &mut Value) {
if crate::core::config::Config::load()
.proxy
.meters_openai_usage()
{
inject_usage_reporting(doc);
}
}
fn inject_usage_reporting(doc: &mut Value) {
if doc.get("stream").and_then(Value::as_bool) != Some(true) {
return;
}
let Some(obj) = doc.as_object_mut() else {
return;
};
let opts = obj
.entry("stream_options")
.or_insert_with(|| Value::Object(serde_json::Map::new()));
if let Some(opts_obj) = opts.as_object_mut() {
opts_obj.entry("include_usage").or_insert(Value::Bool(true));
}
}
#[cfg(test)]
mod tests {
use super::super::compress::compress_tool_result;
use super::*;
#[test]
fn read_file_tool_result_protected() {
let code = (0..60)
.map(|i| format!(" const value{i} = computeValue{i}(ctx, opts);"))
.collect::<Vec<_>>()
.join("\n");
let body = serde_json::json!({
"model": "gpt-5",
"messages": [
{"role": "assistant", "tool_calls": [{"id": "call_1", "type": "function", "function": {"name": "read_file"}}]},
{"role": "tool", "tool_call_id": "call_1", "content": code}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _orig, _comp) = compress_request_body(body, bytes.len());
let parsed: Value = serde_json::from_slice(&out).unwrap();
assert!(
parsed["messages"][1]["content"]
.as_str()
.unwrap()
.contains("value59")
);
}
#[test]
fn json_envelope_tool_output_is_compressed() {
let _iso = crate::core::data_dir::isolated_data_dir();
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",
"messages": [
{"role": "assistant", "tool_calls": [{
"id": "call_1",
"type": "function",
"function": {"name": "Bash"}
}]},
{"role": "tool", "tool_call_id": "call_1", "content": envelope}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "JSON envelope tool output should shrink");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let content = parsed["messages"][1]["content"].as_str().unwrap();
let envelope: Value = serde_json::from_str(content).unwrap();
assert_eq!(envelope["content"][0]["text"].as_str().unwrap(), expected);
}
#[test]
fn injects_include_usage_for_streaming() {
let mut doc = serde_json::json!({"model": "gpt-5.4", "stream": true, "messages": []});
inject_usage_reporting(&mut doc);
assert_eq!(doc["stream_options"]["include_usage"], Value::Bool(true));
}
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
}
#[test]
fn no_injection_for_non_streaming() {
let mut doc = serde_json::json!({"model": "gpt-5.4", "messages": []});
inject_usage_reporting(&mut doc);
assert!(
doc.get("stream_options").is_none(),
"non-streamed requests get usage in the body, no injection needed"
);
}
#[test]
fn respects_client_set_include_usage() {
let mut doc = serde_json::json!({
"model": "gpt-5.4",
"stream": true,
"stream_options": {"include_usage": false},
"messages": []
});
inject_usage_reporting(&mut doc);
assert_eq!(
doc["stream_options"]["include_usage"],
Value::Bool(false),
"an explicit client value must be preserved"
);
}
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")
}
#[test]
fn system_message_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",
"messages": [
{"role": "system", "content": prose},
{"role": "user", "content": "hi"},
{"role": "assistant", "content": prose},
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let parsed: Value = serde_json::from_slice(&out).unwrap();
assert!(
parsed["messages"][0]["content"].as_str().unwrap().len() < prose.len(),
"system message prose must be compressed when enabled"
);
assert_eq!(
parsed["messages"][2]["content"].as_str().unwrap(),
prose,
"assistant turns must pass through verbatim (#710)"
);
}
#[test]
fn openai_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",
"messages": [{"role": "system", "content": prose}, {"role": "user", "content": "hi"}]
})
};
let (a, b) = (mk(), mk());
let la = serde_json::to_vec(&a).unwrap().len();
let lb = 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 effort_control_sets_reasoning_effort_and_off_is_noop() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_EFFORT");
let body = serde_json::json!({
"model": "gpt-5.5", "messages": [{"role": "user", "content": "hi"}]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (off, o, c) = compress_request_body(body.clone(), bytes.len());
assert_eq!(c, o, "effort off must be a passthrough");
assert!(
serde_json::from_slice::<Value>(&off)
.unwrap()
.get("reasoning_effort")
.is_none()
);
crate::core::config::Config::update_global(|cfg| {
cfg.proxy.effort = Some("low".into());
})
.unwrap();
let (on, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
serde_json::from_slice::<Value>(&on).unwrap()["reasoning_effort"],
"low"
);
}
}