Skip to main content

systemprompt_runtime/trace/
service.rs

1//! Read-side query facade over the tracing and audit tables.
2//!
3//! [`TraceQueryService`] is the single entry point for reconstructing a trace
4//! from its constituent rows (logs, AI requests, MCP tool executions, task
5//! execution steps) and for the log/audit browsing surfaces. Each public method
6//! delegates to a focused query module in this directory;
7//! [`get_all_trace_data`] fans the per-source fetches out concurrently.
8//!
9//! [`get_all_trace_data`]: TraceQueryService::get_all_trace_data
10//!
11//! Copyright (c) systemprompt.io — Business Source License 1.1.
12//! See <https://systemprompt.io> for licensing details.
13
14use chrono::{DateTime, Utc};
15use sqlx::PgPool;
16use std::sync::Arc;
17use systemprompt_identifiers::{AiRequestId, TaskId, TraceId};
18
19use systemprompt_logging::models::LogEntry;
20
21use super::TraceError;
22
23pub(super) type Result<T> = std::result::Result<T, TraceError>;
24
25use super::models::{
26    AiRequestDetail, AiRequestFilter, AiRequestListItem, AiRequestStats, AiRequestSummary,
27    AuditLookupResult, AuditPage, AuditToolCallRow, ConversationMessage, ExecutionStepSummary,
28    LevelCount, LinkedMcpCall, LogSearchItem, LogTimeRange, McpExecutionSummary, ModuleCount,
29    ToolExecutionFilter, ToolExecutionItem, TraceEvent, TraceListFilter, TraceListItem,
30};
31use super::{
32    audit_queries, list_queries, log_lookup_queries, log_search_queries, log_summary_queries,
33    queries, request_queries, request_stats_queries, tool_queries,
34};
35
36#[derive(Debug, Clone)]
37pub struct TraceQueryService {
38    pool: Arc<PgPool>,
39}
40
41impl TraceQueryService {
42    pub const fn new(pool: Arc<PgPool>) -> Self {
43        Self { pool }
44    }
45
46    pub async fn get_log_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
47        queries::fetch_log_events(&self.pool, trace_id).await
48    }
49
50    pub async fn get_ai_request_summary(&self, trace_id: &TraceId) -> Result<AiRequestSummary> {
51        queries::fetch_ai_request_summary(&self.pool, trace_id).await
52    }
53
54    pub async fn get_ai_request_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
55        queries::fetch_ai_request_events(&self.pool, trace_id).await
56    }
57
58    pub async fn get_mcp_execution_summary(
59        &self,
60        trace_id: &TraceId,
61    ) -> Result<McpExecutionSummary> {
62        queries::fetch_mcp_execution_summary(&self.pool, trace_id).await
63    }
64
65    pub async fn get_mcp_execution_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
66        queries::fetch_mcp_execution_events(&self.pool, trace_id).await
67    }
68
69    pub async fn get_task_id(&self, trace_id: &TraceId) -> Result<Option<TaskId>> {
70        Ok(queries::fetch_task_id_for_trace(&self.pool, trace_id)
71            .await?
72            .map(TaskId::new))
73    }
74
75    pub async fn get_execution_step_summary(
76        &self,
77        trace_id: &TraceId,
78    ) -> Result<ExecutionStepSummary> {
79        queries::fetch_execution_step_summary(&self.pool, trace_id).await
80    }
81
82    pub async fn get_execution_step_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
83        queries::fetch_execution_step_events(&self.pool, trace_id).await
84    }
85
86    pub async fn get_all_trace_data(
87        &self,
88        trace_id: &TraceId,
89    ) -> Result<(
90        Vec<TraceEvent>,
91        Vec<TraceEvent>,
92        Vec<TraceEvent>,
93        Vec<TraceEvent>,
94        AiRequestSummary,
95        McpExecutionSummary,
96        ExecutionStepSummary,
97        Option<TaskId>,
98    )> {
99        tokio::try_join!(
100            self.get_log_events(trace_id),
101            self.get_ai_request_events(trace_id),
102            self.get_mcp_execution_events(trace_id),
103            self.get_execution_step_events(trace_id),
104            self.get_ai_request_summary(trace_id),
105            self.get_mcp_execution_summary(trace_id),
106            self.get_execution_step_summary(trace_id),
107            self.get_task_id(trace_id),
108        )
109    }
110
111    pub async fn list_traces(&self, filter: &TraceListFilter) -> Result<Vec<TraceListItem>> {
112        list_queries::list_traces(&self.pool, filter).await
113    }
114
115    pub async fn list_tool_executions(
116        &self,
117        filter: &ToolExecutionFilter,
118    ) -> Result<Vec<ToolExecutionItem>> {
119        tool_queries::list_tool_executions(&self.pool, filter).await
120    }
121
122    pub async fn search_logs(
123        &self,
124        pattern: &str,
125        since: Option<DateTime<Utc>>,
126        level: Option<&str>,
127        limit: i64,
128    ) -> Result<Vec<LogSearchItem>> {
129        log_search_queries::search_logs(&self.pool, pattern, since, level, limit).await
130    }
131
132    pub async fn search_tool_executions(
133        &self,
134        pattern: &str,
135        since: Option<DateTime<Utc>>,
136        limit: i64,
137    ) -> Result<Vec<ToolExecutionItem>> {
138        log_search_queries::search_tool_executions(&self.pool, pattern, since, limit).await
139    }
140
141    pub async fn list_ai_requests(
142        &self,
143        filter: &AiRequestFilter,
144    ) -> Result<Vec<AiRequestListItem>> {
145        request_queries::list_ai_requests(&self.pool, filter).await
146    }
147
148    pub async fn get_ai_request_stats(
149        &self,
150        since: Option<DateTime<Utc>>,
151    ) -> Result<AiRequestStats> {
152        request_stats_queries::get_ai_request_stats(&self.pool, since).await
153    }
154
155    pub async fn find_ai_request_detail(&self, id: &str) -> Result<Option<AiRequestDetail>> {
156        request_queries::find_ai_request_detail(&self.pool, id).await
157    }
158
159    pub async fn find_ai_request_for_audit(&self, id: &str) -> Result<Option<AuditLookupResult>> {
160        audit_queries::find_ai_request_for_audit(&self.pool, id).await
161    }
162
163    pub async fn count_audit_messages(&self, request_id: &AiRequestId) -> Result<i64> {
164        audit_queries::count_audit_messages(&self.pool, request_id).await
165    }
166
167    pub async fn count_audit_tool_calls(&self, request_id: &AiRequestId) -> Result<i64> {
168        audit_queries::count_audit_tool_calls(&self.pool, request_id).await
169    }
170
171    pub async fn list_audit_messages(
172        &self,
173        request_id: &AiRequestId,
174        page: AuditPage,
175    ) -> Result<Vec<ConversationMessage>> {
176        audit_queries::list_audit_messages(&self.pool, request_id, page).await
177    }
178
179    pub async fn list_audit_tool_calls(
180        &self,
181        request_id: &AiRequestId,
182        page: AuditPage,
183    ) -> Result<Vec<AuditToolCallRow>> {
184        audit_queries::list_audit_tool_calls(&self.pool, request_id, page).await
185    }
186
187    pub async fn list_linked_mcp_calls(
188        &self,
189        request_id: &AiRequestId,
190    ) -> Result<Vec<LinkedMcpCall>> {
191        audit_queries::list_linked_mcp_calls(&self.pool, request_id).await
192    }
193
194    pub async fn find_log_by_id(&self, id: &str) -> Result<Option<LogEntry>> {
195        log_lookup_queries::find_log_by_id(&self.pool, id).await
196    }
197
198    pub async fn find_log_by_partial_id(&self, id_prefix: &str) -> Result<Option<LogEntry>> {
199        log_lookup_queries::find_log_by_partial_id(&self.pool, id_prefix).await
200    }
201
202    pub async fn find_logs_by_trace_id(&self, trace_id: &TraceId) -> Result<Vec<LogEntry>> {
203        log_lookup_queries::find_logs_by_trace_id(&self.pool, trace_id).await
204    }
205
206    pub async fn list_logs_filtered(
207        &self,
208        since: Option<DateTime<Utc>>,
209        level: Option<&str>,
210        limit: i64,
211    ) -> Result<Vec<LogEntry>> {
212        log_lookup_queries::list_logs_filtered(&self.pool, since, level, limit).await
213    }
214
215    pub async fn count_logs_by_level(
216        &self,
217        since: Option<DateTime<Utc>>,
218    ) -> Result<Vec<LevelCount>> {
219        log_summary_queries::count_logs_by_level(&self.pool, since).await
220    }
221
222    pub async fn top_modules(
223        &self,
224        since: Option<DateTime<Utc>>,
225        limit: i64,
226    ) -> Result<Vec<ModuleCount>> {
227        log_summary_queries::top_modules(&self.pool, since, limit).await
228    }
229
230    pub async fn log_time_range(&self, since: Option<DateTime<Utc>>) -> Result<LogTimeRange> {
231        log_summary_queries::log_time_range(&self.pool, since).await
232    }
233
234    pub async fn total_log_count(&self) -> Result<i64> {
235        log_summary_queries::total_log_count(&self.pool).await
236    }
237}