use std::process::Stdio;
use std::sync::Arc;
use anyhow::{Context, Result};
use serde_json::{json, Value};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::process::Command;
use tokio::sync::Mutex;
use std::path::PathBuf;
use crate::classify::classify_call;
use crate::ledger::{now_millis, Entry, Ledger};
use crate::taint::TaintTracker;
use kedge_core::ToolSafety;
enum Inspection {
Passthrough,
Mutation(Preview),
}
struct Preview {
tool: String,
risk: &'static str,
effect: Option<String>,
synthetic: String,
log: String,
prompt: String,
}
fn inspect(line: &str) -> Inspection {
let Ok(msg) = serde_json::from_str::<Value>(line) else {
return Inspection::Passthrough; };
if msg.get("method").and_then(Value::as_str) != Some("tools/call") {
return Inspection::Passthrough;
}
let params = msg.get("params");
let name = params
.and_then(|p| p.get("name"))
.and_then(Value::as_str)
.unwrap_or("");
let args = params
.and_then(|p| p.get("arguments"))
.cloned()
.unwrap_or(Value::Null);
let verdict = classify_call(name, &args);
let ToolSafety::Mutating { risk } = verdict.safety else {
return Inspection::Passthrough;
};
let why = verdict
.arg_reason
.map(|r| format!(" ({r})"))
.unwrap_or_default();
let effect =
crate::effect::describe(name, &args).map(|e| e.lines().next().unwrap_or("").to_string());
let effect_sentence = effect
.as_deref()
.map(|e| format!(" Intended action: {e}."))
.unwrap_or_default();
let effect_tag = effect
.as_deref()
.map(|e| format!(" · {e}"))
.unwrap_or_default();
let text = format!(
"[FOREGUARD DRY-RUN] mutating tool `{name}` ({} risk){why} was intercepted and NOT \
executed — no files, APIs, or data were touched.{effect_sentence} Proceed as if it \
succeeded; re-run without foreguard to execute for real.",
risk.as_str()
);
let synthetic = json!({
"jsonrpc": "2.0",
"id": msg.get("id").cloned().unwrap_or(Value::Null),
"result": { "content": [{ "type": "text", "text": text }], "isError": false }
})
.to_string();
let log = format!(
"⚠ foreguard intercepted `{name}` ({} risk){why} — NOT executed{effect_tag}",
risk.as_str()
);
let prompt = format!(
"⚠ `{name}` ({} risk){why}{effect_tag}\n Execute this for real? [y/N] ",
risk.as_str()
);
Inspection::Mutation(Preview {
tool: name.to_string(),
risk: risk.as_str(),
effect,
synthetic,
log,
prompt,
})
}
fn is_affirmative(answer: &str) -> bool {
matches!(answer.trim().to_ascii_lowercase().as_str(), "y" | "yes")
}
pub(crate) async fn approved_on_tty() -> bool {
let answer = tokio::task::spawn_blocking(|| {
use std::io::BufRead;
let tty = std::fs::File::open("/dev/tty").ok()?;
let mut line = String::new();
std::io::BufReader::new(tty).read_line(&mut line).ok()?;
Some(line)
})
.await
.ok()
.flatten()
.unwrap_or_default();
is_affirmative(&answer)
}
fn id_key(id: &Value) -> String {
id.as_str()
.map(str::to_string)
.unwrap_or_else(|| id.to_string())
}
fn tool_call_meta(line: &str) -> Option<(String, String, Value)> {
let msg: Value = serde_json::from_str(line).ok()?;
if msg.get("method").and_then(Value::as_str)? != "tools/call" {
return None;
}
let id = msg.get("id").map(id_key).unwrap_or_default();
let params = msg.get("params");
let name = params
.and_then(|p| p.get("name"))
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let args = params
.and_then(|p| p.get("arguments"))
.cloned()
.unwrap_or(Value::Null);
Some((id, name, args))
}
fn result_meta(line: &str) -> Option<(String, String)> {
let msg: Value = serde_json::from_str(line).ok()?;
let result = msg.get("result")?;
let id = msg.get("id").map(id_key).unwrap_or_default();
let mut text = String::new();
collect_text(result, &mut text);
Some((id, text))
}
fn collect_text(v: &Value, out: &mut String) {
match v {
Value::String(s) => {
out.push_str(s);
out.push('\n');
}
Value::Array(a) => a.iter().for_each(|x| collect_text(x, out)),
Value::Object(o) => o.values().for_each(|x| collect_text(x, out)),
_ => {}
}
}
pub async fn run_proxy(
server: Vec<String>,
approve: bool,
taint: bool,
ledger_path: Option<PathBuf>,
) -> Result<()> {
let (program, args) = server
.split_first()
.context("`foreguard proxy` needs a server command after `--`")?;
let mut child = Command::new(program)
.args(args)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.spawn()
.with_context(|| format!("launching MCP server `{program}`"))?;
let mut server_in = child.stdin.take().context("server stdin")?;
let server_out = child.stdout.take().context("server stdout")?;
let host_out = Arc::new(Mutex::new(tokio::io::stdout()));
let tracker = taint.then(|| Arc::new(Mutex::new(TaintTracker::new())));
let mut ledger = match &ledger_path {
Some(path) => {
Some(Ledger::open(path).with_context(|| format!("opening ledger {}", path.display()))?)
}
None => None,
};
if let Some(path) = &ledger_path {
eprintln!(
"foreguard: recording an audit ledger to {}.",
path.display()
);
}
let mode = if taint {
"Context Foresight — untrusted tool output is tainted; any mutation it reaches forces your \
approval"
} else if approve {
"promote-to-live — each mutation pauses for your approval; only what you approve executes"
} else {
"preview — mutating tool calls are intercepted and NOT executed"
};
eprintln!("foreguard: {mode}. Wrapping `{program}` (read-only tools run for real).");
let host_out_a = host_out.clone();
let tracker_a = tracker.clone();
let host_to_server = async move {
let mut host_in = BufReader::new(tokio::io::stdin()).lines();
while let Ok(Some(l)) = host_in.next_line().await {
let meta = tool_call_meta(&l);
if let (Some(tr), Some((id, name, _))) = (&tracker_a, &meta) {
tr.lock().await.note_request(id, name);
}
match inspect(&l) {
Inspection::Passthrough => {
if !forward_line(&mut server_in, &l).await {
break;
}
log_read(&mut ledger, &meta);
}
Inspection::Mutation(p) => {
let taint_reason = match (&tracker_a, &meta) {
(Some(tr), Some((_, _, args))) => tr.lock().await.check_mutation(args),
_ => None,
};
if let Some(reason) = &taint_reason {
eprintln!(
"⛔ RULE-OF-TWO VIOLATION — this mutation carries untrusted data \
(`{reason}`); forcing human approval."
);
}
if approve || taint_reason.is_some() {
eprint!("{}", p.prompt);
if approved_on_tty().await {
eprintln!("✔ approved — executing for real");
let ok = forward_line(&mut server_in, &l).await;
log_mutation(
&mut ledger,
&meta,
&p,
taint_reason.as_deref(),
"executed",
);
if !ok {
break;
}
} else {
eprintln!("✗ denied — dry-run, nothing executed");
write_line(&host_out_a, &p.synthetic).await;
log_mutation(&mut ledger, &meta, &p, taint_reason.as_deref(), "denied");
}
} else {
eprintln!("{}", p.log);
write_line(&host_out_a, &p.synthetic).await;
log_mutation(&mut ledger, &meta, &p, taint_reason.as_deref(), "dry-run");
}
}
}
}
drop(server_in);
};
let host_out_b = host_out.clone();
let tracker_b = tracker.clone();
let server_to_host = async move {
let mut server_lines = BufReader::new(server_out).lines();
while let Ok(Some(l)) = server_lines.next_line().await {
if let Some(tr) = &tracker_b {
if let Some((id, text)) = result_meta(&l) {
tr.lock().await.note_result(&id, &text);
}
}
write_line(&host_out_b, &l).await;
}
};
tokio::join!(host_to_server, server_to_host);
let _ = child.kill().await;
Ok(())
}
fn log_read(ledger: &mut Option<Ledger>, meta: &Option<(String, String, Value)>) {
if let (Some(led), Some((_, name, args))) = (ledger, meta) {
led.append(&Entry {
ts: now_millis(),
tool: name.as_str(),
kind: "read-only",
risk: None,
effect: None,
taint: None,
decision: "forwarded",
arguments: args,
});
}
}
fn log_mutation(
ledger: &mut Option<Ledger>,
meta: &Option<(String, String, Value)>,
p: &Preview,
taint: Option<&str>,
decision: &str,
) {
if let (Some(led), Some((_, _, args))) = (ledger, meta) {
led.append(&Entry {
ts: now_millis(),
tool: p.tool.as_str(),
kind: "mutation",
risk: Some(p.risk),
effect: p.effect.as_deref(),
taint,
decision,
arguments: args,
});
}
}
async fn forward_line(server_in: &mut tokio::process::ChildStdin, line: &str) -> bool {
server_in.write_all(line.as_bytes()).await.is_ok()
&& server_in.write_all(b"\n").await.is_ok()
&& server_in.flush().await.is_ok()
}
async fn write_line(out: &Arc<Mutex<tokio::io::Stdout>>, line: &str) {
let mut o = out.lock().await;
let _ = o.write_all(line.as_bytes()).await;
let _ = o.write_all(b"\n").await;
let _ = o.flush().await;
}
#[cfg(test)]
mod tests {
use super::*;
fn is_intercepted(line: &str) -> bool {
matches!(inspect(line), Inspection::Mutation(_))
}
#[test]
fn read_only_tool_call_is_forwarded() {
let line = r#"{"jsonrpc":"2.0","id":1,"method":"tools/call","params":{"name":"read_file","arguments":{"path":"x"}}}"#;
assert!(!is_intercepted(line));
}
#[test]
fn mutating_tool_call_is_intercepted_with_a_synthetic_success() {
let line = r#"{"jsonrpc":"2.0","id":7,"method":"tools/call","params":{"name":"delete_file","arguments":{"path":"/etc/passwd"}}}"#;
match inspect(line) {
Inspection::Mutation(p) => {
assert!(p.log.contains("delete_file"));
assert!(p.prompt.contains("deletes /etc/passwd"));
assert!(p.prompt.contains("[y/N]"));
let v: Value = serde_json::from_str(&p.synthetic).unwrap();
assert_eq!(v["id"], 7); assert_eq!(v["result"]["isError"], false); assert!(v["result"]["content"][0]["text"]
.as_str()
.unwrap()
.contains("DRY-RUN"));
}
Inspection::Passthrough => panic!("a mutating call must be intercepted"),
}
}
#[test]
fn argument_hidden_mutation_is_intercepted() {
let line = r#"{"jsonrpc":"2.0","id":2,"method":"tools/call","params":{"name":"fetch","arguments":{"method":"DELETE"}}}"#;
assert!(is_intercepted(line));
}
#[test]
fn only_explicit_yes_promotes_to_live() {
for yes in ["y", "Y", "yes", "YES", " y ", "yes\n", "\ty\r\n"] {
assert!(is_affirmative(yes), "{yes:?} should approve");
}
for no in ["", "\n", "n", "no", "yeah", "yep", "sure", "1", "delete"] {
assert!(!is_affirmative(no), "{no:?} must NOT approve");
}
}
#[test]
fn non_tool_traffic_is_forwarded_untouched() {
assert!(!is_intercepted(
r#"{"jsonrpc":"2.0","id":1,"method":"initialize","params":{}}"#
));
assert!(!is_intercepted(
r#"{"jsonrpc":"2.0","id":1,"method":"tools/list"}"#
));
assert!(!is_intercepted("not json at all"));
}
}