Skip to main content

headless_lms_models/
errors.rs

1use crate::prelude::*;
2use headless_lms_utils::error_identifier::{
3    calculate_error_grouping_identifier, calculate_exact_error_identifier, normalize_message,
4    normalize_stack_trace,
5};
6use rand::RngExt;
7use utoipa::ToSchema;
8
9#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, Copy, Type, ToSchema)]
10#[serde(rename_all = "snake_case")]
11#[sqlx(type_name = "error_source", rename_all = "snake_case")]
12pub enum ErrorSource {
13    Backend,
14    Frontend,
15}
16
17#[derive(Debug, Serialize, Deserialize, Clone, ToSchema)]
18pub struct ErrorVariant {
19    pub id: Uuid,
20    pub service: String,
21    pub exact_error_identifier: String,
22    pub error_grouping_identifier: String,
23    pub error_source: ErrorSource,
24    pub example_message: String,
25    pub example_stack_trace: Option<String>,
26    pub normalized_message: String,
27    pub normalized_stack_trace: Option<String>,
28    pub occurrence_count: i32,
29    pub last_seen_at: DateTime<Utc>,
30    pub resolved_at: Option<DateTime<Utc>>,
31    pub created_at: DateTime<Utc>,
32    pub updated_at: DateTime<Utc>,
33    pub deleted_at: Option<DateTime<Utc>>,
34}
35
36#[derive(Debug, Serialize, Deserialize, Clone, ToSchema)]
37pub struct NewErrorReport {
38    pub service: String,
39    pub error_source: Option<ErrorSource>,
40    pub message: String,
41    pub stack_trace: Option<String>,
42    pub path: Option<String>,
43    pub app_version: Option<String>,
44    pub details: Option<serde_json::Value>,
45}
46
47/// Inserts one error occurrence and upserts its variant in a single transaction.
48/// The variant service is treated as immutable after creation; any future service
49/// rename must update both error_variants.service and related error_occurrences.service
50/// in the same transaction.
51pub async fn insert(
52    conn: &mut PgConnection,
53    user_id: Option<Uuid>,
54    report: &NewErrorReport,
55) -> ModelResult<Uuid> {
56    let service = report.service.trim();
57    if service.is_empty() {
58        return Err(model_err!(
59            InvalidRequest,
60            "service must not be empty".to_string()
61        ));
62    }
63
64    let error_source = report.error_source.unwrap_or(ErrorSource::Frontend);
65    let error_source_str = match error_source {
66        ErrorSource::Frontend => "frontend",
67        ErrorSource::Backend => "backend",
68    };
69
70    let normalized_message = normalize_message(&report.message);
71    let normalized_stack_trace = report.stack_trace.as_deref().map(normalize_stack_trace);
72
73    let exact_error_identifier = calculate_exact_error_identifier(
74        service,
75        error_source_str,
76        &normalized_message,
77        normalized_stack_trace.as_deref(),
78    );
79    let error_grouping_identifier =
80        calculate_error_grouping_identifier(service, error_source_str, &normalized_message);
81
82    let mut tx = conn.begin().await?;
83    let variant_id = sqlx::query!(
84        r#"
85INSERT INTO error_variants (
86    id,
87    service,
88    exact_error_identifier,
89    error_grouping_identifier,
90    error_source,
91    example_message,
92    example_stack_trace,
93    normalized_message,
94    normalized_stack_trace
95)
96VALUES (uuid_generate_v4(), $1, $2, $3, $4, $5, $6, $7, $8)
97ON CONFLICT (service, exact_error_identifier, deleted_at) DO UPDATE SET
98    deleted_at = NULL,
99    resolved_at = NULL,
100    occurrence_count = error_variants.occurrence_count + 1,
101    last_seen_at = now(),
102    updated_at = now()
103RETURNING *
104        "#,
105        service,
106        exact_error_identifier,
107        error_grouping_identifier,
108        error_source as ErrorSource,
109        report.message,
110        report.stack_trace,
111        normalized_message,
112        normalized_stack_trace
113    )
114    .fetch_one(&mut *tx)
115    .await?
116    .id;
117
118    let occurrence_id = Uuid::new_v4();
119    sqlx::query!(
120        "
121INSERT INTO error_occurrences (id, error_variant_id, service, user_id, path, app_version, details)
122VALUES ($1, $2, $3, $4, $5, $6, $7)
123        ",
124        occurrence_id,
125        variant_id,
126        service,
127        user_id,
128        report.path,
129        report.app_version,
130        report.details
131    )
132    .execute(&mut *tx)
133    .await?;
134
135    let service_consistent = sqlx::query_scalar!(
136        r#"
137SELECT EXISTS(
138    SELECT 1
139    FROM error_occurrences o
140    JOIN error_variants v ON v.id = o.error_variant_id
141    WHERE o.id = $1
142      AND o.service = v.service
143) AS "exists!"
144        "#,
145        occurrence_id
146    )
147    .fetch_one(&mut *tx)
148    .await?;
149    if !service_consistent {
150        return Err(model_err!(
151            Generic,
152            "error occurrence service must match variant service".to_string()
153        ));
154    }
155
156    tx.commit().await?;
157
158    Ok(variant_id)
159}
160
161/// Reports aren't worth adding to the pool pressure of the outage they may be reporting.
162const BEST_EFFORT_ACQUIRE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
163
164/// [`insert`] for a report that must never affect its caller: a failure to get a connection quickly
165/// or to insert is only logged.
166pub async fn insert_best_effort(pool: &PgPool, report: &NewErrorReport) {
167    let mut conn = match tokio::time::timeout(BEST_EFFORT_ACQUIRE_TIMEOUT, pool.acquire()).await {
168        Ok(Ok(conn)) => conn,
169        Ok(Err(err)) => {
170            warn!(service = %report.service, error = %err, "Error report skipped: could not acquire a connection");
171            return;
172        }
173        Err(_) => {
174            warn!(service = %report.service, "Error report skipped: timed out acquiring a connection");
175            return;
176        }
177    };
178    if let Err(err) = insert(&mut conn, None, report).await {
179        debug!(service = %report.service, error = %err, "Error report insert failed");
180    }
181}
182
183pub async fn get_all_variants(
184    conn: &mut PgConnection,
185    pagination: Pagination,
186) -> ModelResult<Vec<ErrorVariant>> {
187    let res = sqlx::query!(
188        r#"
189SELECT
190    *
191FROM error_variants
192WHERE deleted_at IS NULL
193ORDER BY last_seen_at DESC
194LIMIT $1 OFFSET $2
195        "#,
196        pagination.limit(),
197        pagination.offset()
198    )
199    .map(|r| ErrorVariant {
200        id: r.id,
201        service: r.service,
202        exact_error_identifier: r.exact_error_identifier,
203        error_grouping_identifier: r.error_grouping_identifier,
204        error_source: r.error_source,
205        example_message: r.example_message,
206        example_stack_trace: r.example_stack_trace,
207        normalized_message: r.normalized_message,
208        normalized_stack_trace: r.normalized_stack_trace,
209        occurrence_count: r.occurrence_count,
210        last_seen_at: r.last_seen_at,
211        resolved_at: r.resolved_at,
212        created_at: r.created_at,
213        updated_at: r.updated_at,
214        deleted_at: r.deleted_at,
215    })
216    .fetch_all(conn)
217    .await?;
218    Ok(res)
219}
220
221pub async fn delete_expired(conn: &mut PgConnection) -> ModelResult<()> {
222    use std::collections::HashSet;
223
224    let mut tx = conn.begin().await?;
225    let deleted_variant_ids = sqlx::query!(
226        r#"
227UPDATE error_occurrences
228SET deleted_at = now()
229WHERE created_at < now() - interval '2 months'
230  AND deleted_at IS NULL
231RETURNING *
232        "#
233    )
234    .fetch_all(&mut *tx)
235    .await?
236    .into_iter()
237    .map(|r| r.error_variant_id)
238    .collect::<HashSet<_>>()
239    .into_iter()
240    .collect::<Vec<_>>();
241
242    // Benign race: a concurrent insert between DELETE and this UPDATE can briefly skew aggregates.
243    sqlx::query!(
244        r#"
245UPDATE error_variants v
246SET
247    occurrence_count = (
248        SELECT COUNT(*)::int
249        FROM error_occurrences
250        WHERE error_variant_id = v.id
251          AND deleted_at IS NULL
252    ),
253    last_seen_at = COALESCE(
254        (
255            SELECT MAX(created_at)
256            FROM error_occurrences
257            WHERE error_variant_id = v.id
258              AND deleted_at IS NULL
259        ),
260        v.created_at
261    ),
262    updated_at = now()
263WHERE v.id = ANY($1)
264        "#,
265        &deleted_variant_ids
266    )
267    .execute(&mut *tx)
268    .await?;
269
270    tx.commit().await?;
271    Ok(())
272}
273
274pub async fn maybe_delete_expired(conn: &mut PgConnection) -> ModelResult<()> {
275    if rand::rng().random_range(1..=1000) == 1 {
276        info!("Cleaning up expired errors");
277        delete_expired(conn).await?;
278    }
279    Ok(())
280}