systemprompt_api/services/proxy/audit/
mod.rs1pub 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 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
219async 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}