1use std::collections::{HashMap, HashSet};
8
9use secrecy::ExposeSecret;
10
11use crate::email_deliveries::{
12 EmailSendStatus, EmailSendStatusFacts, EmailSendStatusReport, derive_email_send_status,
13 get_send_statuses, is_hard_send_failure,
14};
15use crate::prelude::*;
16
17#[derive(Debug, Clone)]
18pub struct CreditRegistrationAccountLinkingEmail {
19 pub id: Uuid,
20 pub created_at: DateTime<Utc>,
21 pub updated_at: DateTime<Utc>,
22 pub deleted_at: Option<DateTime<Utc>>,
23 pub student_number: DbSecret,
24 pub sisu_person_id: DbSecret,
25 pub course_id: Uuid,
26 pub emailed_to: DbSecret,
27 pub student_number_verification_token_id: Option<Uuid>,
28 pub email_delivery_id: Option<Uuid>,
29 pub sent_at: DateTime<Utc>,
30}
31
32#[derive(Debug, Clone)]
33pub struct NewAccountLinkingEmail {
34 pub student_number: DbSecret,
35 pub sisu_person_id: DbSecret,
36 pub course_id: Uuid,
37 pub emailed_to: DbSecret,
38 pub student_number_verification_token_id: Option<Uuid>,
39 pub email_delivery_id: Option<Uuid>,
40}
41
42pub async fn claim_send_slot(
45 conn: &mut PgConnection,
46 new: &NewAccountLinkingEmail,
47) -> ModelResult<Option<Uuid>> {
48 let res = sqlx::query!(
49 r#"
50INSERT INTO credit_registration_account_linking_emails (
51 student_number,
52 sisu_person_id,
53 course_id,
54 emailed_to,
55 student_number_verification_token_id,
56 email_delivery_id
57 )
58VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT DO NOTHING
59RETURNING id
60 "#,
61 new.student_number.expose_secret(),
62 new.sisu_person_id.expose_secret(),
63 new.course_id,
64 new.emailed_to.expose_secret(),
65 new.student_number_verification_token_id,
66 new.email_delivery_id,
67 )
68 .fetch_optional(conn)
69 .await?;
70 Ok(res.map(|r| r.id))
71}
72
73#[derive(Debug, Clone)]
75pub struct ExistingLinkingMailFact {
76 pub sisu_person_id: DbSecret,
77 pub course_id: Uuid,
78 pub emailed_to: DbSecret,
79 pub sent_at: DateTime<Utc>,
80}
81
82pub async fn get_existing_facts_for_persons(
85 conn: &mut PgConnection,
86 sisu_person_ids: &[String],
87) -> ModelResult<Vec<ExistingLinkingMailFact>> {
88 let res = sqlx::query_as!(
89 ExistingLinkingMailFact,
90 r#"
91SELECT sisu_person_id,
92 course_id,
93 emailed_to,
94 sent_at
95FROM credit_registration_account_linking_emails
96WHERE sisu_person_id = ANY($1::text [])
97 AND deleted_at IS NULL
98 "#,
99 sisu_person_ids
100 )
101 .fetch_all(conn)
102 .await?;
103 Ok(res)
104}
105
106pub async fn claim_send_slots(
114 conn: &mut PgConnection,
115 new: &[NewAccountLinkingEmail],
116 token_ids: &[Uuid],
117) -> ModelResult<HashSet<Uuid>> {
118 if new.is_empty() {
119 return Ok(HashSet::new());
120 }
121 let student_numbers: Vec<String> = new
122 .iter()
123 .map(|n| n.student_number.expose_secret().to_owned())
124 .collect();
125 let sisu_person_ids: Vec<String> = new
126 .iter()
127 .map(|n| n.sisu_person_id.expose_secret().to_owned())
128 .collect();
129 let course_ids: Vec<Uuid> = new.iter().map(|n| n.course_id).collect();
130 let emailed_tos: Vec<String> = new
131 .iter()
132 .map(|n| n.emailed_to.expose_secret().to_owned())
133 .collect();
134
135 let claimed = sqlx::query_scalar!(
136 r#"
137INSERT INTO credit_registration_account_linking_emails (
138 student_number,
139 sisu_person_id,
140 course_id,
141 emailed_to,
142 student_number_verification_token_id
143 )
144SELECT * FROM UNNEST($1::text [], $2::text [], $3::uuid [], $4::text [], $5::uuid []) ON CONFLICT DO NOTHING
145RETURNING student_number_verification_token_id AS "token_id!"
146 "#,
147 &student_numbers,
148 &sisu_person_ids,
149 &course_ids,
150 &emailed_tos,
151 token_ids,
152 )
153 .fetch_all(conn)
154 .await?;
155 Ok(claimed.into_iter().collect())
156}
157
158pub async fn get_latest_by_course_and_persons(
160 conn: &mut PgConnection,
161 course_id: Uuid,
162 sisu_person_ids: &[String],
163) -> ModelResult<HashMap<String, CreditRegistrationAccountLinkingEmail>> {
164 let res = sqlx::query_as!(
165 CreditRegistrationAccountLinkingEmail,
166 r#"
167SELECT DISTINCT ON (sisu_person_id) *
168FROM credit_registration_account_linking_emails
169WHERE course_id = $1
170 AND sisu_person_id = ANY($2::text [])
171 AND deleted_at IS NULL
172ORDER BY sisu_person_id,
173 sent_at DESC
174 "#,
175 course_id,
176 sisu_person_ids,
177 )
178 .fetch_all(conn)
179 .await?;
180 Ok(res
181 .into_iter()
182 .map(|row| (row.sisu_person_id.expose_secret().to_owned(), row))
183 .collect())
184}
185
186pub async fn get_by_course_id_and_student_number(
191 conn: &mut PgConnection,
192 course_id: Uuid,
193 student_number: &str,
194) -> ModelResult<Vec<CreditRegistrationAccountLinkingEmail>> {
195 let res = sqlx::query_as!(
196 CreditRegistrationAccountLinkingEmail,
197 r#"
198SELECT *
199FROM credit_registration_account_linking_emails
200WHERE course_id = $1
201 AND student_number = $2
202 AND deleted_at IS NULL
203ORDER BY sent_at DESC
204 "#,
205 course_id,
206 student_number,
207 )
208 .fetch_all(conn)
209 .await?;
210 Ok(res)
211}
212
213pub async fn get_by_sisu_person_id(
214 conn: &mut PgConnection,
215 sisu_person_id: &str,
216) -> ModelResult<Vec<CreditRegistrationAccountLinkingEmail>> {
217 let res = sqlx::query_as!(
218 CreditRegistrationAccountLinkingEmail,
219 r#"
220SELECT *
221FROM credit_registration_account_linking_emails
222WHERE sisu_person_id = $1
223 AND deleted_at IS NULL
224ORDER BY sent_at DESC
225 "#,
226 sisu_person_id
227 )
228 .fetch_all(conn)
229 .await?;
230 Ok(res)
231}
232
233pub async fn count_sent_for_person_and_course(
236 conn: &mut PgConnection,
237 sisu_person_id: &str,
238 course_id: Uuid,
239) -> ModelResult<i64> {
240 let count = sqlx::query_scalar!(
241 r#"
242SELECT COUNT(*) AS "count!"
243FROM credit_registration_account_linking_emails
244WHERE sisu_person_id = $1
245 AND course_id = $2
246 AND deleted_at IS NULL
247 "#,
248 sisu_person_id,
249 course_id,
250 )
251 .fetch_one(conn)
252 .await?;
253 Ok(count)
254}
255
256#[derive(Debug, Clone)]
258pub struct LinkingMailToQueue {
259 pub id: Uuid,
260 pub emailed_to: DbSecret,
261 pub student_number: DbSecret,
262 pub first_names: Option<DbSecret>,
263 pub token: DbSecret,
265 pub course_name: String,
266 pub course_language_code: String,
267}
268
269pub async fn claim_unqueued(
278 conn: &mut PgConnection,
279 limit: i64,
280 course_id: Option<Uuid>,
281) -> ModelResult<Vec<LinkingMailToQueue>> {
282 let res = sqlx::query_as!(
283 LinkingMailToQueue,
284 r#"
285SELECT e.id AS "id!",
286 e.emailed_to AS "emailed_to!",
287 e.student_number AS "student_number!",
288 t.first_names AS "first_names?",
289 t.token AS "token!: DbSecret",
290 c.name AS "course_name!",
291 c.language_code AS "course_language_code!"
292FROM credit_registration_account_linking_emails e
293 JOIN student_number_verification_tokens t ON t.id = e.student_number_verification_token_id
294 JOIN courses c ON c.id = e.course_id AND c.deleted_at IS NULL
295WHERE e.email_delivery_id IS NULL
296 AND e.deleted_at IS NULL
297 AND t.deleted_at IS NULL
298 AND t.used_at IS NULL
299 AND t.expires_at > now()
300 AND ($2::uuid IS NULL OR e.course_id = $2)
301ORDER BY e.sent_at
302FOR UPDATE OF e SKIP LOCKED
303LIMIT $1
304 "#,
305 limit,
306 course_id,
307 )
308 .fetch_all(conn)
309 .await?;
310 Ok(res)
311}
312
313pub async fn set_email_delivery_id(
315 conn: &mut PgConnection,
316 id: Uuid,
317 email_delivery_id: Uuid,
318) -> ModelResult<()> {
319 sqlx::query!(
320 r#"
321UPDATE credit_registration_account_linking_emails
322SET email_delivery_id = $2
323WHERE id = $1
324 AND deleted_at IS NULL
325 "#,
326 id,
327 email_delivery_id,
328 )
329 .execute(conn)
330 .await?;
331 Ok(())
332}
333
334pub async fn get_send_status_reports(
337 conn: &mut PgConnection,
338 ids: &[Uuid],
339) -> ModelResult<HashMap<Uuid, EmailSendStatusReport>> {
340 let rows = sqlx::query!(
341 r#"
342SELECT id,
343 email_delivery_id
344FROM credit_registration_account_linking_emails
345WHERE id = ANY($1::uuid [])
346 AND deleted_at IS NULL
347 "#,
348 ids
349 )
350 .fetch_all(&mut *conn)
351 .await?;
352 let delivery_ids: Vec<Uuid> = rows
353 .iter()
354 .filter_map(|row| row.email_delivery_id)
355 .collect();
356 let mut deliveries = get_send_statuses(conn, &delivery_ids).await?;
357 let res = rows
358 .into_iter()
359 .map(|row| {
360 let report = row
361 .email_delivery_id
362 .and_then(|id| deliveries.remove(&id))
363 .unwrap_or_else(not_handed_over_yet);
364 (row.id, report)
365 })
366 .collect();
367 Ok(res)
368}
369
370pub fn not_handed_over_yet() -> EmailSendStatusReport {
373 EmailSendStatusReport {
374 email_send_status: EmailSendStatus::Queued,
375 sent_at: None,
376 last_attempt_at: None,
377 retry_count: 0,
378 next_retry_at: None,
379 failure_code: None,
380 failure_is_transient: None,
381 }
382}
383
384pub async fn count_sent_since(conn: &mut PgConnection, since: DateTime<Utc>) -> ModelResult<i64> {
387 let count = sqlx::query_scalar!(
388 r#"
389SELECT COUNT(*) AS "count!"
390FROM credit_registration_account_linking_emails
391WHERE sent_at >= $1
392 AND deleted_at IS NULL
393 "#,
394 since,
395 )
396 .fetch_one(conn)
397 .await?;
398 Ok(count)
399}
400
401#[derive(Debug, Clone, PartialEq, Default)]
405pub struct LinkingMailSendStatusTotals {
406 pub mails_in_window: i64,
407 pub queued: i64,
408 pub retrying: i64,
409 pub sent: i64,
410 pub send_failed: i64,
411 pub last_send_failed_at: Option<DateTime<Utc>>,
413}
414
415pub async fn get_send_status_totals_since(
416 conn: &mut PgConnection,
417 since: DateTime<Utc>,
418 now: DateTime<Utc>,
419) -> ModelResult<LinkingMailSendStatusTotals> {
420 struct Row {
421 sent_at: DateTime<Utc>,
422 email_delivery_id: Option<Uuid>,
423 delivery_sent: Option<bool>,
424 retryable: Option<bool>,
425 first_failed_at: Option<DateTime<Utc>>,
426 retry_count: Option<i32>,
427 }
428 let rows = sqlx::query_as!(
429 Row,
430 r#"
431SELECT
432 e.sent_at AS "sent_at!",
433 e.email_delivery_id,
434 ed.sent AS "delivery_sent?",
435 ed.retryable AS "retryable?",
436 ed.first_failed_at AS "first_failed_at?",
437 ed.retry_count AS "retry_count?"
438FROM credit_registration_account_linking_emails e
439 LEFT JOIN email_deliveries ed ON ed.id = e.email_delivery_id AND ed.deleted_at IS NULL
440WHERE e.sent_at >= $1
441 AND e.deleted_at IS NULL
442 "#,
443 since,
444 )
445 .fetch_all(conn)
446 .await?;
447
448 let mut totals = LinkingMailSendStatusTotals {
449 mails_in_window: rows.len() as i64,
450 ..Default::default()
451 };
452 for row in rows {
453 let status = match row.email_delivery_id {
454 None => EmailSendStatus::Queued,
455 Some(_) => {
456 let facts = EmailSendStatusFacts {
457 sent: row.delivery_sent.unwrap_or(false),
458 retryable: row.retryable.unwrap_or(false),
459 retry_count: row.retry_count.unwrap_or(0),
460 next_retry_at: None,
461 first_failed_at: row.first_failed_at,
462 last_attempt_at: None,
463 failure_code: None,
464 failure_is_transient: None,
465 };
466 derive_email_send_status(&facts, now).email_send_status
467 }
468 };
469 match status {
470 EmailSendStatus::Queued => totals.queued += 1,
471 EmailSendStatus::Retrying => totals.retrying += 1,
472 EmailSendStatus::Sent => totals.sent += 1,
473 EmailSendStatus::SendFailed => {
474 totals.send_failed += 1;
475 totals.last_send_failed_at = Some(
476 totals
477 .last_send_failed_at
478 .map_or(row.sent_at, |prev| prev.max(row.sent_at)),
479 );
480 }
481 }
482 }
483 Ok(totals)
484}
485
486#[derive(Debug, Clone, PartialEq)]
487pub struct LinkingMailFailureDomain {
488 pub domain: String,
489 pub count: i64,
490}
491
492pub async fn get_send_failure_domains_since(
496 conn: &mut PgConnection,
497 since: DateTime<Utc>,
498 now: DateTime<Utc>,
499) -> ModelResult<Vec<LinkingMailFailureDomain>> {
500 struct Row {
501 emailed_to: DbSecret,
502 retryable: bool,
503 first_failed_at: Option<DateTime<Utc>>,
504 }
505 let rows = sqlx::query_as!(
506 Row,
507 r#"
508SELECT e.emailed_to, ed.retryable, ed.first_failed_at
509FROM credit_registration_account_linking_emails e
510 JOIN email_deliveries ed ON ed.id = e.email_delivery_id AND ed.deleted_at IS NULL
511WHERE e.sent_at >= $1
512 AND e.deleted_at IS NULL
513 AND position('@' IN e.emailed_to) > 0
514 AND NOT ed.sent
515 "#,
516 since,
517 )
518 .fetch_all(conn)
519 .await?;
520
521 let mut counts: HashMap<String, i64> = HashMap::new();
522 for row in rows {
523 if !is_hard_send_failure(row.retryable, row.first_failed_at, now) {
524 continue;
525 }
526 let emailed_to = row.emailed_to.expose_secret();
527 let Some(at) = emailed_to.find('@') else {
528 continue;
529 };
530 *counts.entry(emailed_to[at + 1..].to_string()).or_insert(0) += 1;
531 }
532
533 let mut domains: Vec<LinkingMailFailureDomain> = counts
534 .into_iter()
535 .map(|(domain, count)| LinkingMailFailureDomain { domain, count })
536 .collect();
537 domains.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.domain.cmp(&b.domain)));
538 Ok(domains)
539}
540
541pub async fn count_send_failed_for_course(
544 conn: &mut PgConnection,
545 course_id: Uuid,
546 now: DateTime<Utc>,
547) -> ModelResult<i64> {
548 struct Row {
549 retryable: bool,
550 first_failed_at: Option<DateTime<Utc>>,
551 }
552 let rows = sqlx::query_as!(
553 Row,
554 r#"
555SELECT ed.retryable, ed.first_failed_at
556FROM credit_registration_account_linking_emails e
557 JOIN email_deliveries ed ON ed.id = e.email_delivery_id AND ed.deleted_at IS NULL
558WHERE e.course_id = $1
559 AND e.deleted_at IS NULL
560 AND NOT ed.sent
561 "#,
562 course_id,
563 )
564 .fetch_all(conn)
565 .await?;
566 Ok(rows
567 .into_iter()
568 .filter(|row| is_hard_send_failure(row.retryable, row.first_failed_at, now))
569 .count() as i64)
570}
571
572#[derive(Debug, Clone)]
575pub struct StaleUnclaimedLinkingMails {
576 pub student_number: DbSecret,
577 pub sisu_person_id: DbSecret,
578 pub course_id: Uuid,
579 pub course_name: String,
580 pub mail_count: i64,
581 pub first_sent_at: DateTime<Utc>,
582 pub last_sent_at: DateTime<Utc>,
583 pub mail_ids: Vec<Uuid>,
584 pub addresses: Vec<DbSecret>,
586}
587
588pub async fn get_stale_unclaimed(
589 conn: &mut PgConnection,
590 min_mail_count: i64,
591 limit: i64,
592) -> ModelResult<Vec<StaleUnclaimedLinkingMails>> {
593 let res = sqlx::query_as!(
594 StaleUnclaimedLinkingMails,
595 r#"
596SELECT e.student_number AS "student_number!",
597 e.sisu_person_id AS "sisu_person_id!",
598 e.course_id AS "course_id!",
599 c.name AS "course_name!",
600 COUNT(*) AS "mail_count!",
601 MIN(e.sent_at) AS "first_sent_at!",
602 MAX(e.sent_at) AS "last_sent_at!",
603 ARRAY_AGG(
604 e.id
605 ORDER BY e.sent_at
606 ) AS "mail_ids!",
607 ARRAY_AGG(
608 e.emailed_to
609 ORDER BY e.sent_at
610 ) AS "addresses!: Vec<DbSecret>"
611FROM credit_registration_account_linking_emails e
612 JOIN courses c ON c.id = e.course_id AND c.deleted_at IS NULL
613WHERE e.deleted_at IS NULL
614 AND NOT EXISTS (
615 SELECT 1
616 FROM verified_student_numbers vsn
617 WHERE vsn.sisu_person_id = e.sisu_person_id
618 AND vsn.deleted_at IS NULL
619 )
620GROUP BY e.student_number,
621 e.sisu_person_id,
622 e.course_id,
623 c.name
624HAVING COUNT(*) >= $1
625ORDER BY MAX(e.sent_at) DESC
626LIMIT $2
627 "#,
628 min_mail_count,
629 limit,
630 )
631 .fetch_all(conn)
632 .await?;
633 Ok(res)
634}
635
636pub async fn soft_delete(conn: &mut PgConnection, id: Uuid) -> ModelResult<()> {
638 soft_delete_batch(conn, std::slice::from_ref(&id)).await
639}
640
641pub async fn soft_delete_batch(conn: &mut PgConnection, ids: &[Uuid]) -> ModelResult<()> {
643 if ids.is_empty() {
644 return Ok(());
645 }
646 sqlx::query!(
647 r#"
648UPDATE credit_registration_account_linking_emails
649SET deleted_at = now()
650WHERE id = ANY($1)
651 AND deleted_at IS NULL
652 "#,
653 ids
654 )
655 .execute(conn)
656 .await?;
657 Ok(())
658}
659
660#[cfg(test)]
661mod tests {
662 use super::*;
663 use crate::email_deliveries::insert_email_delivery_to_address;
664 use crate::email_templates::{EmailTemplateNew, EmailTemplateType, insert_email_template};
665 use crate::library::credit_registration::account_linking::{
666 DiscoveredPerson, claim_linking_mails,
667 };
668 use crate::test_helper::*;
669
670 async fn claim_a_mail(conn: &mut PgConnection, course_id: Uuid) -> Uuid {
671 claim_linking_mails(
672 conn,
673 &DiscoveredPerson {
674 sisu_person_id: "hy-hlo-1".to_string().into(),
675 student_number: "012345678".to_string().into(),
676 first_names: Some("Aada".to_string().into()),
677 last_name: Some("Virtanen".to_string().into()),
678 course_id,
679 addresses: vec![DbSecret::new("aada@example.com")],
680 },
681 )
682 .await
683 .unwrap();
684 get_by_sisu_person_id(conn, "hy-hlo-1")
685 .await
686 .unwrap()
687 .pop()
688 .expect("the claim wrote a slot")
689 .id
690 }
691
692 async fn seed_template(conn: &mut PgConnection) -> Uuid {
693 insert_email_template(
694 conn,
695 None,
696 EmailTemplateNew {
697 template_type: EmailTemplateType::CreditRegistrationAccountLinking,
698 language: Some("en".to_string()),
699 content: Some(serde_json::json!([])),
700 subject: Some("Link your student number".to_string()),
701 },
702 None,
703 )
704 .await
705 .unwrap()
706 .id
707 }
708
709 #[tokio::test]
712 async fn a_claimed_slot_leaves_the_queue_once_its_delivery_exists() {
713 insert_data!(:tx, :user, :org, :course);
714 let slot = claim_a_mail(tx.as_mut(), course).await;
715 let template = seed_template(tx.as_mut()).await;
716
717 let claimed = claim_unqueued(tx.as_mut(), 10, None).await.unwrap();
718 assert_eq!(claimed.len(), 1);
719 assert_eq!(claimed[0].id, slot);
720 assert_eq!(claimed[0].emailed_to.expose_secret(), "aada@example.com");
721
722 let delivery = insert_email_delivery_to_address(
723 tx.as_mut(),
724 claimed[0].emailed_to.expose_secret(),
725 template,
726 &serde_json::json!({ "NAME": "Aada" }),
727 )
728 .await
729 .unwrap();
730 set_email_delivery_id(tx.as_mut(), slot, delivery)
731 .await
732 .unwrap();
733
734 assert!(
735 claim_unqueued(tx.as_mut(), 10, None)
736 .await
737 .unwrap()
738 .is_empty()
739 );
740 }
741
742 #[tokio::test]
743 async fn a_slot_with_no_delivery_yet_reports_as_queued() {
744 insert_data!(:tx, :user, :org, :course);
745 let slot = claim_a_mail(tx.as_mut(), course).await;
746 let reports = get_send_status_reports(tx.as_mut(), &[slot]).await.unwrap();
747 assert_eq!(
748 reports.get(&slot).map(|report| report.email_send_status),
749 Some(EmailSendStatus::Queued)
750 );
751 }
752}