Skip to main content

headless_lms_models/
credit_registration_events.rs

1//! Append-only audit trail for the credit registration ledger.
2//!
3//! No retention sweep touches this table, so every Suotar payload must go through
4//! [`scrub_suotar_body`](crate::library::credit_registration::scrub::scrub_suotar_body) at the
5//! write site — redacting on read would leave the raw values on disk.
6use std::collections::HashMap;
7
8use serde_json::Value;
9use utoipa::ToSchema;
10
11use crate::credit_registrations::{CreditRegistrationErrorCode, CreditRegistrationState};
12use crate::prelude::*;
13use crate::suotar_api_calls::SuotarEndpoint;
14
15#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Copy, Hash, Type, ToSchema)]
16#[sqlx(
17    type_name = "credit_registration_event_kind",
18    rename_all = "snake_case"
19)]
20#[serde(rename_all = "snake_case")]
21pub enum CreditRegistrationEventKind {
22    Created,
23    StateChanged,
24    SuotarResponse,
25    RetryScheduled,
26    AdminAction,
27    StudentAction,
28    Cancelled,
29}
30
31/// What Suotar did with one row of a request.
32#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Copy, Hash, Type, ToSchema)]
33#[sqlx(type_name = "suotar_answer", rename_all = "snake_case")]
34#[serde(rename_all = "snake_case")]
35pub enum SuotarAnswer {
36    Answered,
37    /// Suotar answered the request but left this row out.
38    Unanswered,
39    /// Suotar refused the whole request, or it never got there.
40    Refused,
41}
42
43#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
44pub struct CreditRegistrationEvent {
45    pub id: Uuid,
46    pub created_at: DateTime<Utc>,
47    pub updated_at: DateTime<Utc>,
48    pub deleted_at: Option<DateTime<Utc>>,
49    pub credit_registration_id: Uuid,
50    pub kind: CreditRegistrationEventKind,
51    pub from_state: Option<CreditRegistrationState>,
52    pub to_state: Option<CreditRegistrationState>,
53    pub error_code: Option<CreditRegistrationErrorCode>,
54    pub message: Option<String>,
55    pub suotar_api_call_id: Option<Uuid>,
56    pub actor_user_id: Option<Uuid>,
57    pub details: Option<Value>,
58    /// The requestItemId the row went out under in the call behind this event.
59    pub request_item_id: Option<String>,
60    /// Outlives `suotar_api_call_id`, whose call row is swept after 90 days.
61    pub suotar_endpoint: Option<SuotarEndpoint>,
62    pub suotar_requested_at: Option<DateTime<Utc>>,
63    /// Orders the timeline where set: events written in one transaction share `created_at`.
64    pub suotar_answered_at: Option<DateTime<Utc>>,
65    pub suotar_answer: Option<SuotarAnswer>,
66}
67
68#[derive(Debug, Clone, PartialEq)]
69pub struct NewCreditRegistrationEvent {
70    pub credit_registration_id: Uuid,
71    pub kind: CreditRegistrationEventKind,
72    pub from_state: Option<CreditRegistrationState>,
73    pub to_state: Option<CreditRegistrationState>,
74    pub error_code: Option<CreditRegistrationErrorCode>,
75    pub message: Option<String>,
76    pub suotar_api_call_id: Option<Uuid>,
77    pub actor_user_id: Option<Uuid>,
78    /// Build with
79    /// [`suotar_exchange_details`](crate::library::credit_registration::scrub::suotar_exchange_details)
80    /// so it is scrubbed.
81    pub details: Option<Value>,
82    pub request_item_id: Option<String>,
83    pub suotar_endpoint: Option<SuotarEndpoint>,
84    pub suotar_requested_at: Option<DateTime<Utc>>,
85    pub suotar_answered_at: Option<DateTime<Utc>>,
86    pub suotar_answer: Option<SuotarAnswer>,
87}
88
89impl NewCreditRegistrationEvent {
90    pub fn new(credit_registration_id: Uuid, kind: CreditRegistrationEventKind) -> Self {
91        Self {
92            credit_registration_id,
93            kind,
94            from_state: None,
95            to_state: None,
96            error_code: None,
97            message: None,
98            suotar_api_call_id: None,
99            actor_user_id: None,
100            details: None,
101            request_item_id: None,
102            suotar_endpoint: None,
103            suotar_requested_at: None,
104            suotar_answered_at: None,
105            suotar_answer: None,
106        }
107    }
108}
109
110/// Callers that also change `state` must go through `credit_registrations::transition` instead,
111/// which writes both in one transaction.
112pub async fn insert(
113    conn: &mut PgConnection,
114    new: &NewCreditRegistrationEvent,
115) -> ModelResult<Uuid> {
116    let res = sqlx::query!(
117        r#"
118INSERT INTO credit_registration_events (
119    credit_registration_id,
120    kind,
121    from_state,
122    to_state,
123    error_code,
124    message,
125    suotar_api_call_id,
126    actor_user_id,
127    details,
128    request_item_id,
129    suotar_endpoint,
130    suotar_requested_at,
131    suotar_answered_at,
132    suotar_answer
133  )
134VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)
135RETURNING id
136        "#,
137        new.credit_registration_id,
138        new.kind as CreditRegistrationEventKind,
139        new.from_state as Option<CreditRegistrationState>,
140        new.to_state as Option<CreditRegistrationState>,
141        new.error_code as Option<CreditRegistrationErrorCode>,
142        new.message,
143        new.suotar_api_call_id,
144        new.actor_user_id,
145        new.details,
146        new.request_item_id,
147        new.suotar_endpoint as Option<SuotarEndpoint>,
148        new.suotar_requested_at,
149        new.suotar_answered_at,
150        new.suotar_answer as Option<SuotarAnswer>,
151    )
152    .fetch_one(conn)
153    .await?;
154    Ok(res.id)
155}
156
157/// Appends one event per element, in one statement. The column list is [`insert`]'s, so a batch
158/// writes the same rows a loop would.
159pub async fn insert_batch(
160    conn: &mut PgConnection,
161    events: &[NewCreditRegistrationEvent],
162) -> ModelResult<()> {
163    if events.is_empty() {
164        return Ok(());
165    }
166    let ids: Vec<Uuid> = events.iter().map(|e| e.credit_registration_id).collect();
167    let kinds: Vec<CreditRegistrationEventKind> = events.iter().map(|e| e.kind).collect();
168    let from_states: Vec<Option<CreditRegistrationState>> =
169        events.iter().map(|e| e.from_state).collect();
170    let to_states: Vec<Option<CreditRegistrationState>> =
171        events.iter().map(|e| e.to_state).collect();
172    let error_codes: Vec<Option<CreditRegistrationErrorCode>> =
173        events.iter().map(|e| e.error_code).collect();
174    let messages: Vec<Option<String>> = events.iter().map(|e| e.message.clone()).collect();
175    let call_ids: Vec<Option<Uuid>> = events.iter().map(|e| e.suotar_api_call_id).collect();
176    let actors: Vec<Option<Uuid>> = events.iter().map(|e| e.actor_user_id).collect();
177    let details: Vec<Option<Value>> = events.iter().map(|e| e.details.clone()).collect();
178    let request_item_ids: Vec<Option<String>> =
179        events.iter().map(|e| e.request_item_id.clone()).collect();
180    let endpoints: Vec<Option<SuotarEndpoint>> = events.iter().map(|e| e.suotar_endpoint).collect();
181    let requested_ats: Vec<Option<DateTime<Utc>>> =
182        events.iter().map(|e| e.suotar_requested_at).collect();
183    let answered_ats: Vec<Option<DateTime<Utc>>> =
184        events.iter().map(|e| e.suotar_answered_at).collect();
185    let answers: Vec<Option<SuotarAnswer>> = events.iter().map(|e| e.suotar_answer).collect();
186    sqlx::query!(
187        r#"
188INSERT INTO credit_registration_events (
189    credit_registration_id,
190    kind,
191    from_state,
192    to_state,
193    error_code,
194    message,
195    suotar_api_call_id,
196    actor_user_id,
197    details,
198    request_item_id,
199    suotar_endpoint,
200    suotar_requested_at,
201    suotar_answered_at,
202    suotar_answer
203  )
204SELECT *
205FROM UNNEST(
206    $1::uuid [],
207    $2::credit_registration_event_kind [],
208    $3::credit_registration_state [],
209    $4::credit_registration_state [],
210    $5::credit_registration_error_code [],
211    $6::text [],
212    $7::uuid [],
213    $8::uuid [],
214    $9::jsonb [],
215    $10::text [],
216    $11::suotar_endpoint [],
217    $12::timestamptz [],
218    $13::timestamptz [],
219    $14::suotar_answer []
220  )
221        "#,
222        &ids,
223        &kinds as &[CreditRegistrationEventKind],
224        &from_states as &[Option<CreditRegistrationState>],
225        &to_states as &[Option<CreditRegistrationState>],
226        &error_codes as &[Option<CreditRegistrationErrorCode>],
227        &messages as &[Option<String>],
228        &call_ids as &[Option<Uuid>],
229        &actors as &[Option<Uuid>],
230        &details as &[Option<Value>],
231        &request_item_ids as &[Option<String>],
232        &endpoints as &[Option<SuotarEndpoint>],
233        &requested_ats as &[Option<DateTime<Utc>>],
234        &answered_ats as &[Option<DateTime<Utc>>],
235        &answers as &[Option<SuotarAnswer>],
236    )
237    .execute(conn)
238    .await?;
239    Ok(())
240}
241
242/// Appends the same event to many rows in one round trip; [`insert_batch`] for events that differ.
243pub async fn insert_many(
244    conn: &mut PgConnection,
245    credit_registration_ids: &[Uuid],
246    kind: CreditRegistrationEventKind,
247    actor_user_id: Option<Uuid>,
248    message: Option<&str>,
249) -> ModelResult<()> {
250    if credit_registration_ids.is_empty() {
251        return Ok(());
252    }
253    sqlx::query!(
254        r#"
255INSERT INTO credit_registration_events (credit_registration_id, kind, actor_user_id, message)
256SELECT id, $2, $3, $4
257FROM UNNEST($1::uuid []) AS id
258        "#,
259        credit_registration_ids,
260        kind as CreditRegistrationEventKind,
261        actor_user_id,
262        message,
263    )
264    .execute(conn)
265    .await?;
266    Ok(())
267}
268
269/// The per-item timeline, newest first.
270pub async fn get_by_registration_id(
271    conn: &mut PgConnection,
272    credit_registration_id: Uuid,
273) -> ModelResult<Vec<CreditRegistrationEvent>> {
274    let res = sqlx::query_as!(
275        CreditRegistrationEvent,
276        r#"
277SELECT *
278FROM credit_registration_events
279WHERE credit_registration_id = $1
280  AND deleted_at IS NULL
281ORDER BY COALESCE(suotar_answered_at, created_at) DESC,
282  id DESC
283        "#,
284        credit_registration_id
285    )
286    .fetch_all(conn)
287    .await?;
288    Ok(res)
289}
290
291/// The per-item timeline entries one study registry call produced, oldest first: what the answer
292/// did to each row it covered.
293pub async fn get_by_suotar_api_call_id(
294    conn: &mut PgConnection,
295    suotar_api_call_id: Uuid,
296) -> ModelResult<Vec<CreditRegistrationEvent>> {
297    let res = sqlx::query_as!(
298        CreditRegistrationEvent,
299        r#"
300SELECT *
301FROM credit_registration_events
302WHERE suotar_api_call_id = $1
303  AND deleted_at IS NULL
304ORDER BY COALESCE(suotar_answered_at, created_at),
305  id
306        "#,
307        suotar_api_call_id
308    )
309    .fetch_all(conn)
310    .await?;
311    Ok(res)
312}
313
314/// The requestItemId each of `credit_registration_ids` went out under among `request_item_ids`,
315/// which are one call's. A row the call carried but no event recorded is missing from the map.
316pub async fn get_request_item_ids_in_call(
317    conn: &mut PgConnection,
318    credit_registration_ids: &[Uuid],
319    request_item_ids: &[String],
320) -> ModelResult<HashMap<Uuid, String>> {
321    let rows = sqlx::query!(
322        r#"
323SELECT DISTINCT ON (credit_registration_id) credit_registration_id,
324  request_item_id AS "request_item_id!"
325FROM credit_registration_events
326WHERE credit_registration_id = ANY($1)
327  AND request_item_id = ANY($2)
328  AND deleted_at IS NULL
329ORDER BY credit_registration_id,
330  COALESCE(suotar_answered_at, created_at),
331  id
332        "#,
333        credit_registration_ids,
334        request_item_ids,
335    )
336    .fetch_all(conn)
337    .await?;
338    Ok(rows
339        .into_iter()
340        .map(|row| (row.credit_registration_id, row.request_item_id))
341        .collect())
342}
343
344/// The attainment the study registry pointed at when it turned a submission down as no improvement.
345///
346/// Read back off the event the answer was recorded on rather than stored on the ledger row: the row
347/// holds what we sent, and this is what the registry already had.
348#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
349pub struct NotImprovedAttainment {
350    pub grade_id: Option<String>,
351    /// Names the scale `grade_id` is on, without which "1" reads as a one out of five when it means
352    /// a pass.
353    pub grade_scale_id: Option<String>,
354}
355
356/// The registry's verdict for a row it declined as no improvement. `None` for every other row, and
357/// for one whose answer named no attainment.
358pub async fn get_not_improved_attainment(
359    conn: &mut PgConnection,
360    credit_registration_id: Uuid,
361) -> ModelResult<Option<NotImprovedAttainment>> {
362    let found = sqlx::query_scalar!(
363        r#"
364SELECT details #> '{response,result,previousAttainment}' AS "attainment!"
365FROM credit_registration_events
366WHERE credit_registration_id = $1
367  AND to_state = 'not_improved'
368  AND details #> '{response,result,previousAttainment}' IS NOT NULL
369  AND deleted_at IS NULL
370ORDER BY COALESCE(suotar_answered_at, created_at) DESC,
371  id DESC
372LIMIT 1
373        "#,
374        credit_registration_id
375    )
376    .fetch_optional(conn)
377    .await?;
378    Ok(found.map(|attainment| NotImprovedAttainment {
379        grade_id: string_field(&attainment, "gradeId"),
380        grade_scale_id: string_field(&attainment, "gradeScaleId"),
381    }))
382}
383
384fn string_field(value: &Value, key: &str) -> Option<String> {
385    Some(value.get(key)?.as_str()?.to_string())
386}
387
388/// How the study registry answered the items it was asked about in a window.
389#[derive(Debug, Clone, PartialEq, Default)]
390pub struct SuotarItemOutcomeTotals {
391    pub item_count: i64,
392    /// Items that failed on Suotar or Sisu being down or slow, the only proxy we have for Sisu's
393    /// uptime.
394    pub service_unavailable_count: i64,
395    pub last_service_unavailable_at: Option<DateTime<Utc>>,
396}
397
398/// Per-item outcomes in `[since, now)`, counted over events rather than rows: an item that failed
399/// really failed, whatever its row has since become.
400pub async fn count_item_outcomes_since(
401    conn: &mut PgConnection,
402    since: DateTime<Utc>,
403) -> ModelResult<SuotarItemOutcomeTotals> {
404    let res = sqlx::query_as!(
405        SuotarItemOutcomeTotals,
406        r#"
407SELECT COUNT(*) AS "item_count!",
408  COUNT(*) FILTER (
409    WHERE error_code IN ('service_temporarily_unavailable', 'sisu_timeout')
410  ) AS "service_unavailable_count!",
411  MAX(created_at) FILTER (
412    WHERE error_code IN ('service_temporarily_unavailable', 'sisu_timeout')
413  ) AS "last_service_unavailable_at"
414FROM credit_registration_events
415WHERE kind = 'suotar_response'
416  AND created_at >= $1
417  AND deleted_at IS NULL
418        "#,
419        since,
420    )
421    .fetch_one(conn)
422    .await?;
423    Ok(res)
424}
425
426/// One error code's standing over a window and the window before it.
427#[derive(Debug, Clone, PartialEq)]
428pub struct ErrorCodeWindowCounts {
429    pub error_code: CreditRegistrationErrorCode,
430    pub current_count: i64,
431    /// The equally long window immediately before, which is what a spike is measured against.
432    pub previous_count: i64,
433    pub user_count: i64,
434    pub course_count: i64,
435    pub first_seen_at: Option<DateTime<Utc>>,
436    pub last_seen_at: Option<DateTime<Utc>>,
437    /// The endpoints the code arrived on, empty for one recorded without a call of ours.
438    pub endpoints: Vec<SuotarEndpoint>,
439}
440
441/// Error events per code over `[now - 2 * window, now)`, split at `now - window`.
442///
443/// Counts events, so errors on attempts since superseded are included: an `invalid_credits` on
444/// attempt 1 is the configuration bug, whether or not attempt 2 succeeded.
445pub async fn get_error_code_counts_for_window(
446    conn: &mut PgConnection,
447    window_secs: i64,
448) -> ModelResult<Vec<ErrorCodeWindowCounts>> {
449    let res = sqlx::query_as!(
450        ErrorCodeWindowCounts,
451        r#"
452SELECT e.error_code AS "error_code!: CreditRegistrationErrorCode",
453  COUNT(*) FILTER (
454    WHERE e.created_at > now() - MAKE_INTERVAL(secs => $1::double precision)
455  ) AS "current_count!",
456  COUNT(*) FILTER (
457    WHERE e.created_at <= now() - MAKE_INTERVAL(secs => $1::double precision)
458  ) AS "previous_count!",
459  COUNT(DISTINCT cr.user_id) AS "user_count!",
460  COUNT(DISTINCT cr.course_id) AS "course_count!",
461  MIN(e.created_at) AS "first_seen_at",
462  MAX(e.created_at) AS "last_seen_at",
463  COALESCE(
464    ARRAY_AGG(DISTINCT call.endpoint) FILTER (
465      WHERE call.endpoint IS NOT NULL
466    ),
467    '{}'
468  ) AS "endpoints!: Vec<SuotarEndpoint>"
469FROM credit_registration_events e
470  JOIN credit_registrations cr ON cr.id = e.credit_registration_id
471  LEFT JOIN suotar_api_calls call ON call.id = e.suotar_api_call_id
472  AND call.deleted_at IS NULL
473WHERE e.error_code IS NOT NULL
474  AND e.created_at > now() - MAKE_INTERVAL(secs => $1::double precision * 2)
475  AND e.deleted_at IS NULL
476  AND cr.deleted_at IS NULL
477GROUP BY e.error_code
478ORDER BY 2 DESC
479        "#,
480        window_secs as f64,
481    )
482    .fetch_all(conn)
483    .await?;
484    Ok(res)
485}
486
487/// Registrations whose answers named more than one submitted attainment id, which is the shape a
488/// double submission would leave behind.
489///
490/// Per registration, not per completion: a grade improvement is a second registration row and a
491/// second attainment on purpose.
492pub async fn get_ids_with_several_submitted_attainments(
493    conn: &mut PgConnection,
494    limit: i64,
495) -> ModelResult<Vec<Uuid>> {
496    let res = sqlx::query_scalar!(
497        r#"
498SELECT credit_registration_id
499FROM credit_registration_events
500WHERE details #>> '{response,submittedAttainmentId}' IS NOT NULL
501  AND deleted_at IS NULL
502GROUP BY credit_registration_id
503HAVING COUNT(
504    DISTINCT details #>> '{response,submittedAttainmentId}'
505  ) > 1
506LIMIT $1
507        "#,
508        limit,
509    )
510    .fetch_all(conn)
511    .await?;
512    Ok(res)
513}