headless_lms_server/programs/
sync_tmc_users.rs1use 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 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 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 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 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 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}