Skip to main content

systemprompt_runtime/trace/repository/
log_lookup.rs

1//! Log-row lookup by trace/context with typed row mapping.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}