Skip to main content

systemprompt_logging/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
11use 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}