systemprompt_runtime/trace/
service.rs1use 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}