1use async_trait::async_trait;
7use headless_lms_utils::services::suotar::{
8 REFUSED_BEFORE_SENDING_CODE, SuotarCallAudit, SuotarCallFinished, SuotarCallStarted,
9};
10use utoipa::ToSchema;
11
12pub use headless_lms_utils::services::suotar::SuotarEndpoint;
15
16use crate::library::credit_registration::scrub::{scrub_suotar_body, scrub_text};
17use crate::prelude::*;
18
19pub const RETENTION_DAYS: i64 = 90;
21
22pub const FULL_BODY_ITEM_LIMIT: usize = 20;
24
25pub const SAMPLED_BODY_ITEM_COUNT: usize = 5;
27const _: () = assert!(SAMPLED_BODY_ITEM_COUNT <= FULL_BODY_ITEM_LIMIT);
29
30pub const BODY_SAMPLE_MAX_BYTES: usize = 64 * 1024;
32
33#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
34pub struct SuotarApiCall {
35 pub id: Uuid,
36 pub created_at: DateTime<Utc>,
37 pub updated_at: DateTime<Utc>,
38 pub deleted_at: Option<DateTime<Utc>>,
39 pub endpoint: SuotarEndpoint,
40 pub request_item_count: i32,
41 pub http_status: Option<i32>,
42 pub duration_ms: Option<i32>,
43 pub succeeded: bool,
44 pub ok_item_count: i32,
45 pub error_item_count: i32,
46 pub pending_item_count: i32,
47 pub request_level_error_code: Option<String>,
48 pub error_message: Option<String>,
49 pub request_body_sample: Option<serde_json::Value>,
50 pub response_body_sample: Option<serde_json::Value>,
51 pub credit_registration_ids: Vec<Uuid>,
52 pub worker_name: String,
53 pub started_at: DateTime<Utc>,
54 pub request_item_ids: Vec<String>,
56}
57
58#[derive(Debug, Clone, PartialEq)]
59pub struct NewSuotarApiCall {
60 pub endpoint: SuotarEndpoint,
61 pub request_item_count: i32,
62 pub http_status: Option<i32>,
63 pub duration_ms: Option<i32>,
64 pub succeeded: bool,
65 pub ok_item_count: i32,
66 pub error_item_count: i32,
67 pub pending_item_count: i32,
68 pub request_level_error_code: Option<String>,
69 pub error_message: Option<String>,
71 pub request_body_sample: Option<serde_json::Value>,
73 pub response_body_sample: Option<serde_json::Value>,
75 pub credit_registration_ids: Vec<Uuid>,
76 pub worker_name: String,
77 pub started_at: DateTime<Utc>,
78 pub request_item_ids: Vec<String>,
79}
80
81pub async fn insert(conn: &mut PgConnection, new: &NewSuotarApiCall) -> ModelResult<Uuid> {
82 let res = sqlx::query!(
83 r#"
84INSERT INTO suotar_api_calls (
85 endpoint,
86 request_item_count,
87 http_status,
88 duration_ms,
89 succeeded,
90 ok_item_count,
91 error_item_count,
92 pending_item_count,
93 request_level_error_code,
94 error_message,
95 request_body_sample,
96 response_body_sample,
97 credit_registration_ids,
98 worker_name,
99 started_at,
100 request_item_ids
101 )
102VALUES (
103 $1,
104 $2,
105 $3,
106 $4,
107 $5,
108 $6,
109 $7,
110 $8,
111 $9,
112 $10,
113 $11,
114 $12,
115 $13,
116 $14,
117 $15,
118 $16
119 )
120RETURNING id
121 "#,
122 new.endpoint as SuotarEndpoint,
123 new.request_item_count,
124 new.http_status,
125 new.duration_ms,
126 new.succeeded,
127 new.ok_item_count,
128 new.error_item_count,
129 new.pending_item_count,
130 new.request_level_error_code,
131 new.error_message,
132 new.request_body_sample,
133 new.response_body_sample,
134 &new.credit_registration_ids,
135 new.worker_name,
136 new.started_at,
137 &new.request_item_ids,
138 )
139 .fetch_one(conn)
140 .await?;
141 Ok(res.id)
142}
143
144#[derive(Debug, Clone, PartialEq)]
146pub struct FinishedSuotarApiCall {
147 pub http_status: Option<i32>,
148 pub duration_ms: Option<i32>,
149 pub succeeded: bool,
150 pub ok_item_count: i32,
151 pub error_item_count: i32,
152 pub pending_item_count: i32,
153 pub request_level_error_code: Option<String>,
154 pub error_message: Option<String>,
156 pub response_body_sample: Option<serde_json::Value>,
158}
159
160pub async fn finish(
161 conn: &mut PgConnection,
162 id: Uuid,
163 finished: &FinishedSuotarApiCall,
164) -> ModelResult<()> {
165 sqlx::query!(
166 r#"
167UPDATE suotar_api_calls
168SET http_status = $2,
169 duration_ms = $3,
170 succeeded = $4,
171 ok_item_count = $5,
172 error_item_count = $6,
173 pending_item_count = $7,
174 request_level_error_code = $8,
175 error_message = $9,
176 response_body_sample = $10,
177 updated_at = now()
178WHERE id = $1
179 "#,
180 id,
181 finished.http_status,
182 finished.duration_ms,
183 finished.succeeded,
184 finished.ok_item_count,
185 finished.error_item_count,
186 finished.pending_item_count,
187 finished.request_level_error_code,
188 finished.error_message,
189 finished.response_body_sample,
190 )
191 .execute(conn)
192 .await?;
193 Ok(())
194}
195
196fn sample_body(value: &serde_json::Value) -> serde_json::Value {
199 let sampled = match value.as_array() {
200 Some(items) if items.len() > FULL_BODY_ITEM_LIMIT => serde_json::json!({
201 "items": &items[..SAMPLED_BODY_ITEM_COUNT],
202 "totalItemCount": items.len(),
203 }),
204 _ => value.clone(),
205 };
206 let byte_count = serde_json::to_vec(&sampled)
207 .map(|bytes| bytes.len())
208 .unwrap_or(usize::MAX);
209 if byte_count <= BODY_SAMPLE_MAX_BYTES {
210 return sampled;
211 }
212 serde_json::json!({ "omitted": "over the sample size limit", "byteCount": byte_count })
213}
214
215pub struct PgSuotarCallAudit {
220 pool: PgPool,
221 is_waiting_item: fn(SuotarEndpoint, &str) -> bool,
222}
223
224impl PgSuotarCallAudit {
225 pub fn new(pool: PgPool, is_waiting_item: fn(SuotarEndpoint, &str) -> bool) -> Self {
228 Self {
229 pool,
230 is_waiting_item,
231 }
232 }
233}
234
235#[async_trait]
236impl SuotarCallAudit for PgSuotarCallAudit {
237 async fn started(&self, started: SuotarCallStarted) -> Option<Uuid> {
238 let new = NewSuotarApiCall {
239 endpoint: started.endpoint,
240 request_item_count: started.request_item_count.try_into().unwrap_or(i32::MAX),
241 http_status: None,
242 duration_ms: None,
243 succeeded: false,
244 ok_item_count: 0,
245 error_item_count: 0,
246 pending_item_count: 0,
247 request_level_error_code: None,
248 error_message: None,
249 request_body_sample: Some(sample_body(&scrub_suotar_body(&started.request_body))),
250 response_body_sample: None,
251 credit_registration_ids: started.credit_registration_ids,
252 worker_name: started.worker_name,
253 started_at: started.started_at,
254 request_item_ids: started.request_item_ids,
255 };
256 let mut conn = match self.pool.acquire().await {
257 Ok(conn) => conn,
258 Err(error) => {
259 error!("Could not open a connection for a suotar_api_calls row: {error}");
260 return None;
261 }
262 };
263 match insert(&mut conn, &new).await {
264 Ok(id) => Some(id),
265 Err(error) => {
266 error!("Could not insert a suotar_api_calls row: {error}");
267 None
268 }
269 }
270 }
271
272 async fn finished(&self, call_id: Uuid, finished: SuotarCallFinished) {
273 let pending_item_count = finished.endpoint.map_or(0, |endpoint| {
274 finished
275 .error_item_codes
276 .iter()
277 .filter(|code| (self.is_waiting_item)(endpoint, code))
278 .count()
279 });
280 let finished = FinishedSuotarApiCall {
281 http_status: finished.http_status.map(i32::from),
282 duration_ms: Some(finished.duration.as_millis().try_into().unwrap_or(i32::MAX)),
283 succeeded: finished.succeeded,
284 ok_item_count: finished.ok_item_count.try_into().unwrap_or(i32::MAX),
285 error_item_count: finished
286 .error_item_count
287 .saturating_sub(pending_item_count)
288 .try_into()
289 .unwrap_or(i32::MAX),
290 pending_item_count: pending_item_count.try_into().unwrap_or(i32::MAX),
291 request_level_error_code: finished.request_level_error_code,
292 error_message: finished.error_message.map(|message| scrub_text(&message)),
293 response_body_sample: finished
294 .response_body
295 .map(|body| sample_body(&scrub_suotar_body(&body))),
296 };
297 let mut conn = match self.pool.acquire().await {
298 Ok(conn) => conn,
299 Err(error) => {
300 error!(
301 "Could not open a connection to complete suotar_api_calls {call_id}: {error}"
302 );
303 return;
304 }
305 };
306 if let Err(error) = finish(&mut conn, call_id, &finished).await {
307 error!("Could not complete suotar_api_calls {call_id}: {error}");
308 }
309 }
310}
311
312pub async fn get_by_id(conn: &mut PgConnection, id: Uuid) -> ModelResult<SuotarApiCall> {
313 let res = sqlx::query_as!(
314 SuotarApiCall,
315 r#"
316SELECT *
317FROM suotar_api_calls
318WHERE id = $1
319 AND deleted_at IS NULL
320 "#,
321 id
322 )
323 .fetch_one(conn)
324 .await?;
325 Ok(res)
326}
327
328pub async fn get_by_credit_registration_id(
332 conn: &mut PgConnection,
333 credit_registration_id: Uuid,
334 limit: i64,
335) -> ModelResult<Vec<SuotarApiCall>> {
336 let res = sqlx::query_as!(
337 SuotarApiCall,
338 r#"
339SELECT *
340FROM suotar_api_calls
341WHERE credit_registration_ids @> ARRAY [$1::uuid]
342 AND deleted_at IS NULL
343ORDER BY started_at DESC
344LIMIT $2
345 "#,
346 credit_registration_id,
347 limit,
348 )
349 .fetch_all(conn)
350 .await?;
351 Ok(res)
352}
353
354#[derive(Debug, Clone, Default)]
356pub struct SuotarApiCallFilters {
357 pub endpoint: Option<SuotarEndpoint>,
358 pub succeeded: Option<bool>,
359 pub worker_name: Option<String>,
361 pub started_after: Option<DateTime<Utc>>,
362 pub started_before: Option<DateTime<Utc>>,
363 pub credit_registration_id: Option<Uuid>,
366}
367
368pub struct SuotarApiCallPageRow {
374 pub id: Uuid,
375 pub endpoint: SuotarEndpoint,
376 pub request_item_count: i32,
377 pub http_status: Option<i32>,
378 pub duration_ms: Option<i32>,
379 pub succeeded: bool,
380 pub ok_item_count: i32,
381 pub error_item_count: i32,
382 pub pending_item_count: i32,
383 pub request_level_error_code: Option<String>,
384 pub credit_registration_ids: Vec<Uuid>,
385 pub worker_name: String,
386 pub started_at: DateTime<Utc>,
387 pub total_count: i64,
388}
389
390pub async fn get_page(
392 conn: &mut PgConnection,
393 filters: &SuotarApiCallFilters,
394 limit: i64,
395 offset: i64,
396) -> ModelResult<Vec<SuotarApiCallPageRow>> {
397 let res = sqlx::query_as!(
398 SuotarApiCallPageRow,
399 r#"
400SELECT id,
401 endpoint AS "endpoint!",
402 request_item_count,
403 http_status,
404 duration_ms,
405 succeeded,
406 ok_item_count,
407 error_item_count,
408 pending_item_count,
409 request_level_error_code,
410 credit_registration_ids,
411 worker_name,
412 started_at,
413 COUNT(*) OVER () AS "total_count!"
414FROM suotar_api_calls
415WHERE deleted_at IS NULL
416 AND (
417 $1::suotar_endpoint IS NULL
418 OR endpoint = $1
419 )
420 AND ($2::bool IS NULL OR succeeded = $2)
421 AND ($3::text IS NULL OR worker_name = $3)
422 AND ($4::timestamptz IS NULL OR started_at >= $4)
423 AND ($5::timestamptz IS NULL OR started_at <= $5)
424 AND (
425 $6::uuid IS NULL
426 OR credit_registration_ids @> ARRAY [$6::uuid]
427 )
428ORDER BY started_at DESC,
429 id
430LIMIT $7 OFFSET $8
431 "#,
432 filters.endpoint as Option<SuotarEndpoint>,
433 filters.succeeded,
434 filters.worker_name.as_deref(),
435 filters.started_after,
436 filters.started_before,
437 filters.credit_registration_id,
438 limit,
439 offset,
440 )
441 .fetch_all(conn)
442 .await?;
443 Ok(res)
444}
445
446pub async fn get_worker_names(conn: &mut PgConnection) -> ModelResult<Vec<String>> {
449 let res = sqlx::query_scalar!(
450 r#"
451SELECT DISTINCT worker_name
452FROM suotar_api_calls
453WHERE deleted_at IS NULL
454ORDER BY worker_name
455 "#,
456 )
457 .fetch_all(conn)
458 .await?;
459 Ok(res)
460}
461
462#[derive(Debug, Clone, PartialEq)]
468pub struct SuotarEndpointStatsForWindow {
469 pub window_secs: i64,
470 pub endpoint: SuotarEndpoint,
471 pub call_count: i64,
472 pub failed_call_count: i64,
473 pub in_flight_count: i64,
474 pub ok_item_count: i64,
475 pub error_item_count: i64,
476 pub pending_item_count: i64,
477 pub p50_duration_ms: Option<i32>,
478 pub p95_duration_ms: Option<i32>,
479 pub last_success_at: Option<DateTime<Utc>>,
480 pub last_failure_at: Option<DateTime<Utc>>,
481 pub last_request_level_error_code: Option<String>,
482}
483
484pub async fn get_endpoint_stats_for_windows(
487 conn: &mut PgConnection,
488 window_secs: &[i64],
489) -> ModelResult<Vec<SuotarEndpointStatsForWindow>> {
490 let now = Utc::now();
491 let window_secs = window_secs.to_vec();
492 let since: Vec<DateTime<Utc>> = window_secs
493 .iter()
494 .map(|secs| now - chrono::Duration::seconds(*secs))
495 .collect();
496 let rows = sqlx::query_as!(
497 SuotarEndpointStatsForWindow,
498 r#"
499WITH windows AS (
500 SELECT * FROM UNNEST($1::bigint [], $2::timestamptz []) AS w(window_secs, since)
501)
502SELECT w.window_secs AS "window_secs!",
503 c.endpoint,
504 COUNT(*) FILTER (WHERE c.duration_ms IS NOT NULL) AS "call_count!",
505 COUNT(*) FILTER (
506 WHERE c.duration_ms IS NOT NULL
507 AND NOT c.succeeded
508 ) AS "failed_call_count!",
509 COUNT(*) FILTER (WHERE c.duration_ms IS NULL) AS "in_flight_count!",
510 COALESCE(SUM(c.ok_item_count), 0) AS "ok_item_count!",
511 COALESCE(SUM(c.error_item_count), 0) AS "error_item_count!",
512 COALESCE(SUM(c.pending_item_count), 0) AS "pending_item_count!",
513 PERCENTILE_DISC(0.5) WITHIN GROUP (
514 ORDER BY c.duration_ms
515 ) AS "p50_duration_ms",
516 PERCENTILE_DISC(0.95) WITHIN GROUP (
517 ORDER BY c.duration_ms
518 ) AS "p95_duration_ms",
519 MAX(c.started_at) FILTER (WHERE c.succeeded) AS "last_success_at",
520 MAX(c.started_at) FILTER (
521 WHERE c.duration_ms IS NOT NULL
522 AND NOT c.succeeded
523 ) AS "last_failure_at",
524 (
525 ARRAY_AGG(
526 c.request_level_error_code
527 ORDER BY c.started_at DESC
528 ) FILTER (
529 WHERE c.duration_ms IS NOT NULL
530 AND NOT c.succeeded
531 AND c.request_level_error_code IS NOT NULL
532 )
533 ) [1] AS "last_request_level_error_code"
534FROM windows w
535 JOIN suotar_api_calls c ON c.started_at >= w.since
536 AND c.deleted_at IS NULL
537GROUP BY w.window_secs,
538 c.endpoint
539 "#,
540 &window_secs,
541 &since,
542 )
543 .fetch_all(conn)
544 .await?;
545 Ok(rows)
546}
547
548#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, ToSchema)]
550pub struct SuotarEndpointDailyCost {
551 pub day: chrono::NaiveDate,
552 pub endpoint: SuotarEndpoint,
553 pub call_count: i64,
554 pub failed_call_count: i64,
555 pub item_count: i64,
556 pub max_items_per_call: i32,
557 pub p50_duration_ms: Option<i32>,
558 pub p95_duration_ms: Option<i32>,
559}
560
561pub async fn get_daily_costs_since(
563 conn: &mut PgConnection,
564 since: DateTime<Utc>,
565) -> ModelResult<Vec<SuotarEndpointDailyCost>> {
566 let rows = sqlx::query_as!(
567 SuotarEndpointDailyCost,
568 r#"
569SELECT (started_at AT TIME ZONE 'UTC')::date AS "day!",
570 endpoint AS "endpoint!",
571 COUNT(*) AS "call_count!",
572 COUNT(*) FILTER (
573 WHERE NOT succeeded
574 ) AS "failed_call_count!",
575 COALESCE(SUM(request_item_count), 0) AS "item_count!",
576 COALESCE(MAX(request_item_count), 0) AS "max_items_per_call!",
577 PERCENTILE_DISC(0.5) WITHIN GROUP (
578 ORDER BY duration_ms
579 ) AS p50_duration_ms,
580 PERCENTILE_DISC(0.95) WITHIN GROUP (
581 ORDER BY duration_ms
582 ) AS p95_duration_ms
583FROM suotar_api_calls
584WHERE started_at >= $1
585 AND duration_ms IS NOT NULL
586 AND deleted_at IS NULL
587GROUP BY 1,
588 2
589ORDER BY 1 DESC,
590 endpoint::text
591 "#,
592 since,
593 )
594 .fetch_all(conn)
595 .await?;
596 Ok(rows)
597}
598
599#[derive(Debug, Clone, PartialEq)]
601pub struct SuotarEndpointStanding {
602 pub endpoint: SuotarEndpoint,
603 pub last_success_at: Option<DateTime<Utc>>,
604 pub last_failure_at: Option<DateTime<Utc>>,
605 pub consecutive_failures: i64,
607}
608
609pub async fn get_endpoint_standings(
612 conn: &mut PgConnection,
613) -> ModelResult<Vec<SuotarEndpointStanding>> {
614 let since = Utc::now() - chrono::Duration::days(RETENTION_DAYS);
615 let rows = sqlx::query_as!(
616 SuotarEndpointStanding,
617 r#"
618WITH last_success AS (
619 SELECT endpoint,
620 MAX(started_at) AS at
621 FROM suotar_api_calls
622 WHERE succeeded
623 AND started_at >= $1
624 AND deleted_at IS NULL
625 GROUP BY endpoint
626)
627SELECT c.endpoint AS "endpoint!: SuotarEndpoint",
628 ls.at AS "last_success_at",
629 MAX(c.started_at) FILTER (
630 WHERE c.duration_ms IS NOT NULL
631 AND NOT c.succeeded
632 ) AS "last_failure_at",
633 COUNT(*) FILTER (
634 WHERE c.duration_ms IS NOT NULL
635 AND NOT c.succeeded
636 AND (
637 ls.at IS NULL
638 OR c.started_at > ls.at
639 )
640 ) AS "consecutive_failures!"
641FROM suotar_api_calls c
642 LEFT JOIN last_success ls ON ls.endpoint = c.endpoint
643WHERE c.started_at >= $1
644 AND c.deleted_at IS NULL
645GROUP BY c.endpoint,
646 ls.at
647 "#,
648 since,
649 )
650 .fetch_all(conn)
651 .await?;
652 Ok(rows)
653}
654
655#[derive(Debug, Clone, PartialEq)]
657pub struct SuotarFailureRun {
658 pub count: i64,
659 pub last_at: Option<DateTime<Utc>>,
660}
661
662pub async fn count_credential_rejections_since(
664 conn: &mut PgConnection,
665 since: DateTime<Utc>,
666) -> ModelResult<SuotarFailureRun> {
667 let row = sqlx::query_as!(
668 SuotarFailureRun,
669 r#"
670SELECT COUNT(*) AS "count!",
671 MAX(started_at) AS "last_at"
672FROM suotar_api_calls
673WHERE started_at >= $1
674 AND deleted_at IS NULL
675 AND (
676 http_status IN (401, 403)
677 OR request_level_error_code = 'unauthorized'
678 )
679 "#,
680 since,
681 )
682 .fetch_one(conn)
683 .await?;
684 Ok(row)
685}
686
687pub async fn count_unreachable_run_since(
691 conn: &mut PgConnection,
692 since: DateTime<Utc>,
693) -> ModelResult<SuotarFailureRun> {
694 let row = sqlx::query_as!(
695 SuotarFailureRun,
696 r#"
697WITH last_success AS (
698 SELECT MAX(started_at) AS at
699 FROM suotar_api_calls
700 WHERE succeeded
701 AND started_at >= $1
702 AND deleted_at IS NULL
703)
704SELECT COUNT(*) AS "count!",
705 MAX(c.started_at) AS "last_at"
706FROM suotar_api_calls c
707 CROSS JOIN last_success ls
708WHERE c.started_at >= $1
709 AND c.deleted_at IS NULL
710 AND c.duration_ms IS NOT NULL
711 AND NOT c.succeeded
712 AND c.request_level_error_code IS DISTINCT FROM $2
713 AND (
714 c.http_status IS NULL
715 OR c.http_status >= 500
716 )
717 AND (
718 ls.at IS NULL
719 OR c.started_at > ls.at
720 )
721 "#,
722 since,
723 REFUSED_BEFORE_SENDING_CODE,
724 )
725 .fetch_one(conn)
726 .await?;
727 Ok(row)
728}
729
730pub async fn delete_older_than(
737 conn: &mut PgConnection,
738 cutoff: DateTime<Utc>,
739 limit: i64,
740) -> ModelResult<u64> {
741 let res = sqlx::query!(
742 r#"
743DELETE FROM suotar_api_calls
744WHERE id IN (
745 SELECT id
746 FROM suotar_api_calls
747 WHERE started_at < $1
748 ORDER BY started_at
749 LIMIT $2
750 )
751 "#,
752 cutoff,
753 limit,
754 )
755 .execute(conn)
756 .await?;
757 Ok(res.rows_affected())
758}