systemprompt_logging/trace/
service.rs1use chrono::{DateTime, Utc};
12use sqlx::PgPool;
13use std::sync::Arc;
14use systemprompt_identifiers::{AiRequestId, TaskId, TraceId};
15
16use crate::models::{LogEntry, LoggingError};
17
18pub(super) type Result<T> = std::result::Result<T, LoggingError>;
19
20use super::models::{
21 AiRequestDetail, AiRequestListItem, AiRequestStats, AiRequestSummary, AuditLookupResult,
22 AuditToolCallRow, ConversationMessage, ExecutionStepSummary, LevelCount, LinkedMcpCall,
23 LogSearchItem, LogTimeRange, McpExecutionSummary, ModuleCount, ToolExecutionFilter,
24 ToolExecutionItem, TraceEvent, TraceListFilter, TraceListItem,
25};
26use super::{
27 audit_queries, list_queries, log_lookup_queries, log_search_queries, log_summary_queries,
28 queries, request_queries, tool_queries,
29};
30
31#[derive(Debug, Clone)]
32pub struct TraceQueryService {
33 pool: Arc<PgPool>,
34}
35
36impl TraceQueryService {
37 pub const fn new(pool: Arc<PgPool>) -> Self {
38 Self { pool }
39 }
40
41 pub async fn get_log_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
42 queries::fetch_log_events(&self.pool, trace_id).await
43 }
44
45 pub async fn get_ai_request_summary(&self, trace_id: &TraceId) -> Result<AiRequestSummary> {
46 queries::fetch_ai_request_summary(&self.pool, trace_id).await
47 }
48
49 pub async fn get_ai_request_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
50 queries::fetch_ai_request_events(&self.pool, trace_id).await
51 }
52
53 pub async fn get_mcp_execution_summary(
54 &self,
55 trace_id: &TraceId,
56 ) -> Result<McpExecutionSummary> {
57 queries::fetch_mcp_execution_summary(&self.pool, trace_id).await
58 }
59
60 pub async fn get_mcp_execution_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
61 queries::fetch_mcp_execution_events(&self.pool, trace_id).await
62 }
63
64 pub async fn get_task_id(&self, trace_id: &TraceId) -> Result<Option<TaskId>> {
65 Ok(queries::fetch_task_id_for_trace(&self.pool, trace_id)
66 .await?
67 .map(TaskId::new))
68 }
69
70 pub async fn get_execution_step_summary(
71 &self,
72 trace_id: &TraceId,
73 ) -> Result<ExecutionStepSummary> {
74 queries::fetch_execution_step_summary(&self.pool, trace_id).await
75 }
76
77 pub async fn get_execution_step_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
78 queries::fetch_execution_step_events(&self.pool, trace_id).await
79 }
80
81 pub async fn get_all_trace_data(
82 &self,
83 trace_id: &TraceId,
84 ) -> Result<(
85 Vec<TraceEvent>,
86 Vec<TraceEvent>,
87 Vec<TraceEvent>,
88 Vec<TraceEvent>,
89 AiRequestSummary,
90 McpExecutionSummary,
91 ExecutionStepSummary,
92 Option<TaskId>,
93 )> {
94 tokio::try_join!(
95 self.get_log_events(trace_id),
96 self.get_ai_request_events(trace_id),
97 self.get_mcp_execution_events(trace_id),
98 self.get_execution_step_events(trace_id),
99 self.get_ai_request_summary(trace_id),
100 self.get_mcp_execution_summary(trace_id),
101 self.get_execution_step_summary(trace_id),
102 self.get_task_id(trace_id),
103 )
104 }
105
106 pub async fn list_traces(&self, filter: &TraceListFilter) -> Result<Vec<TraceListItem>> {
107 list_queries::list_traces(&self.pool, filter).await
108 }
109
110 pub async fn list_tool_executions(
111 &self,
112 filter: &ToolExecutionFilter,
113 ) -> Result<Vec<ToolExecutionItem>> {
114 tool_queries::list_tool_executions(&self.pool, filter).await
115 }
116
117 pub async fn search_logs(
118 &self,
119 pattern: &str,
120 since: Option<DateTime<Utc>>,
121 level: Option<&str>,
122 limit: i64,
123 ) -> Result<Vec<LogSearchItem>> {
124 log_search_queries::search_logs(&self.pool, pattern, since, level, limit).await
125 }
126
127 pub async fn search_tool_executions(
128 &self,
129 pattern: &str,
130 since: Option<DateTime<Utc>>,
131 limit: i64,
132 ) -> Result<Vec<ToolExecutionItem>> {
133 log_search_queries::search_tool_executions(&self.pool, pattern, since, limit).await
134 }
135
136 pub async fn list_ai_requests(
137 &self,
138 since: Option<DateTime<Utc>>,
139 model: Option<&str>,
140 provider: Option<&str>,
141 limit: i64,
142 ) -> Result<Vec<AiRequestListItem>> {
143 request_queries::list_ai_requests(&self.pool, since, model, provider, limit).await
144 }
145
146 pub async fn get_ai_request_stats(
147 &self,
148 since: Option<DateTime<Utc>>,
149 ) -> Result<AiRequestStats> {
150 request_queries::get_ai_request_stats(&self.pool, since).await
151 }
152
153 pub async fn find_ai_request_detail(&self, id: &str) -> Result<Option<AiRequestDetail>> {
154 request_queries::find_ai_request_detail(&self.pool, id).await
155 }
156
157 pub async fn find_ai_request_for_audit(&self, id: &str) -> Result<Option<AuditLookupResult>> {
158 audit_queries::find_ai_request_for_audit(&self.pool, id).await
159 }
160
161 pub async fn list_audit_messages(
162 &self,
163 request_id: &AiRequestId,
164 ) -> Result<Vec<ConversationMessage>> {
165 audit_queries::list_audit_messages(&self.pool, request_id).await
166 }
167
168 pub async fn list_audit_tool_calls(
169 &self,
170 request_id: &AiRequestId,
171 ) -> Result<Vec<AuditToolCallRow>> {
172 audit_queries::list_audit_tool_calls(&self.pool, request_id).await
173 }
174
175 pub async fn list_linked_mcp_calls(
176 &self,
177 request_id: &AiRequestId,
178 ) -> Result<Vec<LinkedMcpCall>> {
179 audit_queries::list_linked_mcp_calls(&self.pool, request_id).await
180 }
181
182 pub async fn find_log_by_id(&self, id: &str) -> Result<Option<LogEntry>> {
183 log_lookup_queries::find_log_by_id(&self.pool, id).await
184 }
185
186 pub async fn find_log_by_partial_id(&self, id_prefix: &str) -> Result<Option<LogEntry>> {
187 log_lookup_queries::find_log_by_partial_id(&self.pool, id_prefix).await
188 }
189
190 pub async fn find_logs_by_trace_id(&self, trace_id: &TraceId) -> Result<Vec<LogEntry>> {
191 log_lookup_queries::find_logs_by_trace_id(&self.pool, trace_id).await
192 }
193
194 pub async fn list_logs_filtered(
195 &self,
196 since: Option<DateTime<Utc>>,
197 level: Option<&str>,
198 limit: i64,
199 ) -> Result<Vec<LogEntry>> {
200 log_lookup_queries::list_logs_filtered(&self.pool, since, level, limit).await
201 }
202
203 pub async fn count_logs_by_level(
204 &self,
205 since: Option<DateTime<Utc>>,
206 ) -> Result<Vec<LevelCount>> {
207 log_summary_queries::count_logs_by_level(&self.pool, since).await
208 }
209
210 pub async fn top_modules(
211 &self,
212 since: Option<DateTime<Utc>>,
213 limit: i64,
214 ) -> Result<Vec<ModuleCount>> {
215 log_summary_queries::top_modules(&self.pool, since, limit).await
216 }
217
218 pub async fn log_time_range(&self, since: Option<DateTime<Utc>>) -> Result<LogTimeRange> {
219 log_summary_queries::log_time_range(&self.pool, since).await
220 }
221
222 pub async fn total_log_count(&self) -> Result<i64> {
223 log_summary_queries::total_log_count(&self.pool).await
224 }
225}