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