Skip to main content

systemprompt_api/services/proxy/audit/
mod.rs

1//! Per-tool audit for external MCP servers served over the HTTP gateway.
2//!
3//! A client-mediated `tools/call` to an external provider has no backend
4//! process to record it, so the gateway taps the forwarded request/response,
5//! writes one `mcp_tool_executions` row under the calling user, and hands the
6//! result to the artifact ingest. The execution id is minted before the
7//! response leaves, and stamped into its `_meta`, so the client's own later
8//! report of the same result carries the exact server key. `record` composes
9//! the tap over the upstream body; the tap owns an [`McpAudit`] and finalizes
10//! it (once) on stream EOF or drop, writing on the process's
11//! [`BackgroundTasks`] so shutdown drains the row.
12//!
13//! The execution row is written first, then paired with the model's intent
14//! through [`IntentClaimService`], mirroring
15//! [`systemprompt_mcp::McpToolExecutor`]: a client-supplied tool-call id is
16//! claimed as an exact pairing, otherwise the newest unclaimed intent for the
17//! tool in the calling session is claimed as an inferred one. Failing to claim
18//! leaves the execution unpaired rather than failing the call, which has
19//! already returned to the client.
20//!
21//! Copyright (c) systemprompt.io — Business Source License 1.1.
22//! See <https://systemprompt.io> for licensing details.
23
24pub mod jsonrpc;
25pub mod tap;
26
27use std::sync::Arc;
28
29use chrono::{DateTime, Utc};
30use serde_json::Value;
31use systemprompt_identifiers::{AiToolCallId, McpExecutionId, McpServerId};
32use systemprompt_mcp::models::{ExecutionStatus, ToolExecutionRequest, ToolExecutionResult};
33use systemprompt_mcp::{
34    ArtifactIngest, INTENT_CLAIM_WINDOW_SECONDS, IngestRequest, IntentClaimService, from_wire_value,
35};
36use systemprompt_models::RequestContext;
37use systemprompt_models::mcp::{Correlation, ExecutionSource};
38use systemprompt_traits::BackgroundTasks;
39
40pub(crate) use jsonrpc::{ToolCallFrame, classify_tool_call, parse_tool_call};
41pub(crate) use tap::record;
42
43use jsonrpc::{ToolCallInvocation, ToolCallOutcome};
44
45#[derive(Debug)]
46pub struct McpAudit {
47    intent_claims: IntentClaimService,
48    ingest: Option<Arc<ArtifactIngest>>,
49    context: RequestContext,
50    server_name: McpServerId,
51    invocation: ToolCallInvocation,
52    started_at: DateTime<Utc>,
53    mcp_execution_id: McpExecutionId,
54    background: BackgroundTasks,
55}
56
57#[derive(Debug, Clone)]
58pub struct AuditSinks {
59    pub intent_claims: IntentClaimService,
60    pub ingest: Option<Arc<ArtifactIngest>>,
61    pub background: BackgroundTasks,
62}
63
64impl McpAudit {
65    pub fn new(
66        sinks: AuditSinks,
67        context: RequestContext,
68        server_name: McpServerId,
69        invocation: ToolCallInvocation,
70    ) -> Self {
71        let AuditSinks {
72            intent_claims,
73            ingest,
74            background,
75        } = sinks;
76        Self {
77            intent_claims,
78            ingest,
79            context,
80            server_name,
81            invocation,
82            started_at: Utc::now(),
83            mcp_execution_id: McpExecutionId::generate(),
84            background,
85        }
86    }
87
88    const fn request_id(&self) -> &Value {
89        &self.invocation.id
90    }
91
92    pub const fn mcp_execution_id(&self) -> &McpExecutionId {
93        &self.mcp_execution_id
94    }
95
96    fn finalize(self, outcome: Option<ToolCallOutcome>) {
97        let (output, error_message, result) = match outcome {
98            Some(o) => (o.output, o.error_message, o.result),
99            None => (
100                None,
101                Some("external MCP tool call produced no parseable result".to_owned()),
102                None,
103            ),
104        };
105
106        let request = ToolExecutionRequest {
107            tool_name: self.invocation.tool_name,
108            server_name: self.server_name.clone(),
109            input: self.invocation.arguments,
110            started_at: self.started_at,
111            context: self.context,
112            request_method: Some("mcp".to_owned()),
113            request_source: Some(String::from(self.server_name)),
114            ai_tool_call_id: None,
115            source: ExecutionSource::Proxy,
116        };
117        let result_row = ToolExecutionResult {
118            status: ExecutionStatus::from_error(error_message.is_some()).to_string(),
119            error_message,
120            output,
121            output_schema: None,
122            started_at: self.started_at,
123            completed_at: Some(Utc::now()),
124        };
125
126        let intent_claims = self.intent_claims;
127        let ingest = self.ingest;
128        let mcp_execution_id = self.mcp_execution_id;
129        self.background.spawn("mcp_proxy_audit", async move {
130            let mut request = request;
131            request.ai_tool_call_id = request.context.ai_tool_call_id().cloned();
132            let exact = request.ai_tool_call_id.clone();
133            let correlation = if exact.is_some() {
134                Correlation::Exact
135            } else {
136                Correlation::Inferred
137            };
138            if let Err(e) = intent_claims
139                .executions()
140                .log_execution_sync_with_id(&mcp_execution_id, &request, &result_row, correlation)
141                .await
142            {
143                tracing::warn!(
144                    tool = %request.tool_name,
145                    server = %request.server_name,
146                    error = %e,
147                    "Failed to record external MCP tool execution"
148                );
149                return;
150            }
151            // Why: no client sends `x-ai-tool-call-id`, so an external-server
152            // execution arrives with nothing to join it to the inference turn
153            // that asked for it. The in-process executor claims the newest
154            // unclaimed intent for the tool in this session; without the same
155            // claim here, every proxied call stayed unpaired — and eight of
156            // nine configured servers are external.
157            match exact {
158                Some(call_id) => {
159                    claim_exact(&intent_claims, &request, &call_id, &mcp_execution_id).await;
160                },
161                None => {
162                    request.ai_tool_call_id =
163                        claim_inferred(&intent_claims, &request, &mcp_execution_id).await;
164                },
165            }
166            if let Some(ingest) = ingest {
167                ingest_proxied_result(&ingest, &request, result, mcp_execution_id).await;
168            }
169        });
170    }
171}
172
173async fn claim_exact(
174    intent_claims: &IntentClaimService,
175    request: &ToolExecutionRequest,
176    call_id: &AiToolCallId,
177    mcp_execution_id: &McpExecutionId,
178) {
179    if let Err(e) = intent_claims.claim_exact(call_id, mcp_execution_id).await {
180        tracing::warn!(
181            tool = %request.tool_name,
182            server = %request.server_name,
183            %mcp_execution_id,
184            %call_id,
185            error = %e,
186            "Proxy exact intent claim failed"
187        );
188    }
189}
190
191async fn claim_inferred(
192    intent_claims: &IntentClaimService,
193    request: &ToolExecutionRequest,
194    mcp_execution_id: &McpExecutionId,
195) -> Option<AiToolCallId> {
196    match intent_claims
197        .claim_inferred(
198            request.context.session_id(),
199            &request.tool_name,
200            mcp_execution_id,
201            INTENT_CLAIM_WINDOW_SECONDS,
202        )
203        .await
204    {
205        Ok(claimed) => claimed,
206        Err(e) => {
207            tracing::warn!(
208                tool = %request.tool_name,
209                server = %request.server_name,
210                %mcp_execution_id,
211                error = %e,
212                "Proxy intent claim failed"
213            );
214            None
215        },
216    }
217}
218
219// JSON: MCP `tools/call` result — open-shaped per the MCP spec, ingested as
220// artifacts.
221async fn ingest_proxied_result(
222    ingest: &ArtifactIngest,
223    request: &ToolExecutionRequest,
224    result: Option<Value>,
225    mcp_execution_id: McpExecutionId,
226) {
227    let Some(wire) = result.as_ref().and_then(from_wire_value) else {
228        return;
229    };
230    let ingest_request = IngestRequest {
231        result: wire,
232        tool_name: request.tool_name.clone(),
233        server_name: Some(request.server_name.clone()),
234        ai_tool_call_id: request.ai_tool_call_id.clone(),
235        mcp_execution_id: Some(mcp_execution_id),
236        ctx: request.context.clone(),
237        skill: None,
238        source: ExecutionSource::Proxy,
239        started_at: Some(request.started_at),
240        input: Some(request.input.clone()),
241    };
242    if let Err(e) = ingest.ingest(ingest_request).await {
243        tracing::warn!(
244            tool = %request.tool_name,
245            server = %request.server_name,
246            error = %e,
247            "Failed to ingest external MCP tool result as an artifact"
248        );
249    }
250}