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.
11//!
12//! Copyright (c) systemprompt.io — Business Source License 1.1.
13//! See <https://systemprompt.io> for licensing details.
14
15pub mod jsonrpc;
16pub mod tap;
17
18use std::sync::Arc;
19
20use chrono::{DateTime, Utc};
21use serde_json::Value;
22use systemprompt_identifiers::McpExecutionId;
23use systemprompt_mcp::models::{ExecutionStatus, ToolExecutionRequest, ToolExecutionResult};
24use systemprompt_mcp::repository::ToolUsageRepository;
25use systemprompt_mcp::{ArtifactIngest, IngestRequest, from_wire_value};
26use systemprompt_models::RequestContext;
27use systemprompt_models::mcp::{Correlation, ExecutionSource};
28
29pub(crate) use jsonrpc::parse_tool_call;
30pub(crate) use tap::record;
31
32use jsonrpc::{ToolCallInvocation, ToolCallOutcome};
33
34#[derive(Debug)]
35pub struct McpAudit {
36    repo: Arc<ToolUsageRepository>,
37    ingest: Option<Arc<ArtifactIngest>>,
38    context: RequestContext,
39    server_name: String,
40    invocation: ToolCallInvocation,
41    started_at: DateTime<Utc>,
42    mcp_execution_id: McpExecutionId,
43}
44
45impl McpAudit {
46    pub fn new(
47        repo: Arc<ToolUsageRepository>,
48        ingest: Option<Arc<ArtifactIngest>>,
49        context: RequestContext,
50        server_name: String,
51        invocation: ToolCallInvocation,
52    ) -> Self {
53        Self {
54            repo,
55            ingest,
56            context,
57            server_name,
58            invocation,
59            started_at: Utc::now(),
60            mcp_execution_id: McpExecutionId::new(uuid::Uuid::new_v4().to_string()),
61        }
62    }
63
64    const fn request_id(&self) -> &Value {
65        &self.invocation.id
66    }
67
68    pub const fn mcp_execution_id(&self) -> &McpExecutionId {
69        &self.mcp_execution_id
70    }
71
72    fn finalize(self, outcome: Option<ToolCallOutcome>) {
73        let (output, error_message, result) = match outcome {
74            Some(o) => (o.output, o.error_message, o.result),
75            None => (
76                None,
77                Some("external MCP tool call produced no parseable result".to_owned()),
78                None,
79            ),
80        };
81
82        let request = ToolExecutionRequest {
83            tool_name: self.invocation.tool_name,
84            server_name: self.server_name.clone(),
85            input: self.invocation.arguments,
86            started_at: self.started_at,
87            context: self.context,
88            request_method: Some("mcp".to_owned()),
89            request_source: Some(self.server_name),
90            ai_tool_call_id: None,
91            source: ExecutionSource::Proxy,
92        };
93        let result_row = ToolExecutionResult {
94            status: ExecutionStatus::from_error(error_message.is_some()).to_string(),
95            error_message,
96            output,
97            output_schema: None,
98            started_at: self.started_at,
99            completed_at: Utc::now(),
100        };
101
102        let repo = self.repo;
103        let ingest = self.ingest;
104        let mcp_execution_id = self.mcp_execution_id;
105        tokio::spawn(async move {
106            let mut request = request;
107            request.ai_tool_call_id = request.context.ai_tool_call_id().cloned();
108            if let Err(e) = repo
109                .log_execution_sync_with_id(
110                    &mcp_execution_id,
111                    &request,
112                    &result_row,
113                    Correlation::Exact,
114                )
115                .await
116            {
117                tracing::warn!(
118                    tool = %request.tool_name,
119                    server = %request.server_name,
120                    error = %e,
121                    "Failed to record external MCP tool execution"
122                );
123                return;
124            }
125            if let Some(ingest) = ingest {
126                ingest_proxied_result(&ingest, &request, result, mcp_execution_id).await;
127            }
128        });
129    }
130}
131
132async fn ingest_proxied_result(
133    ingest: &ArtifactIngest,
134    request: &ToolExecutionRequest,
135    result: Option<Value>,
136    mcp_execution_id: McpExecutionId,
137) {
138    let Some(wire) = result.as_ref().and_then(from_wire_value) else {
139        return;
140    };
141    let ingest_request = IngestRequest {
142        result: wire,
143        tool_name: request.tool_name.clone(),
144        server_name: Some(request.server_name.clone()),
145        ai_tool_call_id: request.ai_tool_call_id.clone(),
146        mcp_execution_id: Some(mcp_execution_id),
147        ctx: request.context.clone(),
148        skill: None,
149        source: ExecutionSource::Proxy,
150        started_at: Some(request.started_at),
151        input: Some(request.input.clone()),
152    };
153    if let Err(e) = ingest.ingest(ingest_request).await {
154        tracing::warn!(
155            tool = %request.tool_name,
156            server = %request.server_name,
157            error = %e,
158            "Failed to ingest external MCP tool result as an artifact"
159        );
160    }
161}