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