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::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 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 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 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 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 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}