Skip to main content

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