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.anthropic_upstream();
forward::forward_request(
State(state),
req,
&upstream,
"/v1/messages",
compress_request_body,
"Anthropic",
&[],
)
.await
}
pub(super) 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 inject_breakpoint = cfg.proxy.cache_breakpoint_enabled();
let align_volatile = cfg.proxy.cache_aligner_enabled();
let relocate_volatile = cfg.proxy.cache_align_relocate_enabled();
let cache_economics = cfg.proxy.cache_policy_enabled();
let arm = super::holdout::assign(
&super::holdout::anthropic_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_anthropic(&mut doc, effort);
}
if cfg.proxy.verbosity_steer_enabled() {
modified |= super::verbosity::apply_anthropic(&mut doc);
}
}
if !live_compress
&& mode == HistoryMode::Off
&& system_aggr.is_none()
&& user_aggr.is_none()
&& !modified
&& !inject_breakpoint
&& !align_volatile
&& !relocate_volatile
&& !cache_economics
{
let out = serde_json::to_vec(&doc).unwrap_or_default();
return (out, original_size, original_size);
}
let mut prose_segments: u64 = 0;
let cached = doc
.get("messages")
.and_then(|m| m.as_array())
.map_or(0, |m| super::history_prune::cached_prefix_len(m));
if cache_economics && let Some(m) = doc.get("messages").and_then(|m| m.as_array()) {
super::cache_attribution::record_request(m, cached);
}
let repack = cfg.proxy.repacks_cold_prefix()
&& doc
.get("messages")
.and_then(|m| m.as_array())
.is_some_and(|m| {
super::cold_prefix::repack_decision(m, cached)
&& (!cache_economics
|| super::cache_policy::worth_repacking(doc.get("system"), m, cached))
});
let protect = if repack { 0 } else { cached };
if let Some(a) = system_aggr
&& protect == 0
&& let Some(system) = doc.get_mut("system")
&& (repack || !prose::value_has_cache_control(system))
{
let n = prose::compress_system_value(system, a);
if n > 0 {
prose_segments += u64::from(n);
modified = true;
}
}
if let Some(messages) = doc.get_mut("messages").and_then(|m| m.as_array_mut()) {
let tool_names = tool_kind::anthropic_tool_names(messages);
let boundary = super::history_prune::prune_boundary(mode, messages.len());
modified |=
super::history_prune::prune_history_range(messages, protect, boundary, &tool_names);
for msg in messages.iter_mut() {
let role = msg.get("role").and_then(|r| r.as_str()).unwrap_or("");
if role != "user" {
continue;
}
if let Some(content) = msg.get_mut("content").and_then(|c| c.as_array_mut()) {
for block in content.iter_mut() {
if block.get("type").and_then(|t| t.as_str()) != Some("tool_result") {
continue;
}
let name = block
.get("tool_use_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);
let excluded =
name.is_some_and(|n| cfg.proxy.is_tool_live_compress_excluded(n));
if live_compress
&& !excluded
&& let Some(inner_content) = block.get_mut("content")
{
modified |= compress_content_field(inner_content, name, kind);
}
}
}
}
if let Some(a) = user_aggr {
let end = boundary.min(messages.len());
let start = protect.min(end);
for msg in &mut messages[start..end] {
if msg.get("role").and_then(|r| r.as_str()) == Some("user")
&& let Some(content) = msg.get_mut("content").and_then(|c| c.as_array_mut())
{
prose_segments += u64::from(prose::compress_text_blocks(content, a));
}
}
}
}
if prose_segments > 0 {
modified = true;
}
if align_volatile
&& cached == 0
&& let Some(system) = doc.get("system")
&& !prose::value_has_cache_control(system)
&& let Some(text) = super::cache_aligner::system_text(system)
{
let scan = super::cache_aligner::scan_volatile(&text);
cache_safety::record_volatile_system(scan.fields as u64);
}
if relocate_volatile
&& arm == super::holdout::Arm::Treatment
&& cached == 0
&& doc
.get("system")
.is_some_and(|s| !prose::value_has_cache_control(s))
{
let relocated = super::cache_aligner::apply_anthropic_relocate(&mut doc);
if relocated > 0 {
modified = true;
cache_safety::record_volatile_relocated(relocated as u64);
}
}
if inject_breakpoint
&& cached == 0
&& doc
.get("system")
.is_some_and(|s| !prose::value_has_cache_control(s))
&& super::cache_breakpoint::inject_anthropic_system(&mut doc)
{
modified = true;
cache_safety::record_breakpoint_injected();
}
if repack {
cache_safety::record_cold_repack();
}
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_content_field(
content: &mut Value,
tool_name: Option<&str>,
kind: ToolResultKind,
) -> bool {
match content {
Value::String(s) => super::tool_output::compress_text(s, tool_name, kind),
Value::Array(arr) => {
let mut modified = false;
for item in arr.iter_mut() {
if item.get("type").and_then(|t| t.as_str()) == Some("text")
&& let Some(Value::String(text)) = item.get_mut("text")
{
modified |= super::tool_output::compress_text(text, tool_name, kind);
}
}
modified
}
_ => false,
}
}
#[cfg(test)]
mod tests {
use super::super::compress::compress_tool_result;
use super::*;
fn source_file_body() -> Vec<u8> {
let code = (0..60)
.map(|i| format!(" let binding_{i} = compute_value_{i}(context, options);"))
.collect::<Vec<_>>()
.join("\n");
let body = serde_json::json!({
"model": "claude-opus-4-8",
"messages": [
{
"role": "assistant",
"content": [{"type": "tool_use", "id": "toolu_1", "name": "Read", "input": {"file_path": "src/app.rs"}}]
},
{
"role": "user",
"content": [{"type": "tool_result", "tool_use_id": "toolu_1", "content": code}]
}
]
});
serde_json::to_vec(&body).unwrap()
}
#[test]
fn read_tool_result_is_never_truncated() {
let bytes = source_file_body();
let body: Value = serde_json::from_slice(&bytes).unwrap();
let (out, _orig, _comp) = compress_request_body(body, bytes.len());
let parsed: Value = serde_json::from_slice(&out).unwrap();
let content = parsed["messages"][1]["content"][0]["content"]
.as_str()
.unwrap();
assert!(
content.contains("binding_59"),
"the full source body must survive — refactors need it intact"
);
assert!(!content.contains("lines omitted"));
}
fn forge_log_body(tool_name: &str) -> Value {
let mut log = String::new();
for i in 0..90 {
log.push_str(&format!(
"INFO processing item {i}: ok, latency={i}ms, queue depth normal, retries 0\n"
));
}
serde_json::json!({
"messages": [
{"role": "assistant", "content": [{"type": "tool_use", "id": "f1", "name": tool_name, "input": {}}]},
{"role": "user", "content": [{"type": "tool_result", "tool_use_id": "f1", "content": log}]}
]
})
}
#[test]
fn forge_shell_tool_result_compresses() {
let body = forge_log_body("forge_shell");
let bytes = serde_json::to_vec(&body).unwrap();
let (_out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "foreign shell output must be compressed");
}
#[test]
fn foreign_read_tool_protects_source() {
let code = (0..60)
.map(|i| format!(" let binding_{i} = compute_value_{i}(context, options);"))
.collect::<Vec<_>>()
.join("\n");
let body = serde_json::json!({
"messages": [
{"role": "assistant", "content": [{"type": "tool_use", "id": "r1", "name": "forge_read", "input": {"path": "src/app.rs"}}]},
{"role": "user", "content": [{"type": "tool_result", "tool_use_id": "r1", "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();
let content = parsed["messages"][1]["content"][0]["content"]
.as_str()
.unwrap();
assert!(
content.contains("binding_59"),
"source body must survive intact"
);
}
#[test]
fn compress_request_body_is_deterministic() {
let _lock = crate::core::data_dir::test_env_lock();
let bytes = serde_json::to_vec(&forge_log_body("Bash")).unwrap();
let a = compress_request_body(serde_json::from_slice(&bytes).unwrap(), bytes.len()).0;
let b = compress_request_body(serde_json::from_slice(&bytes).unwrap(), bytes.len()).0;
assert_eq!(a, b, "identical input must yield byte-identical output");
}
fn big_log() -> String {
(0..200)
.map(|i| format!("[info] processed item {i:04} ok, latency {i}ms, queue normal"))
.collect::<Vec<_>>()
.join("\n")
}
#[test]
fn inband_ccr_emit_echo_splice_round_trip() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CCR_INBAND");
crate::core::config::Config::update_global(|c| {
c.proxy.ccr_inband = Some(true);
})
.unwrap();
let log = big_log();
let emit = serde_json::json!({
"messages": [
{"role": "assistant", "content": [{"type": "tool_use", "id": "t1", "name": "bash", "input": {}}]},
{"role": "user", "content": [{"type": "tool_result", "tool_use_id": "t1", "content": log}]}
]
});
let bytes = serde_json::to_vec(&emit).unwrap();
let (out, _o, _c) = compress_request_body(emit, bytes.len());
let emitted: Value = serde_json::from_slice(&out).unwrap();
let stub = emitted["messages"][1]["content"][0]["content"]
.as_str()
.unwrap();
assert!(
stub.contains("<lc_expand:"),
"in-band stub must advertise an echo-able marker: {stub}"
);
assert!(
!stub.contains("/tee/proxy_"),
"in-band stub must not leak the unreachable local tee path: {stub}"
);
let start = stub.find("<lc_expand:").unwrap();
let end = stub[start..].find('>').unwrap() + start + 1;
let marker = &stub[start..end];
let echo = serde_json::json!({
"messages": [
{"role": "user", "content": [{"type": "text", "text": "look again"}]},
{"role": "assistant", "content": format!("revisiting that output: {marker}")}
]
});
let bytes = serde_json::to_vec(&echo).unwrap();
let (out, _o, _c) = compress_request_body(echo, bytes.len());
let spliced: Value = serde_json::from_slice(&out).unwrap();
let assistant = spliced["messages"][1]["content"].as_str().unwrap();
assert!(
assistant.contains("processed item 0007 ok")
&& assistant.contains("processed item 0199 ok"),
"the verbatim original must be spliced back in full: {assistant}"
);
assert!(
!assistant.contains("<lc_expand:"),
"the marker must be consumed by the splice"
);
}
#[test]
fn inband_marker_less_turn_is_byte_identical_on_or_off() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CCR_INBAND");
let body = serde_json::json!({
"messages": [
{"role": "user", "content": [{"type": "text", "text": "hello there"}]},
{"role": "assistant", "content": "hi — how can I help?"}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| c.proxy.ccr_inband = Some(false)).unwrap();
let off = compress_request_body(body.clone(), bytes.len()).0;
crate::core::config::Config::update_global(|c| c.proxy.ccr_inband = Some(true)).unwrap();
let on = compress_request_body(body, bytes.len()).0;
assert_eq!(
off, on,
"a marker-less request must be byte-identical whether in-band is on or off"
);
}
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_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);
c.proxy.role_aggressiveness.user = Some(0.6);
})
.unwrap();
let prose = big_prose();
let assistant_text = big_prose();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": prose,
"messages": [
{"role": "user", "content": [{"type": "text", "text": prose}]},
{"role": "assistant", "content": assistant_text},
]
});
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["system"].as_str().unwrap().len() < prose.len(),
"system prose must be compressed when enabled"
);
assert_eq!(
parsed["messages"][1]["content"].as_str().unwrap(),
assistant_text,
"assistant turns must pass through verbatim (#710)"
);
}
#[test]
fn 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);
})
.unwrap();
let prose = big_prose();
let mut messages = Vec::new();
for i in 0..30 {
let role = if i % 2 == 0 { "user" } else { "assistant" };
messages.push(serde_json::json!({
"role": role,
"content": [{"type": "text", "text": prose}]
}));
}
let body = serde_json::json!({ "messages": messages });
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();
let frozen_user = parsed["messages"][0]["content"][0]["text"]
.as_str()
.unwrap();
assert!(
frozen_user.len() < prose.len(),
"user prose in the frozen region must be compressed"
);
assert_eq!(
parsed["messages"][1]["content"][0]["text"]
.as_str()
.unwrap(),
prose,
"assistant prose is never compressed"
);
let live_tail_user = parsed["messages"][28]["content"][0]["text"]
.as_str()
.unwrap();
assert_eq!(
live_tail_user, prose,
"user prose in the live tail (>= boundary) must be preserved for quality"
);
}
#[test]
fn client_cached_prefix_disables_system_prose() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.9);
})
.unwrap();
let prose = big_prose();
let body = serde_json::json!({
"system": prose,
"messages": [
{"role": "user", "content": [
{"type": "text", "text": "hi", "cache_control": {"type": "ephemeral"}}
]},
{"role": "assistant", "content": "ok"}
]
});
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_eq!(
parsed["system"].as_str().unwrap(),
prose,
"system must stay verbatim when the client caches a message prefix (#448)"
);
}
#[test]
fn 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!({"system": prose, "messages": [{"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,
"prose compression must be byte-identical for identical input (#498)"
);
}
#[test]
fn bash_tool_result_still_compresses() {
let log = {
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..90 {
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
};
let body = serde_json::json!({
"messages": [
{"role": "assistant", "content": [{"type": "tool_use", "id": "t1", "name": "Bash", "input": {}}]},
{"role": "user", "content": [{"type": "tool_result", "tool_use_id": "t1", "content": log}]}
]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (_out, orig, comp) = compress_request_body(body, bytes.len());
assert!(comp < orig, "shell output must still be compressed");
}
#[test]
fn json_envelope_tool_result_is_compressed() {
let _iso = crate::core::data_dir::isolated_data_dir();
let log = long_git_status();
let expected = compress_tool_result(&log, Some("Bash"));
let envelope = serde_json::to_string(&serde_json::json!({
"content": [{"type": "text", "text": log}],
"isError": false,
}))
.unwrap();
let body = serde_json::json!({
"messages": [
{"role": "assistant", "content": [{
"type": "tool_use",
"id": "t1",
"name": "Bash",
"input": {}
}]},
{"role": "user", "content": [{
"type": "tool_result",
"tool_use_id": "t1",
"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 result should shrink");
let parsed: Value = serde_json::from_slice(&out).unwrap();
let content = parsed["messages"][1]["content"][0]["content"]
.as_str()
.unwrap();
let envelope: Value = serde_json::from_str(content).unwrap();
assert_eq!(envelope["content"][0]["text"].as_str().unwrap(), expected);
}
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 cached_prefix_body(first_text: &str, prose: &str) -> (Vec<Value>, Value) {
let messages = vec![
serde_json::json!({"role": "user", "content": [
{"type": "text", "text": first_text, "cache_control": {"type": "ephemeral"}}
]}),
serde_json::json!({"role": "assistant", "content": "ok"}),
];
let body = serde_json::json!({ "system": prose, "messages": messages.clone() });
(messages, body)
}
#[test]
fn cold_prefix_repack_rewrites_protected_system_prose_when_enabled() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_COLD_PREFIX_REPACK");
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.9);
c.proxy.cold_prefix_repack = Some(true);
})
.unwrap();
let prose = big_prose().repeat(6);
let (messages, body) = cached_prefix_body("cold-repack-enabled-session", &prose);
super::super::cold_prefix::test_seed_last_touch(&messages, 3 * 60 * 60);
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["system"].as_str().unwrap().len() < prose.len(),
"a predicted-cold prefix must let the proxy repack the otherwise-protected system prose"
);
}
#[test]
fn cold_prefix_repack_skipped_for_subcacheable_prefix_by_default() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_COLD_PREFIX_REPACK");
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_POLICY");
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.9);
c.proxy.cold_prefix_repack = Some(true);
})
.unwrap();
let prose = big_prose();
let (messages, body) = cached_prefix_body("cold-repack-subcacheable-session", &prose);
super::super::cold_prefix::test_seed_last_touch(&messages, 3 * 60 * 60);
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_eq!(
parsed["system"].as_str().unwrap(),
prose,
"the net-cost gate must skip repacking a sub-cacheable prefix (premium default)"
);
}
#[test]
fn cold_prefix_repack_off_by_default_keeps_prefix_protected() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_COLD_PREFIX_REPACK");
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.9);
c.proxy.cold_prefix_repack = Some(false);
})
.unwrap();
let prose = big_prose();
let (messages, body) = cached_prefix_body("cold-repack-disabled-session", &prose);
super::super::cold_prefix::test_seed_last_touch(&messages, 24 * 60 * 60);
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_eq!(
parsed["system"].as_str().unwrap(),
prose,
"with repack off the cached prefix stays byte-stable regardless of idle time (#448)"
);
}
#[test]
fn cold_prefix_repack_protects_warm_prefix_even_when_enabled() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_COLD_PREFIX_REPACK");
crate::core::config::Config::update_global(|c| {
c.proxy.role_aggressiveness.system = Some(0.9);
c.proxy.cold_prefix_repack = Some(true);
})
.unwrap();
let prose = big_prose();
let (messages, body) = cached_prefix_body("cold-repack-warm-session", &prose);
super::super::cold_prefix::test_seed_last_touch(&messages, 60);
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_eq!(
parsed["system"].as_str().unwrap(),
prose,
"a warm prefix must stay protected even with repack enabled — only LARGE gaps trigger"
);
}
#[test]
fn cache_policy_attribution_is_measurement_only() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_POLICY");
crate::test_env::remove_var("LEAN_CTX_PROXY_COLD_PREFIX_REPACK");
let prose = big_prose();
let (_messages, body) = cached_prefix_body("cache-policy-measurement-session", &prose);
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.cache_policy = Some(false);
})
.unwrap();
let (off, _o, _c) = compress_request_body(body.clone(), bytes.len());
crate::core::config::Config::update_global(|c| {
c.proxy.cache_policy = Some(true);
})
.unwrap();
let (on, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
off, on,
"miss attribution is measurement-only: the wire bytes must not change"
);
}
#[test]
fn effort_control_dials_adaptive_thinking_only() {
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("medium".into());
})
.unwrap();
let adaptive = serde_json::json!({
"model": "claude-opus-4-8",
"thinking": {"type": "adaptive"},
"messages": [{"role": "user", "content": "hi"}]
});
let bytes = serde_json::to_vec(&adaptive).unwrap();
let (out, _o, _c) = compress_request_body(adaptive, bytes.len());
assert_eq!(
serde_json::from_slice::<Value>(&out).unwrap()["output_config"]["effort"],
"medium"
);
let plain = serde_json::json!({
"model": "claude-opus-4-8",
"messages": [{"role": "user", "content": "hi"}]
});
let bytes = serde_json::to_vec(&plain).unwrap();
let (out, _o, _c) = compress_request_body(plain, bytes.len());
assert!(
serde_json::from_slice::<Value>(&out)
.unwrap()
.get("output_config")
.is_none()
);
}
#[test]
fn verbosity_steer_applies_to_treatment_skips_control() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_VERBOSITY_STEER");
crate::test_env::remove_var("LEAN_CTX_PROXY_OUTPUT_HOLDOUT");
crate::test_env::remove_var("LEAN_CTX_PROXY_EFFORT");
let req = serde_json::json!({
"model": "claude-opus-4-8",
"messages": [{"role": "user", "content": "Summarize the design."}]
});
let bytes = serde_json::to_vec(&req).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.verbosity_steer = Some(true);
c.proxy.output_holdout = Some(0.0);
})
.unwrap();
let (out, _o, _c) = compress_request_body(req.clone(), bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert!(
v["messages"][0]["content"]
.as_str()
.unwrap()
.contains(crate::proxy::verbosity::STEER),
"treatment arm must receive the verbosity steer"
);
crate::core::config::Config::update_global(|c| {
c.proxy.output_holdout = Some(1.0);
})
.unwrap();
let (out2, _o, _c) = compress_request_body(req, bytes.len());
let v2: Value = serde_json::from_slice(&out2).unwrap();
assert!(
!v2["messages"][0]["content"]
.as_str()
.unwrap()
.contains(crate::proxy::verbosity::STEER),
"control arm must NOT be steered (measurement baseline)"
);
}
fn cacheable_system() -> String {
"You are a careful, senior software engineer who writes maintainable code. ".repeat(400)
}
#[test]
fn cache_breakpoint_injected_on_unanchored_system_when_opt_in() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_BREAKPOINT");
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| c.proxy.cache_breakpoint = Some(true))
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert_eq!(
v["system"][0]["cache_control"]["type"], "ephemeral",
"an unanchored system prompt must receive one ephemeral breakpoint"
);
assert!(
v["system"][0]["text"]
.as_str()
.unwrap()
.contains("senior software engineer"),
"the system text must be preserved verbatim under the marker"
);
}
#[test]
fn cache_breakpoint_off_by_default_is_byte_unchanged() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_BREAKPOINT");
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
out, bytes,
"default-off must leave the request byte-identical"
);
}
#[test]
fn cache_breakpoint_respects_client_anchor() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_BREAKPOINT");
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system(),
"messages": [{
"role": "user",
"content": [{
"type": "text",
"text": "hello",
"cache_control": {"type": "ephemeral"}
}]
}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| c.proxy.cache_breakpoint = Some(true))
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert!(
v["system"].is_string(),
"with a client anchor present, system must be left untouched (no second breakpoint)"
);
}
#[test]
fn cache_aligner_measures_without_mutating_body() {
let _iso = crate::core::data_dir::isolated_data_dir();
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_ALIGNER");
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_BREAKPOINT");
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": "Today is 2026-06-22. Session 550e8400-e29b-41d4-a716-446655440000.",
"messages": [{"role": "user", "content": "Hello."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| c.proxy.cache_aligner = Some(true)).unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
out, bytes,
"cache-aligner telemetry must never mutate the request body"
);
}
fn cacheable_system_with_date() -> String {
format!("Today is 2026-06-27. {}", cacheable_system())
}
fn clear_relocate_env() {
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_ALIGN_RELOCATE");
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_BREAKPOINT");
crate::test_env::remove_var("LEAN_CTX_PROXY_CACHE_ALIGNER");
crate::test_env::remove_var("LEAN_CTX_PROXY_OUTPUT_HOLDOUT");
}
#[test]
fn cache_align_relocate_moves_volatiles_to_tail_when_opt_in() {
let _iso = crate::core::data_dir::isolated_data_dir();
clear_relocate_env();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system_with_date(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.cache_align_relocate = Some(true);
c.proxy.output_holdout = Some(0.0);
})
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert!(v["system"].is_array(), "system reshaped into a block array");
assert_eq!(v["system"][0]["cache_control"]["type"], "ephemeral");
assert!(
!v["system"][0]["text"]
.as_str()
.unwrap()
.contains("2026-06-27"),
"the volatile date must leave the cacheable prefix"
);
assert!(
v["system"][1].get("cache_control").is_none(),
"the relocated tail block stays uncached"
);
assert!(
v["system"][1]["text"]
.as_str()
.unwrap()
.contains("2026-06-27"),
"the date must be re-stated in the tail"
);
}
#[test]
fn cache_align_relocate_off_by_default_is_byte_unchanged() {
let _iso = crate::core::data_dir::isolated_data_dir();
clear_relocate_env();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system_with_date(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
assert_eq!(
out, bytes,
"default-off must leave the request byte-identical"
);
}
#[test]
fn cache_align_relocate_skips_control_arm() {
let _iso = crate::core::data_dir::isolated_data_dir();
clear_relocate_env();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system_with_date(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.cache_align_relocate = Some(true);
c.proxy.output_holdout = Some(1.0);
})
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert!(
v["system"].is_string(),
"control arm must not be relocated (measurement baseline)"
);
}
#[test]
fn cache_align_relocate_respects_client_anchor() {
let _iso = crate::core::data_dir::isolated_data_dir();
clear_relocate_env();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system_with_date(),
"messages": [{
"role": "user",
"content": [{
"type": "text",
"text": "hello",
"cache_control": {"type": "ephemeral"}
}]
}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.cache_align_relocate = Some(true);
c.proxy.output_holdout = Some(0.0);
})
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
assert!(
v["system"].is_string(),
"with a client anchor present, system must be left untouched"
);
}
#[test]
fn cache_align_relocate_composes_with_breakpoint_to_one_anchor() {
let _iso = crate::core::data_dir::isolated_data_dir();
clear_relocate_env();
let body = serde_json::json!({
"model": "claude-opus-4-8",
"system": cacheable_system_with_date(),
"messages": [{"role": "user", "content": "Refactor the parser."}]
});
let bytes = serde_json::to_vec(&body).unwrap();
crate::core::config::Config::update_global(|c| {
c.proxy.cache_align_relocate = Some(true);
c.proxy.cache_breakpoint = Some(true);
c.proxy.output_holdout = Some(0.0);
})
.unwrap();
let (out, _o, _c) = compress_request_body(body, bytes.len());
let v: Value = serde_json::from_slice(&out).unwrap();
let blocks = v["system"].as_array().expect("system is a block array");
assert_eq!(blocks.len(), 2, "stable block + volatile tail");
assert_eq!(blocks[0]["cache_control"]["type"], "ephemeral");
assert!(
blocks[1].get("cache_control").is_none(),
"no second breakpoint — the tail stays uncached"
);
}
}