Skip to main content

headless_lms_server/programs/
sync_tmc_users.rs

1/*!
2Syncs tmc users
3*/
4use std::env;
5
6use crate::config::program_config::ProgramConfig;
7use crate::domain::email_ownership_verification::queue_verification_email_best_effort;
8use crate::setup_tracing;
9use anyhow::Context;
10
11use chrono::DateTime;
12use dotenvy::dotenv;
13use headless_lms_models as models;
14use models::users::{get_users_ids_in_db_from_upstream_ids, update_email_for_user};
15
16use serde::{Deserialize, Serialize};
17use sqlx::{PgConnection, PgPool};
18
19const URL: &str = "https://tmc.mooc.fi/api/v8/users/recently_changed_user_details";
20
21#[derive(Debug, Serialize, Deserialize)]
22pub struct TMCRecentChanges {
23    pub changes: Vec<Change>,
24}
25
26#[derive(Debug, Serialize, Deserialize)]
27pub struct Change {
28    pub change_type: String,
29    pub new_value: Option<String>,
30    pub old_value: Option<String>,
31    pub created_at: String,
32    pub id: i32,
33    pub user_id: Option<i32>,
34}
35
36pub async fn main() -> anyhow::Result<()> {
37    // TODO: Audit that the environment access only happens in single-threaded code.
38    unsafe { env::set_var("RUST_LOG", "info,actix_web=info,sqlx=warn") };
39    dotenv().ok();
40    setup_tracing()?;
41    let database_url = ProgramConfig::database_url_with_default();
42    // The same variable the server reads into
43    // `ApplicationConfiguration::enable_email_ownership_verification`.
44    let email_ownership_verification_enabled =
45        ProgramConfig::bool_flag("ENABLE_EMAIL_OWNERSHIP_VERIFICATION");
46    let recent_changes = fetch_recently_changed_user_details().await?;
47    let db_pool = PgPool::connect(&database_url).await?;
48    let mut conn = db_pool.acquire().await?;
49    delete_users(&mut conn, &recent_changes).await?;
50    update_users(
51        &mut conn,
52        email_ownership_verification_enabled,
53        &recent_changes,
54    )
55    .await?;
56    Ok(())
57}
58
59pub async fn update_users(
60    conn: &mut PgConnection,
61    email_ownership_verification_enabled: bool,
62    recent_changes: &TMCRecentChanges,
63) -> anyhow::Result<()> {
64    let email_update_list = recent_changes
65        .changes
66        .iter()
67        .filter(|c| c.change_type == "email_changed")
68        .collect::<Vec<_>>();
69
70    info!("Updating emails for {} users", email_update_list.len());
71
72    // Parse dates and filter out invalid entries
73    let mut parsed_updates: Vec<(DateTime<chrono::FixedOffset>, &Change)> = Vec::new();
74    for change in &email_update_list {
75        match DateTime::parse_from_rfc3339(change.created_at.as_str()) {
76            Ok(date) => parsed_updates.push((date, change)),
77            Err(e) => {
78                error!("Error converting date: '{}'", change.created_at);
79                error!("Error: {}", e);
80                return Err(anyhow::anyhow!(
81                    "Error converting date: '{}'",
82                    change.created_at
83                ));
84            }
85        }
86    }
87
88    parsed_updates.sort_by_key(|a| a.0);
89
90    let email_update_list: Vec<&Change> = parsed_updates
91        .into_iter()
92        .map(|(_, change)| change)
93        .collect();
94
95    for change in email_update_list {
96        if let Some(user_id) = change.user_id {
97            // `user_details.email` is CHECKed to contain an '@', so a change carrying no new value
98            // cannot be applied at all.
99            let Some(new_email) = change.new_value.as_deref() else {
100                error!(
101                    "TMC email change {} for user {user_id} carries no new value",
102                    change.id
103                );
104                continue;
105            };
106            match update_email_for_user(&mut *conn, &user_id, new_email.to_string()).await {
107                Ok(changed_user_id) => {
108                    // The `clear_email_verification` trigger just dropped the old address's proof.
109                    queue_verification_email_best_effort(
110                        &mut *conn,
111                        email_ownership_verification_enabled,
112                        changed_user_id,
113                    )
114                    .await;
115                }
116                Err(e) => {
117                    error!("Error updating user with id {}", user_id);
118                    error!("Error: {}", e);
119                }
120            };
121        };
122    }
123
124    info!("Update done");
125    Ok(())
126}
127
128pub async fn delete_users(
129    conn: &mut PgConnection,
130    recent_changes: &TMCRecentChanges,
131) -> anyhow::Result<()> {
132    let to_delete = recent_changes
133        .changes
134        .iter()
135        .filter(|c| c.change_type == "deleted")
136        .filter_map(|c| c.user_id)
137        .collect::<Vec<_>>();
138    info!("Making sure {} users are deleted", to_delete.len());
139    let user_ids_in_db = get_users_ids_in_db_from_upstream_ids(&mut *conn, &to_delete).await?;
140    info!("{} users need to be deleted", to_delete.len());
141    for id in user_ids_in_db {
142        models::users::delete_user(&mut *conn, id).await?;
143    }
144    info!("Deletions done");
145    Ok(())
146}
147
148pub async fn fetch_recently_changed_user_details() -> anyhow::Result<TMCRecentChanges> {
149    let access_token = ProgramConfig::required("TMC_ACCESS_TOKEN")?;
150    let ratelimit_api_key = ProgramConfig::required("RATELIMIT_PROTECTION_SAFE_API_KEY")?;
151    let client = reqwest::Client::new();
152    let res = client
153        .get(URL)
154        .header("RATELIMIT-PROTECTION-SAFE-API-KEY", ratelimit_api_key)
155        .header(reqwest::header::CONTENT_TYPE, "application/json")
156        .header(reqwest::header::ACCEPT, "application/json")
157        .bearer_auth(&access_token)
158        .send()
159        .await
160        .context("Failed to send request to https://tmc.mooc.fi")?;
161    if res.status().is_success() {
162        let res: TMCRecentChanges = res.json().await?;
163        info!("Fetched {} changes", res.changes.len());
164        Ok(res)
165    } else {
166        let response_body = res.bytes().await?.to_vec();
167        let response_body_string = String::from_utf8_lossy(&response_body);
168        error!(
169            ?response_body_string,
170            "Failed to fetch recently changed user details",
171        );
172        Err(anyhow::anyhow!(
173            "Failed to get recently changed user details from TMC"
174        ))
175    }
176}