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