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//!
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 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}