use anyhow::{Context, Result};
use serde_json::Value;
use std::io::{BufRead, BufReader, BufWriter, Write};
use std::process::{Command, Stdio};
use crate::chunker::count_tokens;
use crate::compress::{compress_output, now_ts};
use crate::store::{log_hook_event, HookEvent};
const MIN_RESULT_TOKENS: usize = 200;
fn proxy_enabled() -> bool {
!std::env::var("TOKENIX_MCP_PROXY").is_ok_and(|v| v == "0")
}
pub fn compress_result_in_place(response: &mut Value) -> Option<(usize, usize)> {
let content = response
.get_mut("result")?
.get_mut("content")?
.as_array_mut()?;
let mut original = 0usize;
let mut compressed = 0usize;
let mut changed = false;
for block in content.iter_mut() {
if block.get("type").and_then(Value::as_str) != Some("text") {
continue;
}
let Some(text) = block.get("text").and_then(Value::as_str) else {
continue;
};
let before = count_tokens(text);
if before < MIN_RESULT_TOKENS {
original += before;
compressed += before;
continue;
}
let new_text = compress_output(text);
let after = count_tokens(&new_text);
original += before;
if after < before {
compressed += after;
changed = true;
block["text"] = Value::String(new_text);
} else {
compressed += before;
}
}
changed.then_some((original, compressed))
}
pub fn run_proxy(name: &str, command: &[String]) -> Result<i32> {
let (program, args) = command
.split_first()
.context("mcp-proxy needs a server command after `--`")?;
let mut child = Command::new(program)
.args(args)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.spawn()
.with_context(|| format!("failed to start MCP server `{program}`"))?;
let mut child_stdin = child.stdin.take().context("child stdin")?;
let child_stdout = child.stdout.take().context("child stdout")?;
std::thread::spawn(move || {
let stdin = std::io::stdin();
let mut reader = stdin.lock();
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) | Err(_) => break,
Ok(_) => {
if child_stdin.write_all(line.as_bytes()).is_err()
|| child_stdin.flush().is_err()
{
break;
}
}
}
}
});
let repo_root = crate::store::find_project_root(
&std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
);
let stdout = std::io::stdout();
let mut out = BufWriter::new(stdout.lock());
let mut reader = BufReader::new(child_stdout);
let mut line = String::new();
let enabled = proxy_enabled();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) | Err(_) => break,
Ok(_) => {}
}
let mut forwarded = line.clone();
if enabled {
if let Ok(mut response) = serde_json::from_str::<Value>(line.trim()) {
if let Some((original, compressed)) = compress_result_in_place(&mut response) {
if let Ok(rewritten) = serde_json::to_string(&response) {
forwarded = format!("{rewritten}\n");
let _ = log_hook_event(
&repo_root,
&HookEvent {
ts: now_ts(),
tool: "MCP".to_string(),
action: "intercepted".to_string(),
phase: "proxy".to_string(),
reason: format!("compressed {name} tool result"),
saved_tokens: (original as i64 - compressed as i64).max(0),
actual_tokens: compressed as i64,
original_estimate: original as i64,
input_preview: String::new(),
command: format!("mcp:{name}"),
},
);
}
}
}
}
if out.write_all(forwarded.as_bytes()).is_err() || out.flush().is_err() {
break;
}
}
let status = child.wait()?;
Ok(status.code().unwrap_or(0))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn text_result(text: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": 7,
"result": { "content": [{ "type": "text", "text": text }] }
})
}
#[test]
fn compresses_a_base64_heavy_tool_result() {
let blob = "Zm9vQmFyBaz1".repeat(200);
let mut response = text_result(&format!("generated image:\ndata:image/png;base64,{blob}"));
let (original, compressed) =
compress_result_in_place(&mut response).expect("should compress");
assert!(compressed < original / 2, "{compressed} vs {original}");
let text = response["result"]["content"][0]["text"].as_str().unwrap();
assert!(text.contains("omitted"));
assert!(text.contains("generated image"), "signal must survive");
}
#[test]
fn leaves_short_results_untouched() {
let mut response = text_result("ok: 3 files changed");
assert!(compress_result_in_place(&mut response).is_none());
assert_eq!(
response["result"]["content"][0]["text"], "ok: 3 files changed",
"short results must be byte-identical"
);
}
#[test]
fn never_rewrites_image_blocks() {
let mut response = json!({
"result": { "content": [{ "type": "image", "data": "iVBORw0KGgo".repeat(200), "mimeType": "image/png" }] }
});
let before = response.clone();
assert!(compress_result_in_place(&mut response).is_none());
assert_eq!(response, before);
}
#[test]
fn ignores_non_result_messages() {
let mut request = json!({"jsonrpc": "2.0", "id": 1, "method": "tools/list"});
assert!(compress_result_in_place(&mut request).is_none());
let mut error = json!({"jsonrpc": "2.0", "id": 1, "error": {"code": -32601}});
assert!(compress_result_in_place(&mut error).is_none());
}
#[test]
fn keeps_result_when_compression_would_not_help() {
let prose = "The quick brown fox jumps over the lazy dog. ".repeat(120);
let mut response = text_result(&prose);
let outcome = compress_result_in_place(&mut response);
assert!(outcome.is_none() || response["result"]["content"][0]["text"] != Value::Null);
assert!(response["result"]["content"][0]["text"]
.as_str()
.unwrap()
.contains("quick brown fox"));
}
}