Skip to main content

headless_lms_server/programs/
sync_tmc_users.rs

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