systemprompt_logging/trace/
service.rs1use chrono::{DateTime, Utc};
15use sqlx::PgPool;
16use std::sync::Arc;
17use systemprompt_identifiers::{AiRequestId, TaskId, TraceId};
18
19use crate::models::{LogEntry, LoggingError};
20
21pub(super) type Result<T> = std::result::Result<T, LoggingError>;
22
23use super::models::{
24 AiRequestDetail, AiRequestFilter, AiRequestListItem, AiRequestStats, AiRequestSummary,
25 AuditLookupResult, AuditToolCallRow, ConversationMessage, ExecutionStepSummary, LevelCount,
26 LinkedMcpCall, LogSearchItem, LogTimeRange, McpExecutionSummary, ModuleCount,
27 ToolExecutionFilter, ToolExecutionItem, TraceEvent, TraceListFilter, TraceListItem,
28};
29use super::{
30 audit_queries, list_queries, log_lookup_queries, log_search_queries, log_summary_queries,
31 queries, request_queries, tool_queries,
32};
33
34#[derive(Debug, Clone)]
35pub struct TraceQueryService {
36 pool: Arc<PgPool>,
37}
38
39impl TraceQueryService {
40 pub const fn new(pool: Arc<PgPool>) -> Self {
41 Self { pool }
42 }
43
44 pub async fn get_log_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
45 queries::fetch_log_events(&self.pool, trace_id).await
46 }
47
48 pub async fn get_ai_request_summary(&self, trace_id: &TraceId) -> Result<AiRequestSummary> {
49 queries::fetch_ai_request_summary(&self.pool, trace_id).await
50 }
51
52 pub async fn get_ai_request_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
53 queries::fetch_ai_request_events(&self.pool, trace_id).await
54 }
55
56 pub async fn get_mcp_execution_summary(
57 &self,
58 trace_id: &TraceId,
59 ) -> Result<McpExecutionSummary> {
60 queries::fetch_mcp_execution_summary(&self.pool, trace_id).await
61 }
62
63 pub async fn get_mcp_execution_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
64 queries::fetch_mcp_execution_events(&self.pool, trace_id).await
65 }
66
67 pub async fn get_task_id(&self, trace_id: &TraceId) -> Result<Option<TaskId>> {
68 Ok(queries::fetch_task_id_for_trace(&self.pool, trace_id)
69 .await?
70 .map(TaskId::new))
71 }
72
73 pub async fn get_execution_step_summary(
74 &self,
75 trace_id: &TraceId,
76 ) -> Result<ExecutionStepSummary> {
77 queries::fetch_execution_step_summary(&self.pool, trace_id).await
78 }
79
80 pub async fn get_execution_step_events(&self, trace_id: &TraceId) -> Result<Vec<TraceEvent>> {
81 queries::fetch_execution_step_events(&self.pool, trace_id).await
82 }
83
84 pub async fn get_all_trace_data(
85 &self,
86 trace_id: &TraceId,
87 ) -> Result<(
88 Vec<TraceEvent>,
89 Vec<TraceEvent>,
90 Vec<TraceEvent>,
91 Vec<TraceEvent>,
92 AiRequestSummary,
93 McpExecutionSummary,
94 ExecutionStepSummary,
95 Option<TaskId>,
96 )> {
97 tokio::try_join!(
98 self.get_log_events(trace_id),
99 self.get_ai_request_events(trace_id),
100 self.get_mcp_execution_events(trace_id),
101 self.get_execution_step_events(trace_id),
102 self.get_ai_request_summary(trace_id),
103 self.get_mcp_execution_summary(trace_id),
104 self.get_execution_step_summary(trace_id),
105 self.get_task_id(trace_id),
106 )
107 }
108
109 pub async fn list_traces(&self, filter: &TraceListFilter) -> Result<Vec<TraceListItem>> {
110 list_queries::list_traces(&self.pool, filter).await
111 }
112
113 pub async fn list_tool_executions(
114 &self,
115 filter: &ToolExecutionFilter,
116 ) -> Result<Vec<ToolExecutionItem>> {
117 tool_queries::list_tool_executions(&self.pool, filter).await
118 }
119
120 pub async fn search_logs(
121 &self,
122 pattern: &str,
123 since: Option<DateTime<Utc>>,
124 level: Option<&str>,
125 limit: i64,
126 ) -> Result<Vec<LogSearchItem>> {
127 log_search_queries::search_logs(&self.pool, pattern, since, level, limit).await
128 }
129
130 pub async fn search_tool_executions(
131 &self,
132 pattern: &str,
133 since: Option<DateTime<Utc>>,
134 limit: i64,
135 ) -> Result<Vec<ToolExecutionItem>> {
136 log_search_queries::search_tool_executions(&self.pool, pattern, since, limit).await
137 }
138
139 pub async fn list_ai_requests(
140 &self,
141 filter: &AiRequestFilter,
142 ) -> Result<Vec<AiRequestListItem>> {
143 request_queries::list_ai_requests(&self.pool, filter).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}