Skip to main content

systemprompt_api/services/proxy/audit/
jsonrpc.rs

1//! Minimal JSON-RPC / MCP frame parsing for the tool-call audit tap.
2//!
3//! The gateway forwards MCP frames verbatim; to audit a `tools/call` it parses
4//! the tool name and arguments from the request and the result from the
5//! response, matching them by JSON-RPC id. The `arguments`, `result`, and
6//! `content` payloads are `serde_json::Value` because MCP defines them as
7//! open-shaped at the wire boundary. The matching response frame is also
8//! stamped with the execution id the tap minted, under the systemprompt
9//! `_meta` key, so a client that reports the result later carries the exact
10//! server key.
11//!
12//! Copyright (c) systemprompt.io — Business Source License 1.1.
13//! See <https://systemprompt.io> for licensing details.
14
15use serde::Deserialize;
16use serde_json::{Value, json};
17use systemprompt_models::artifacts::EXECUTION_META_KEY;
18
19const TOOLS_CALL_METHOD: &str = "tools/call";
20
21#[derive(Deserialize)]
22struct RequestFrame {
23    #[serde(default)]
24    id: Option<Value>,
25    method: String,
26    #[serde(default)]
27    params: Option<ToolCallParams>,
28}
29
30#[derive(Deserialize)]
31struct ToolCallParams {
32    name: String,
33    #[serde(default)]
34    arguments: Option<Value>,
35}
36
37#[derive(Debug)]
38pub struct ToolCallInvocation {
39    pub id: Value,
40    pub tool_name: String,
41    pub arguments: Value,
42}
43
44pub fn parse_tool_call(body: &[u8]) -> Option<ToolCallInvocation> {
45    let frame: RequestFrame = serde_json::from_slice(body).ok()?;
46    if frame.method != TOOLS_CALL_METHOD {
47        return None;
48    }
49    let params = frame.params?;
50    Some(ToolCallInvocation {
51        id: frame.id.unwrap_or(Value::Null),
52        tool_name: params.name,
53        arguments: params.arguments.unwrap_or(Value::Null),
54    })
55}
56
57#[derive(Deserialize)]
58struct ResponseFrame {
59    #[serde(default)]
60    result: Option<ToolCallResult>,
61    #[serde(default)]
62    error: Option<Value>,
63}
64
65#[derive(Deserialize)]
66struct ToolCallResult {
67    #[serde(default, rename = "isError")]
68    is_error: bool,
69    #[serde(default, rename = "structuredContent")]
70    structured_content: Option<Value>,
71    #[serde(default)]
72    content: Option<Value>,
73}
74
75#[derive(Debug)]
76pub struct ToolCallOutcome {
77    pub output: Option<Value>,
78    pub error_message: Option<String>,
79    pub result: Option<Value>,
80}
81
82pub fn parse_response_frame(data: &str, request_id: &Value) -> Option<ToolCallOutcome> {
83    let frame: Value = serde_json::from_str(data).ok()?;
84    if frame.get("id") != Some(request_id) {
85        return None;
86    }
87    let parsed: ResponseFrame = serde_json::from_value(frame.clone()).ok()?;
88    if let Some(error) = parsed.error {
89        return Some(ToolCallOutcome {
90            error_message: Some(error.to_string()),
91            output: Some(error),
92            result: None,
93        });
94    }
95    let result = parsed.result?;
96    let output = result.structured_content.or(result.content);
97    let error_message = result
98        .is_error
99        .then(|| "MCP tool call returned isError".to_owned());
100    Some(ToolCallOutcome {
101        output,
102        error_message,
103        result: frame.get("result").cloned(),
104    })
105}
106
107pub fn frame_matches(data: &str, request_id: &Value) -> bool {
108    serde_json::from_str::<Value>(data)
109        .ok()
110        .is_some_and(|frame| frame.get("id") == Some(request_id))
111}
112
113pub fn stamp_execution(data: &str, mcp_execution_id: &str) -> Option<String> {
114    let mut frame: Value = serde_json::from_str(data).ok()?;
115    let result = frame.get_mut("result")?.as_object_mut()?;
116    let meta = result
117        .entry("_meta")
118        .or_insert_with(|| json!({}))
119        .as_object_mut()?;
120    let execution = meta
121        .entry(EXECUTION_META_KEY)
122        .or_insert_with(|| json!({}))
123        .as_object_mut()?;
124    execution
125        .entry("mcp_execution_id")
126        .or_insert_with(|| Value::String(mcp_execution_id.to_owned()));
127    serde_json::to_string(&frame).ok()
128}
129
130pub fn extract_sse_data(frame: &str) -> Option<String> {
131    let mut data = String::new();
132    for line in frame.lines() {
133        if let Some(rest) = line.strip_prefix("data:") {
134            if !data.is_empty() {
135                data.push('\n');
136            }
137            data.push_str(rest.trim_start());
138        }
139    }
140    (!data.is_empty()).then_some(data)
141}
142
143pub fn replace_sse_data(frame: &str, data: &str) -> String {
144    let mut out = String::with_capacity(frame.len() + data.len());
145    let mut wrote = false;
146    for line in frame.trim_end_matches('\n').lines() {
147        if line.starts_with("data:") {
148            if !wrote {
149                out.push_str("data: ");
150                out.push_str(data);
151                out.push('\n');
152                wrote = true;
153            }
154            continue;
155        }
156        out.push_str(line);
157        out.push('\n');
158    }
159    if !wrote {
160        out.push_str("data: ");
161        out.push_str(data);
162        out.push('\n');
163    }
164    out.push('\n');
165    out
166}