systemprompt_api/services/proxy/audit/
mod.rs1pub mod jsonrpc;
22pub mod tap;
23
24use std::sync::Arc;
25
26use chrono::{DateTime, Utc};
27use serde_json::Value;
28use systemprompt_identifiers::McpExecutionId;
29use systemprompt_mcp::models::{ExecutionStatus, ToolExecutionRequest, ToolExecutionResult};
30use systemprompt_mcp::repository::ToolUsageRepository;
31use systemprompt_mcp::{
32 ArtifactIngest, INTENT_CLAIM_WINDOW_SECONDS, IngestRequest, from_wire_value,
33};
34use systemprompt_models::RequestContext;
35use systemprompt_models::mcp::{Correlation, ExecutionSource};
36
37pub(crate) use jsonrpc::parse_tool_call;
38pub(crate) use tap::record;
39
40use jsonrpc::{ToolCallInvocation, ToolCallOutcome};
41
42#[derive(Debug)]
43pub struct McpAudit {
44 repo: Arc<ToolUsageRepository>,
45 ingest: Option<Arc<ArtifactIngest>>,
46 context: RequestContext,
47 server_name: String,
48 invocation: ToolCallInvocation,
49 started_at: DateTime<Utc>,
50 mcp_execution_id: McpExecutionId,
51}
52
53impl McpAudit {
54 pub fn new(
55 repo: Arc<ToolUsageRepository>,
56 ingest: Option<Arc<ArtifactIngest>>,
57 context: RequestContext,
58 server_name: String,
59 invocation: ToolCallInvocation,
60 ) -> Self {
61 Self {
62 repo,
63 ingest,
64 context,
65 server_name,
66 invocation,
67 started_at: Utc::now(),
68 mcp_execution_id: McpExecutionId::new(uuid::Uuid::new_v4().to_string()),
69 }
70 }
71
72 const fn request_id(&self) -> &Value {
73 &self.invocation.id
74 }
75
76 pub const fn mcp_execution_id(&self) -> &McpExecutionId {
77 &self.mcp_execution_id
78 }
79
80 fn finalize(self, outcome: Option<ToolCallOutcome>) {
81 let (output, error_message, result) = match outcome {
82 Some(o) => (o.output, o.error_message, o.result),
83 None => (
84 None,
85 Some("external MCP tool call produced no parseable result".to_owned()),
86 None,
87 ),
88 };
89
90 let request = ToolExecutionRequest {
91 tool_name: self.invocation.tool_name,
92 server_name: self.server_name.clone(),
93 input: self.invocation.arguments,
94 started_at: self.started_at,
95 context: self.context,
96 request_method: Some("mcp".to_owned()),
97 request_source: Some(self.server_name),
98 ai_tool_call_id: None,
99 source: ExecutionSource::Proxy,
100 };
101 let result_row = ToolExecutionResult {
102 status: ExecutionStatus::from_error(error_message.is_some()).to_string(),
103 error_message,
104 output,
105 output_schema: None,
106 started_at: self.started_at,
107 completed_at: Some(Utc::now()),
108 };
109
110 let repo = self.repo;
111 let ingest = self.ingest;
112 let mcp_execution_id = self.mcp_execution_id;
113 tokio::spawn(async move {
114 let mut request = request;
115 request.ai_tool_call_id = request.context.ai_tool_call_id().cloned();
116 let correlation = if request.ai_tool_call_id.is_some() {
123 Correlation::Exact
124 } else {
125 claim_intent(&repo, &mut request, &mcp_execution_id).await
126 };
127 if let Err(e) = repo
128 .log_execution_sync_with_id(&mcp_execution_id, &request, &result_row, correlation)
129 .await
130 {
131 tracing::warn!(
132 tool = %request.tool_name,
133 server = %request.server_name,
134 error = %e,
135 "Failed to record external MCP tool execution"
136 );
137 return;
138 }
139 if let Some(ingest) = ingest {
140 ingest_proxied_result(&ingest, &request, result, mcp_execution_id).await;
141 }
142 });
143 }
144}
145
146async fn claim_intent(
147 repo: &ToolUsageRepository,
148 request: &mut ToolExecutionRequest,
149 mcp_execution_id: &McpExecutionId,
150) -> Correlation {
151 match repo
152 .claim_unclaimed_intent(
153 request.context.session_id(),
154 &request.tool_name,
155 mcp_execution_id,
156 INTENT_CLAIM_WINDOW_SECONDS,
157 )
158 .await
159 {
160 Ok(Some(call_id)) => {
161 request.ai_tool_call_id = Some(call_id);
162 Correlation::Inferred
163 },
164 Ok(None) => Correlation::Inferred,
165 Err(e) => {
166 tracing::warn!(
167 tool = %request.tool_name,
168 server = %request.server_name,
169 %mcp_execution_id,
170 error = %e,
171 "Proxy intent claim failed"
172 );
173 Correlation::Inferred
174 },
175 }
176}
177
178async fn ingest_proxied_result(
179 ingest: &ArtifactIngest,
180 request: &ToolExecutionRequest,
181 result: Option<Value>,
182 mcp_execution_id: McpExecutionId,
183) {
184 let Some(wire) = result.as_ref().and_then(from_wire_value) else {
185 return;
186 };
187 let ingest_request = IngestRequest {
188 result: wire,
189 tool_name: request.tool_name.clone(),
190 server_name: Some(request.server_name.clone()),
191 ai_tool_call_id: request.ai_tool_call_id.clone(),
192 mcp_execution_id: Some(mcp_execution_id),
193 ctx: request.context.clone(),
194 skill: None,
195 source: ExecutionSource::Proxy,
196 started_at: Some(request.started_at),
197 input: Some(request.input.clone()),
198 };
199 if let Err(e) = ingest.ingest(ingest_request).await {
200 tracing::warn!(
201 tool = %request.tool_name,
202 server = %request.server_name,
203 error = %e,
204 "Failed to ingest external MCP tool result as an artifact"
205 );
206 }
207}