systemprompt_runtime/trace/repository/
log_lookup.rs1use chrono::{DateTime, Utc};
7use systemprompt_identifiers::{
8 ClientId, ContextId, InstanceId, LogId, SessionId, TaskId, TraceId, UserId,
9};
10
11use super::{Result, TraceRepository};
12use systemprompt_logging::models::{LogEntry, LogLevel};
13
14struct LogRow {
15 id: LogId,
16 timestamp: DateTime<Utc>,
17 level: String,
18 module: String,
19 message: String,
20 metadata: Option<String>,
21 user_id: UserId,
22 session_id: SessionId,
23 task_id: Option<TaskId>,
24 trace_id: TraceId,
25 context_id_text: Option<String>,
26 client_id: Option<ClientId>,
27 instance_id: Option<InstanceId>,
28}
29
30fn row_to_entry(r: LogRow) -> LogEntry {
31 LogEntry {
32 id: r.id,
33 timestamp: r.timestamp,
34 level: r.level.parse().unwrap_or(LogLevel::Info),
35 module: r.module,
36 message: r.message,
37 metadata: r.metadata.as_ref().and_then(|m| {
38 serde_json::from_str(m)
39 .map_err(|e| {
40 tracing::warn!(error = %e, raw = %m, "Failed to parse log metadata JSON");
41 e
42 })
43 .ok()
44 }),
45 user_id: r.user_id,
46 session_id: r.session_id,
47 task_id: r.task_id,
48 trace_id: r.trace_id,
49 context_id: r.context_id_text.and_then(|s| {
50 ContextId::try_new(&s)
51 .map_err(|e| {
52 tracing::warn!(error = %e, raw = %s, "Skipping non-UUID context_id from log row");
53 e
54 })
55 .ok()
56 }),
57 client_id: r.client_id,
58 instance_id: r.instance_id,
59 }
60}
61
62impl TraceRepository {
63 pub async fn find_log_by_id(&self, id: &LogId) -> Result<Option<LogEntry>> {
64 let row = sqlx::query_as!(
65 LogRow,
66 r#"
67 SELECT
68 id as "id!: LogId", timestamp as "timestamp!", level as "level!", module as "module!",
69 message as "message!", metadata,
70 user_id as "user_id!: UserId",
71 session_id as "session_id!: SessionId",
72 task_id as "task_id: TaskId",
73 trace_id as "trace_id!: TraceId",
74 context_id as "context_id_text",
75 client_id as "client_id: ClientId",
76 instance_id as "instance_id: InstanceId"
77 FROM logs WHERE id = $1
78 "#,
79 id.as_str()
80 )
81 .fetch_optional(&*self.pool)
82 .await?;
83
84 Ok(row.map(row_to_entry))
85 }
86
87 pub async fn find_log_by_partial_id(&self, id_prefix: &str) -> Result<Option<LogEntry>> {
88 let pattern = format!("{id_prefix}%");
89 let row = sqlx::query_as!(
90 LogRow,
91 r#"
92 SELECT
93 id as "id!: LogId", timestamp as "timestamp!", level as "level!", module as "module!",
94 message as "message!", metadata,
95 user_id as "user_id!: UserId",
96 session_id as "session_id!: SessionId",
97 task_id as "task_id: TaskId",
98 trace_id as "trace_id!: TraceId",
99 context_id as "context_id_text",
100 client_id as "client_id: ClientId",
101 instance_id as "instance_id: InstanceId"
102 FROM logs
103 WHERE id LIKE $1
104 ORDER BY timestamp DESC
105 LIMIT 1
106 "#,
107 pattern
108 )
109 .fetch_optional(&*self.pool)
110 .await?;
111
112 Ok(row.map(row_to_entry))
113 }
114
115 pub async fn find_logs_by_trace_id(&self, trace_id: &TraceId) -> Result<Vec<LogEntry>> {
116 let rows = sqlx::query_as!(
117 LogRow,
118 r#"
119 SELECT
120 id as "id!: LogId", timestamp as "timestamp!", level as "level!", module as "module!",
121 message as "message!", metadata,
122 user_id as "user_id!: UserId",
123 session_id as "session_id!: SessionId",
124 task_id as "task_id: TaskId",
125 trace_id as "trace_id!: TraceId",
126 context_id as "context_id_text",
127 client_id as "client_id: ClientId",
128 instance_id as "instance_id: InstanceId"
129 FROM logs
130 WHERE trace_id = $1
131 ORDER BY timestamp ASC
132 "#,
133 trace_id.as_str()
134 )
135 .fetch_all(&*self.pool)
136 .await?;
137
138 if !rows.is_empty() {
139 return Ok(rows.into_iter().map(row_to_entry).collect());
140 }
141
142 let pattern = format!("{}%", trace_id.as_str());
143 let rows = sqlx::query_as!(
144 LogRow,
145 r#"
146 SELECT
147 id as "id!: LogId", timestamp as "timestamp!", level as "level!", module as "module!",
148 message as "message!", metadata,
149 user_id as "user_id!: UserId",
150 session_id as "session_id!: SessionId",
151 task_id as "task_id: TaskId",
152 trace_id as "trace_id!: TraceId",
153 context_id as "context_id_text",
154 client_id as "client_id: ClientId",
155 instance_id as "instance_id: InstanceId"
156 FROM logs
157 WHERE trace_id LIKE $1
158 ORDER BY timestamp ASC
159 LIMIT 100
160 "#,
161 pattern
162 )
163 .fetch_all(&*self.pool)
164 .await?;
165
166 Ok(rows.into_iter().map(row_to_entry).collect())
167 }
168
169 pub async fn list_logs_filtered(
170 &self,
171 since: Option<DateTime<Utc>>,
172 level: Option<LogLevel>,
173 limit: i64,
174 ) -> Result<Vec<LogEntry>> {
175 let rows = sqlx::query_as!(
176 LogRow,
177 r#"
178 SELECT
179 id as "id!: LogId", timestamp as "timestamp!", level as "level!", module as "module!",
180 message as "message!", metadata,
181 user_id as "user_id!: UserId",
182 session_id as "session_id!: SessionId",
183 task_id as "task_id: TaskId",
184 trace_id as "trace_id!: TraceId",
185 context_id as "context_id_text",
186 client_id as "client_id: ClientId",
187 instance_id as "instance_id: InstanceId"
188 FROM logs
189 WHERE ($1::TIMESTAMPTZ IS NULL OR timestamp >= $1)
190 AND ($2::TEXT IS NULL OR UPPER(level) = $2)
191 ORDER BY timestamp DESC
192 LIMIT $3
193 "#,
194 since,
195 level.map(LogLevel::as_str),
196 limit
197 )
198 .fetch_all(&*self.pool)
199 .await?;
200
201 Ok(rows.into_iter().map(row_to_entry).collect())
202 }
203}