1use crate::config::program_config::ProgramConfig;
2use crate::prelude::*;
3use crate::setup_tracing;
4use dotenvy::dotenv;
5use headless_lms_models::marketing_consents::MarketingMailingListAccessToken;
6use headless_lms_models::marketing_consents::UserEmailSubscription;
7use headless_lms_models::marketing_consents::UserMarketingConsentWithDetails;
8use headless_lms_utils::http::REQWEST_CLIENT;
9use secrecy::ExposeSecret;
10use serde_json::json;
11use sqlx::{PgConnection, PgPool};
12use std::time::{Duration, Instant};
13use uuid::Uuid;
14
15mod batch_client;
16mod mailchimp_ops;
17mod policy_tags;
18
19use batch_client::MAX_MAILCHIMP_BATCH_SIZE;
20use mailchimp_ops::{MailchimpExecutor, MailchimpOperation};
21use policy_tags::sync_policy_tags_for_users;
22use reqwest::Method;
23
24#[derive(Debug, Deserialize)]
25struct MailchimpField {
26 field_id: String,
27 field_name: String,
28}
29
30#[derive(Debug)]
31struct FieldSchema {
32 tag: &'static str,
33 name: &'static str,
34 default_value: &'static str,
35}
36
37const REQUIRED_FIELDS: &[FieldSchema] = &[
38 FieldSchema {
39 tag: "FNAME",
40 name: "First Name",
41 default_value: "",
42 },
43 FieldSchema {
44 tag: "LNAME",
45 name: "Last Name",
46 default_value: "",
47 },
48 FieldSchema {
49 tag: "MARKETING",
50 name: "Accepts Marketing",
51 default_value: "disallowed",
52 },
53 FieldSchema {
54 tag: "LOCALE",
55 name: "Locale",
56 default_value: "en",
57 },
58 FieldSchema {
59 tag: "GRADUATED",
60 name: "Graduated",
61 default_value: "",
62 },
63 FieldSchema {
64 tag: "COURSEID",
65 name: "Course ID",
66 default_value: "",
67 },
68 FieldSchema {
69 tag: "LANGGRPID",
70 name: "Course language Group ID",
71 default_value: "",
72 },
73 FieldSchema {
74 tag: "USERID",
75 name: "User ID",
76 default_value: "",
77 },
78 FieldSchema {
79 tag: "RESEARCH",
80 name: "Research consent",
81 default_value: "false",
82 },
83];
84
85const FIELDS_EXCLUDED_FROM_REMOVING: &[&str] = &["PHONE", "PACE", "COUNTRY", "MMERGE9"];
87const REMOVE_UNSUPPORTED_FIELDS: bool = false;
88const PROCESS_UNSUBSCRIBES_INTERVAL_SECS: u64 = 10_800;
89
90const SYNC_INTERVAL_SECS: u64 = 10;
91const PRINT_STILL_RUNNING_MESSAGE_TICKS_THRESHOLD: u32 = 60;
92
93const BATCH_POLL_INTERVAL_SECS: u64 = 10;
94const BATCH_POLL_TIMEOUT_SECS: u64 = 300;
95const BATCH_RESULT_DOWNLOAD_TIMEOUT_SECS: u64 = 600;
98
99#[derive(Debug)]
100struct SyncUser {
101 details: UserMarketingConsentWithDetails,
102 payload: serde_json::Value,
103}
104
105#[derive(Debug)]
106struct EmailSyncResult {
107 user_id: Uuid,
108 user_mailchimp_id: String,
109}
110
111pub async fn main() -> anyhow::Result<()> {
113 initialize_environment()?;
114
115 let config = initialize_configuration().await?;
116
117 let db_pool = initialize_database_pool(&config.database_url).await?;
118 let mut conn = db_pool.acquire().await?;
119
120 let mut interval = tokio::time::interval(Duration::from_secs(SYNC_INTERVAL_SECS));
121 let mut ticks = 0;
122
123 let access_tokens =
124 headless_lms_models::marketing_consents::fetch_all_marketing_mailing_list_access_tokens(
125 &mut conn,
126 )
127 .await?;
128
129 for token in &access_tokens {
131 if let Err(e) = ensure_mailchimp_schema(
132 &token.mailchimp_mailing_list_id,
133 &token.server_prefix,
134 &token.access_token,
135 )
136 .await
137 {
138 error!(
139 "Failed to set up Mailchimp schema for list '{}': {:?}",
140 token.mailchimp_mailing_list_id, e
141 );
142 return Err(e);
143 }
144 }
145
146 info!("Starting mailchimp syncer (periodic reconciliation loop).");
147
148 let mut last_time_unsubscribes_processed = Instant::now();
149 let mut last_time_tags_synced = Instant::now();
150
151 loop {
152 interval.tick().await;
153 ticks += 1;
154
155 if ticks >= PRINT_STILL_RUNNING_MESSAGE_TICKS_THRESHOLD {
156 ticks = 0;
157 info!("Still syncing.");
158 }
159 let mut process_unsubscribes = false;
160 if last_time_unsubscribes_processed.elapsed().as_secs()
161 >= PROCESS_UNSUBSCRIBES_INTERVAL_SECS
162 {
163 process_unsubscribes = true;
164 last_time_unsubscribes_processed = Instant::now();
165 };
166
167 if last_time_tags_synced.elapsed().as_secs() >= 3600 {
169 info!("Stage: tag catalog sync (keep local tag metadata aligned with Mailchimp).");
170 for token in &access_tokens {
171 if let Err(e) = sync_tags_from_mailchimp(
172 &mut conn,
173 &token.mailchimp_mailing_list_id,
174 &token.access_token,
175 &token.server_prefix,
176 token.id,
177 token.course_language_group_id,
178 )
179 .await
180 {
181 error!(
182 "Failed to sync tags for list '{}': {:?}",
183 token.mailchimp_mailing_list_id, e
184 );
185 }
186 }
187 last_time_tags_synced = Instant::now();
188 }
189
190 if let Err(e) = sync_contacts(&mut conn, &config, process_unsubscribes).await {
191 error!("Error during synchronization: {:?}", e);
192 if let Ok(sqlx::Error::Io(..)) = e.downcast::<sqlx::Error>() {
193 info!("syncer may have lost its connection to the db, trying to reconnect");
195 conn = db_pool.acquire().await?;
196 }
197 }
198 }
199}
200
201fn initialize_environment() -> anyhow::Result<()> {
203 dotenv().ok();
204 ProgramConfig::ensure_default_rust_log_for_workers();
205 setup_tracing()?;
206 Ok(())
207}
208
209struct SyncerConfig {
211 database_url: String,
212}
213
214async fn initialize_configuration() -> anyhow::Result<SyncerConfig> {
216 let database_url = ProgramConfig::database_url_with_default();
217
218 Ok(SyncerConfig { database_url })
219}
220
221async fn initialize_database_pool(database_url: &str) -> anyhow::Result<PgPool> {
223 PgPool::connect(database_url).await.map_err(|e| {
224 anyhow::anyhow!(
225 "Failed to connect to the database at {}: {:?}",
226 database_url,
227 e
228 )
229 })
230}
231
232async fn ensure_mailchimp_schema(
234 list_id: &str,
235 server_prefix: &str,
236 access_token: &DbSecret,
237) -> anyhow::Result<()> {
238 let existing_fields =
239 fetch_current_mailchimp_fields(list_id, server_prefix, access_token).await?;
240
241 if REMOVE_UNSUPPORTED_FIELDS {
242 for field in existing_fields.iter() {
244 if !REQUIRED_FIELDS
245 .iter()
246 .any(|r| r.tag == field.field_name.as_str())
247 && !FIELDS_EXCLUDED_FROM_REMOVING.contains(&field.field_name.as_str())
248 {
249 match remove_field_from_mailchimp(
250 list_id,
251 &field.field_id,
252 server_prefix,
253 access_token,
254 )
255 .await
256 {
257 Err(e) => {
258 warn!("Could not remove field '{}': {}", field.field_name, e);
259 }
260 _ => {
261 info!("Removed field '{}'", field.field_name);
262 }
263 }
264 }
265 }
266 }
267
268 for required_field in REQUIRED_FIELDS.iter() {
270 if !existing_fields
271 .iter()
272 .any(|f| f.field_name == required_field.tag)
273 {
274 match add_field_to_mailchimp(list_id, required_field, server_prefix, access_token).await
275 {
276 Err(e) => {
277 warn!(
278 "Failed to add required field '{}': {}",
279 required_field.name, e
280 );
281 }
282 _ => {
283 info!(
284 "Successfully added required field '{}'",
285 required_field.name
286 );
287 }
288 }
289 } else {
290 info!(
291 "Field '{}' already exists, skipping addition.",
292 required_field.name
293 );
294 }
295 }
296
297 Ok(())
298}
299
300async fn fetch_current_mailchimp_fields(
302 list_id: &str,
303 server_prefix: &str,
304 access_token: &DbSecret,
305) -> Result<Vec<MailchimpField>, anyhow::Error> {
306 let url = format!(
307 "https://{}.api.mailchimp.com/3.0/lists/{}/merge-fields",
308 server_prefix, list_id
309 );
310
311 let response = REQWEST_CLIENT
312 .get(&url)
313 .header(
314 "Authorization",
315 format!("apikey {}", access_token.expose_secret()),
316 )
317 .send()
318 .await?;
319
320 if response.status().is_success() {
321 let json = response.json::<serde_json::Value>().await?;
322
323 let fields: Vec<MailchimpField> = json["merge_fields"]
324 .as_array()
325 .unwrap_or(&vec![])
326 .iter()
327 .filter_map(|field| {
328 let field_id = field["merge_id"].as_u64();
329 let field_name = field["tag"].as_str();
330
331 if let (Some(field_id), Some(field_name)) = (field_id, field_name) {
332 Some(MailchimpField {
333 field_id: field_id.to_string(),
334 field_name: field_name.to_string(),
335 })
336 } else {
337 None
338 }
339 })
340 .collect();
341
342 Ok(fields)
343 } else {
344 let error_text = response.text().await?;
345 error!("Error fetching merge fields: {}", error_text);
346 Err(anyhow::anyhow!("Failed to fetch current Mailchimp fields."))
347 }
348}
349
350async fn add_field_to_mailchimp(
352 list_id: &str,
353 field_schema: &FieldSchema,
354 server_prefix: &str,
355 access_token: &DbSecret,
356) -> anyhow::Result<()> {
357 let url = format!(
358 "https://{}.api.mailchimp.com/3.0/lists/{}/merge-fields",
359 server_prefix, list_id
360 );
361
362 let body = json!({
363 "tag": field_schema.tag,
364 "name": field_schema.name,
365 "type": "text",
366 "default_value": field_schema.default_value,
367 });
368
369 let response = REQWEST_CLIENT
370 .post(&url)
371 .header(
372 "Authorization",
373 format!("apikey {}", access_token.expose_secret()),
374 )
375 .json(&body)
376 .send()
377 .await?;
378
379 if response.status().is_success() {
380 Ok(())
381 } else {
382 let status = response.status();
383 let error_text = response
384 .text()
385 .await
386 .unwrap_or_else(|_| "No additional error info.".to_string());
387 Err(anyhow::anyhow!(
388 "Failed to add field to Mailchimp. Status: {}. Error: {}",
389 status,
390 error_text
391 ))
392 }
393}
394
395async fn remove_field_from_mailchimp(
397 list_id: &str,
398 field_id: &str,
399 server_prefix: &str,
400 access_token: &DbSecret,
401) -> anyhow::Result<()> {
402 let url = format!(
403 "https://{}.api.mailchimp.com/3.0/lists/{}/merge-fields/{}",
404 server_prefix, list_id, field_id
405 );
406
407 let response = REQWEST_CLIENT
408 .delete(&url)
409 .header(
410 "Authorization",
411 format!("apikey {}", access_token.expose_secret()),
412 )
413 .send()
414 .await?;
415
416 if response.status().is_success() {
417 Ok(())
418 } else {
419 let status = response.status();
420 let error_text = response
421 .text()
422 .await
423 .unwrap_or_else(|_| "No additional error info.".to_string());
424 Err(anyhow::anyhow!(
425 "Failed to remove field from Mailchimp. Status: {}. Error: {}",
426 status,
427 error_text
428 ))
429 }
430}
431
432pub async fn sync_tags_from_mailchimp(
434 conn: &mut PgConnection,
435 list_id: &str,
436 access_token: &DbSecret,
437 server_prefix: &str,
438 marketing_mailing_list_access_token_id: Uuid,
439 course_language_group_id: Uuid,
440) -> anyhow::Result<()> {
441 let url = format!(
442 "https://{}.api.mailchimp.com/3.0/lists/{}/tag-search",
443 server_prefix, list_id
444 );
445
446 let response = REQWEST_CLIENT
447 .get(&url)
448 .header(
449 "Authorization",
450 format!("apikey {}", access_token.expose_secret()),
451 )
452 .send()
453 .await?;
454
455 let response_json = response.json::<serde_json::Value>().await?;
457
458 let mailchimp_tags = match response_json.get("tags") {
459 Some(tags) if tags.is_array() => tags
460 .as_array()
461 .ok_or_else(|| anyhow::anyhow!("tags field is not an array despite is_array() check"))?
462 .iter()
463 .filter_map(|tag| {
464 let name = tag.get("name")?.as_str()?.to_string();
465
466 let id = match tag.get("id") {
467 Some(serde_json::Value::Number(num)) => num.to_string(),
468 Some(serde_json::Value::String(str_id)) => str_id.clone(),
469 _ => return None,
470 };
471
472 Some((id, name))
473 })
474 .collect::<Vec<(String, String)>>(),
475 _ => {
476 warn!("No tags found for list '{}', skipping sync.", list_id);
477 return Ok(());
478 }
479 };
480
481 let db_tags = headless_lms_models::marketing_consents::fetch_tags_with_course_language_group_id_and_marketing_mailing_list_access_token_id(
483 conn,
484 course_language_group_id,
485 marketing_mailing_list_access_token_id,
486 )
487 .await?;
488
489 for (tag_id, tag_name) in &mailchimp_tags {
491 if let Some(db_tag) = db_tags.iter().find(|db_tag| {
492 db_tag.get("id").and_then(|v| v.as_str()).map(|v| v.trim()) == Some(tag_id.trim())
493 }) {
494 let db_tag_name = db_tag
495 .get("tag_name")
496 .and_then(|v| v.as_str())
497 .unwrap_or_default();
498 if db_tag_name != tag_name {
499 headless_lms_models::marketing_consents::upsert_tag(
500 conn,
501 course_language_group_id,
502 marketing_mailing_list_access_token_id,
503 tag_id.clone(),
504 tag_name.clone(),
505 )
506 .await?;
507 }
508 }
509 }
510
511 for db_tag in db_tags.iter() {
513 let db_tag_id = db_tag
514 .get("id")
515 .and_then(|v| v.as_str())
516 .map(|v| v.trim().to_string())
517 .unwrap_or_default();
518
519 if !mailchimp_tags
521 .iter()
522 .any(|(tag_id, _)| tag_id.trim() == db_tag_id.trim())
523 {
524 if db_tag_id.is_empty() {
525 warn!("Skipping tag deletion due to missing ID: {:?}", db_tag);
526 continue;
527 }
528 headless_lms_models::marketing_consents::delete_tag(
529 conn,
530 db_tag_id.clone(),
531 course_language_group_id,
532 )
533 .await?;
534 }
535 }
536
537 Ok(())
538}
539
540async fn sync_contacts(
543 conn: &mut PgConnection,
544 _config: &SyncerConfig,
545 process_unsubscribes: bool,
546) -> anyhow::Result<()> {
547 let access_tokens =
548 headless_lms_models::marketing_consents::fetch_all_marketing_mailing_list_access_tokens(
549 conn,
550 )
551 .await?;
552
553 let mut successfully_synced_user_ids = Vec::new();
554
555 for token in access_tokens {
557 let course_language_group_slug =
558 match headless_lms_models::course_language_groups::get_slug_by_id(
559 conn,
560 token.course_language_group_id,
561 )
562 .await
563 {
564 Ok(Some(s)) => Some(s),
565 Ok(None) => None,
566 Err(e) => {
567 error!(
568 course_language_group_id = %token.course_language_group_id,
569 "Failed to get course language group slug: {:?}",
570 e
571 );
572 return Err(e.into());
573 }
574 };
575
576 if process_unsubscribes {
578 info!(
579 "Stage: unsubscribe sync (apply Mailchimp compliance/unsubscribe changes locally)."
580 );
581 let mailchimp_data = fetch_unsubscribed_users_from_mailchimp_in_chunks(
582 &token.mailchimp_mailing_list_id,
583 &token.server_prefix,
584 &token.access_token,
585 1000,
586 )
587 .await?;
588
589 info!(
590 "Processing Mailchimp data for list: {}",
591 token.mailchimp_mailing_list_id
592 );
593
594 process_unsubscribed_users_from_mailchimp(conn, mailchimp_data).await?;
595 }
596
597 let users_with_unsynced_emails =
599 headless_lms_models::marketing_consents::fetch_all_unsynced_updated_emails(
600 conn,
601 token.course_language_group_id,
602 )
603 .await?;
604
605 info!(
606 "Stage: email updates (ensure member identifiers stay correct). Found {} unsynced user email(s) for course language group: {}",
607 users_with_unsynced_emails.len(),
608 token.course_language_group_id
609 );
610
611 if !users_with_unsynced_emails.is_empty() {
612 let email_sync_results = update_emails_in_mailchimp(
613 users_with_unsynced_emails,
614 &token.mailchimp_mailing_list_id,
615 &token.server_prefix,
616 &token.access_token,
617 )
618 .await?;
619
620 let email_synced_user_ids: Vec<Uuid> =
621 email_sync_results.iter().map(|r| r.user_id).collect();
622 successfully_synced_user_ids.extend(email_synced_user_ids.clone());
623
624 if !email_sync_results.is_empty() {
625 let mailchimp_id_by_user: std::collections::HashMap<Uuid, String> =
626 email_sync_results
627 .iter()
628 .map(|r| (r.user_id, r.user_mailchimp_id.clone()))
629 .collect();
630 let user_details =
631 headless_lms_models::marketing_consents::fetch_user_marketing_consents_with_details_by_user_ids(
632 conn,
633 token.course_language_group_id,
634 &email_synced_user_ids,
635 )
636 .await?;
637 match sync_policy_tags_for_users(
638 &token,
639 &user_details,
640 &mailchimp_id_by_user,
641 Duration::from_secs(BATCH_POLL_TIMEOUT_SECS),
642 Duration::from_secs(BATCH_POLL_INTERVAL_SECS),
643 )
644 .await
645 {
646 Ok(results) => {
647 let (ok, failed): (Vec<_>, Vec<_>) =
648 results.into_iter().partition(|r| r.success);
649 if !failed.is_empty() {
650 let sample: Vec<String> = failed
651 .iter()
652 .take(5)
653 .map(|r| {
654 format!(
655 "{}: {}",
656 r.user_id,
657 r.error.as_deref().unwrap_or("unknown error")
658 )
659 })
660 .collect();
661 warn!(
662 "Policy tag sync after email updates for list '{}' had {} failure(s) out of {}. Sample: {:?}",
663 token.mailchimp_mailing_list_id,
664 failed.len(),
665 ok.len() + failed.len(),
666 sample
667 );
668 }
669 }
670 Err(e) => {
671 error!(
672 "Failed to sync policy tags after email updates for list '{}': {:?}",
673 token.mailchimp_mailing_list_id, e
674 );
675 }
676 }
677 }
678 }
679
680 let tag_objects = headless_lms_models::marketing_consents::fetch_tags_with_course_language_group_id_and_marketing_mailing_list_access_token_id(conn, token.course_language_group_id, token.id).await?;
681
682 let unsynced_users_details =
684 headless_lms_models::marketing_consents::fetch_all_unsynced_user_marketing_consents_by_course_language_group_id(
685 conn,
686 token.course_language_group_id,
687 )
688 .await?;
689
690 info!(
691 "Stage: member upsert (merge fields + consent). Found {} unsynced user consent(s) for course language group: {}",
692 unsynced_users_details.len(),
693 token.course_language_group_id
694 );
695
696 if !unsynced_users_details.is_empty() {
697 let consent_synced_user_ids =
698 send_users_to_mailchimp(conn, &token, &unsynced_users_details, tag_objects).await?;
699
700 if let Some(ref slug) = course_language_group_slug
701 && !consent_synced_user_ids.is_empty()
702 {
703 let mailchimp_id_mapping =
704 headless_lms_models::marketing_consents::fetch_user_mailchimp_id_mapping(
705 conn,
706 token.course_language_group_id,
707 &consent_synced_user_ids,
708 )
709 .await?;
710 if let Err(e) = sync_completed_tag_for_members(
711 &unsynced_users_details,
712 &consent_synced_user_ids,
713 &mailchimp_id_mapping,
714 slug,
715 &token,
716 )
717 .await
718 {
719 error!(
720 "Failed to sync completed tag for list '{}': {:?}",
721 token.mailchimp_mailing_list_id, e
722 );
723 }
724 }
725
726 successfully_synced_user_ids.extend(consent_synced_user_ids);
728 }
729 }
730
731 if !successfully_synced_user_ids.is_empty() {
733 match headless_lms_models::marketing_consents::update_synced_to_mailchimp_at_to_all_synced_users(
734 conn,
735 &successfully_synced_user_ids,
736 )
737 .await
738 {
739 Ok(_) => {
740 info!(
741 "Stage: mark synced (avoid repeat work). Successfully updated synced status for {} users.",
742 successfully_synced_user_ids.len()
743 );
744 }
745 Err(e) => {
746 error!(
747 "Failed to update synced status for {} users: {:?}",
748 successfully_synced_user_ids.len(),
749 e
750 );
751 }
752 }
753 }
754
755 Ok(())
756}
757
758pub async fn send_users_to_mailchimp(
760 conn: &mut PgConnection,
761 token: &MarketingMailingListAccessToken,
762 users_details: &[UserMarketingConsentWithDetails],
763 tag_objects: Vec<serde_json::Value>,
764) -> anyhow::Result<Vec<Uuid>> {
765 let mut users_to_sync = vec![];
766 let mut sent_user_ids = Vec::new();
767 let mut successfully_synced_user_ids = Vec::new();
768 let mut user_id_contact_id_pairs = Vec::new();
769
770 for user in users_details {
772 if let Some(ref subscription) = user.email_subscription_in_mailchimp
774 && subscription == "subscribed"
775 {
776 sent_user_ids.push(user.user_id);
777 let user_details = json!({
778 "email_address": user.email,
779 "status": user.email_subscription_in_mailchimp,
780 "merge_fields": {
781 "FNAME": user.first_name.clone().unwrap_or("".to_string()),
782 "LNAME": user.last_name.clone().unwrap_or("".to_string()),
783 "MARKETING": if user.consent { "allowed" } else { "disallowed" },
784 "LOCALE": user.locale,
785 "GRADUATED": user.completed_course_at.map(|cca| cca.to_rfc3339()).unwrap_or("".to_string()),
786 "USERID": user.user_id,
787 "COURSEID": user.course_id,
788 "LANGGRPID": user.course_language_group_id,
789 "RESEARCH" : if user.research_consent.unwrap_or(false) { "allowed" } else { "disallowed" },
790 "COUNTRY" : user.country.clone().unwrap_or("".to_string()),
791 },
792 "tags": tag_objects.iter().map(|tag| tag["name"].clone()).collect::<Vec<_>>()
793 });
794 users_to_sync.push(SyncUser {
795 details: user.clone(),
796 payload: user_details,
797 });
798 }
799 }
800
801 if users_to_sync.is_empty() {
802 info!("No new users to sync.");
803 return Ok(vec![]);
804 }
805
806 let url = format!(
807 "https://{}.api.mailchimp.com/3.0/lists/{}",
808 token.server_prefix, token.mailchimp_mailing_list_id
809 );
810
811 let total_chunks = users_to_sync.len().div_ceil(MAX_MAILCHIMP_BATCH_SIZE);
812 info!(
813 "Syncing {} members to list '{}' in {} chunk(s)",
814 users_to_sync.len(),
815 token.mailchimp_mailing_list_id,
816 total_chunks
817 );
818
819 for (chunk_index, chunk) in users_to_sync.chunks(MAX_MAILCHIMP_BATCH_SIZE).enumerate() {
820 info!(
821 "Syncing users chunk {}/{} ({} members)",
822 chunk_index + 1,
823 total_chunks,
824 chunk.len()
825 );
826 let chunk_members: Vec<serde_json::Value> =
827 chunk.iter().map(|user| user.payload.clone()).collect();
828 let batch_request = json!({
829 "members": chunk_members,
830 "update_existing": true
831 });
832
833 let response = REQWEST_CLIENT
834 .post(&url)
835 .header("Content-Type", "application/json")
836 .header(
837 "Authorization",
838 format!("apikey {}", token.access_token.expose_secret()),
839 )
840 .json(&batch_request)
841 .send()
842 .await?;
843
844 if !response.status().is_success() {
845 let status = response.status();
846 let error_text = response.text().await?;
847 return Err(anyhow::anyhow!(
848 "Error syncing users to Mailchimp. Status: {}. Error: {}",
849 status,
850 error_text
851 ));
852 }
853
854 let response_data: serde_json::Value = response.json().await?;
855 let mut chunk_contact_count = 0;
856 let mut chunk_mailchimp_id_by_user: std::collections::HashMap<Uuid, String> =
857 std::collections::HashMap::new();
858 for user in chunk {
859 if let Some(ref mailchimp_id) = user.details.user_mailchimp_id {
860 chunk_mailchimp_id_by_user
861 .entry(user.details.user_id)
862 .or_insert_with(|| mailchimp_id.clone());
863 }
864 }
865 for key in &["new_members", "updated_members"] {
866 if let Some(members) = response_data[key].as_array() {
867 for member in members {
868 if let Some(user_id) = member["merge_fields"]["USERID"].as_str() {
869 if let Ok(uuid) = uuid::Uuid::parse_str(user_id) {
870 successfully_synced_user_ids.push(uuid);
871 }
872 if let Some(contact_id) = member["contact_id"].as_str() {
873 user_id_contact_id_pairs
874 .push((user_id.to_string(), contact_id.to_string()));
875 if let Ok(uuid) = uuid::Uuid::parse_str(user_id) {
876 chunk_mailchimp_id_by_user.insert(uuid, contact_id.to_string());
877 }
878 chunk_contact_count += 1;
879 }
880 }
881 }
882 }
883 }
884 if let Some(errors) = response_data["errors"].as_array()
885 && !errors.is_empty()
886 {
887 let sample: Vec<String> = errors
888 .iter()
889 .take(5)
890 .map(|e| {
891 let email = e
892 .get("email_address")
893 .and_then(|v| v.as_str())
894 .unwrap_or("?");
895 let msg = e.get("error").and_then(|v| v.as_str()).unwrap_or_else(|| {
896 e.get("message").and_then(|v| v.as_str()).unwrap_or("?")
897 });
898 format!("{}: {}", email, msg)
899 })
900 .collect();
901 warn!(
902 "Mailchimp batch subscribe chunk {}/{} returned {} error(s) (e.g. unsubscribed/rejected). Sample: {:?}",
903 chunk_index + 1,
904 total_chunks,
905 errors.len(),
906 sample
907 );
908 }
909 info!(
910 "Chunk {}/{}: {} contact_id(s) from new_members/updated_members",
911 chunk_index + 1,
912 total_chunks,
913 chunk_contact_count
914 );
915
916 let chunk_users: Vec<UserMarketingConsentWithDetails> =
917 chunk.iter().map(|user| user.details.clone()).collect();
918 match sync_policy_tags_for_users(
919 token,
920 &chunk_users,
921 &chunk_mailchimp_id_by_user,
922 Duration::from_secs(BATCH_POLL_TIMEOUT_SECS),
923 Duration::from_secs(BATCH_POLL_INTERVAL_SECS),
924 )
925 .await
926 {
927 Ok(results) => {
928 let (ok, failed): (Vec<_>, Vec<_>) = results.into_iter().partition(|r| r.success);
929 if !failed.is_empty() {
930 let sample: Vec<String> = failed
931 .iter()
932 .take(5)
933 .map(|r| {
934 format!(
935 "{}: {}",
936 r.user_id,
937 r.error.as_deref().unwrap_or("unknown error")
938 )
939 })
940 .collect();
941 warn!(
942 "Policy tag sync for list '{}' chunk {}/{} had {} failure(s) out of {}. Sample: {:?}",
943 token.mailchimp_mailing_list_id,
944 chunk_index + 1,
945 total_chunks,
946 failed.len(),
947 ok.len() + failed.len(),
948 sample
949 );
950 }
951 }
952 Err(e) => {
953 error!(
954 "Failed to sync policy tags for list '{}' chunk {}/{}: {:?}",
955 token.mailchimp_mailing_list_id,
956 chunk_index + 1,
957 total_chunks,
958 e
959 );
960 }
961 }
962 }
963
964 let got_contact_id_set: std::collections::HashSet<Uuid> = user_id_contact_id_pairs
965 .iter()
966 .filter_map(|(uid, _)| Uuid::parse_str(uid).ok())
967 .collect();
968 let no_contact_id_user_ids: Vec<Uuid> = sent_user_ids
969 .iter()
970 .filter(|id| !got_contact_id_set.contains(id))
971 .copied()
972 .collect();
973
974 if !no_contact_id_user_ids.is_empty() {
975 let sample_len = no_contact_id_user_ids.len().min(10);
976 warn!(
977 "Mailchimp did not return contact_id for {} member(s) (likely unsubscribed, removed, or rejected). Marking synced_to_mailchimp_at to stop retry. First {} user_id(s): {:?}",
978 no_contact_id_user_ids.len(),
979 sample_len,
980 &no_contact_id_user_ids[..sample_len]
981 );
982 if let Err(e) = headless_lms_models::marketing_consents::update_synced_to_mailchimp_at_to_all_synced_users(
983 conn,
984 &no_contact_id_user_ids,
985 )
986 .await
987 {
988 error!(
989 "Failed to update synced_to_mailchimp_at for no-contact_id users: {:?}",
990 e
991 );
992 }
993 }
994
995 info!(
996 "Batch subscribe list '{}': sent {} member(s), got {} contact_id(s){}",
997 token.mailchimp_mailing_list_id,
998 sent_user_ids.len(),
999 user_id_contact_id_pairs.len(),
1000 if no_contact_id_user_ids.is_empty() {
1001 String::new()
1002 } else {
1003 format!(
1004 ", {} without contact_id (marked synced to stop retry)",
1005 no_contact_id_user_ids.len()
1006 )
1007 }
1008 );
1009
1010 if !user_id_contact_id_pairs.is_empty() {
1011 headless_lms_models::marketing_consents::update_user_mailchimp_id_at_to_all_synced_users(
1012 conn,
1013 user_id_contact_id_pairs,
1014 )
1015 .await?;
1016 }
1017
1018 Ok(successfully_synced_user_ids)
1019}
1020
1021async fn sync_completed_tag_for_members(
1023 users_details: &[UserMarketingConsentWithDetails],
1024 successfully_synced_user_ids: &[Uuid],
1025 mailchimp_id_by_user: &std::collections::HashMap<Uuid, String>,
1026 slug: &str,
1027 token: &MarketingMailingListAccessToken,
1028) -> anyhow::Result<()> {
1029 let tag_name = format!("{}-completed", slug);
1030 let success_set: std::collections::HashSet<_> = successfully_synced_user_ids.iter().collect();
1031
1032 let mut operations = Vec::new();
1033 for user in users_details {
1034 if !success_set.contains(&user.user_id) {
1035 continue;
1036 }
1037 let Some(user_mailchimp_id) = mailchimp_id_by_user.get(&user.user_id) else {
1038 continue;
1039 };
1040 let status = if user.completed_course_at.is_some() {
1041 "active"
1042 } else {
1043 "inactive"
1044 };
1045 operations.push(MailchimpOperation {
1046 method: Method::POST,
1047 path: format!(
1048 "/lists/{}/members/{}/tags",
1049 token.mailchimp_mailing_list_id, user_mailchimp_id
1050 ),
1051 body: Some(json!({
1052 "tags": [
1053 { "name": tag_name, "status": status }
1054 ]
1055 })),
1056 operation_id: Some(user.user_id.to_string()),
1057 });
1058 }
1059
1060 if operations.is_empty() {
1061 return Ok(());
1062 }
1063
1064 let ops_count = operations.len();
1065 info!(
1066 "Stage: completion tags (mark course completion). Preparing {} operation(s) for list '{}'",
1067 ops_count, token.mailchimp_mailing_list_id
1068 );
1069
1070 let timeout = Duration::from_secs(BATCH_POLL_TIMEOUT_SECS);
1071 let poll_interval = Duration::from_secs(BATCH_POLL_INTERVAL_SECS);
1072 let executor = MailchimpExecutor::new(timeout, poll_interval);
1073
1074 let start_time = Instant::now();
1075 let results = executor.execute(token, operations).await?;
1076 let failures: Vec<_> = results.iter().filter(|r| !r.is_success()).collect();
1077 if !failures.is_empty() {
1078 let sample: Vec<String> = failures
1079 .iter()
1080 .take(5)
1081 .map(|r| {
1082 format!(
1083 "{}: {}",
1084 r.operation_id.as_deref().unwrap_or("?"),
1085 r.error.as_deref().unwrap_or("unknown error")
1086 )
1087 })
1088 .collect();
1089 warn!(
1090 "Completion tag sync for list '{}' had {} failure(s) out of {}. Sample: {:?}",
1091 token.mailchimp_mailing_list_id,
1092 failures.len(),
1093 results.len(),
1094 sample
1095 );
1096 return Err(anyhow::anyhow!("completion tag sync failed"));
1097 }
1098
1099 info!(
1100 "Completed sync of {} completion tags for list '{}' in {:.2}s",
1101 ops_count,
1102 token.mailchimp_mailing_list_id,
1103 start_time.elapsed().as_secs_f64()
1104 );
1105
1106 Ok(())
1107}
1108
1109async fn update_emails_in_mailchimp(
1111 users: Vec<UserEmailSubscription>,
1112 list_id: &str,
1113 server_prefix: &str,
1114 access_token: &DbSecret,
1115) -> anyhow::Result<Vec<EmailSyncResult>> {
1116 let mut successfully_synced_users = Vec::new();
1117 let mut failed_user_ids = Vec::new();
1118
1119 for user in users {
1120 if let Some(ref user_mailchimp_id) = user.user_mailchimp_id {
1121 if let Some(ref status) = user.email_subscription_in_mailchimp
1122 && status != "subscribed"
1123 {
1124 continue; }
1126
1127 let url = format!(
1128 "https://{}.api.mailchimp.com/3.0/lists/{}/members/{}",
1129 server_prefix, list_id, user_mailchimp_id
1130 );
1131
1132 let body = serde_json::json!({
1134 "email_address": &user.email,
1135 "status": &user.email_subscription_in_mailchimp,
1136 });
1137
1138 let update_response = REQWEST_CLIENT
1140 .put(&url)
1141 .header(
1142 "Authorization",
1143 format!("apikey {}", access_token.expose_secret()),
1144 )
1145 .json(&body)
1146 .send()
1147 .await?;
1148
1149 if update_response.status().is_success() {
1150 successfully_synced_users.push(EmailSyncResult {
1151 user_id: user.user_id,
1152 user_mailchimp_id: user_mailchimp_id.clone(),
1153 });
1154 } else {
1155 failed_user_ids.push(user.user_id);
1156 }
1157 } else {
1158 continue;
1159 }
1160 }
1161
1162 if !failed_user_ids.is_empty() {
1163 info!("Failed to update the following users:");
1164 for user_id in &failed_user_ids {
1165 error!("User ID: {}", user_id);
1166 }
1167 }
1168
1169 Ok(successfully_synced_users)
1170}
1171
1172async fn fetch_unsubscribed_users_from_mailchimp_in_chunks(
1174 list_id: &str,
1175 server_prefix: &str,
1176 access_token: &DbSecret,
1177 chunk_size: usize,
1178) -> anyhow::Result<Vec<(String, String, String, String)>> {
1179 let mut all_data = Vec::new();
1180 let mut offset = 0;
1181
1182 loop {
1183 let url = format!(
1184 "https://{}.api.mailchimp.com/3.0/lists/{}/members?offset={}&count={}&fields=members.merge_fields,members.status,members.last_changed&status=unsubscribed,non-subscribed",
1185 server_prefix, list_id, offset, chunk_size
1186 );
1187
1188 let response = REQWEST_CLIENT
1189 .get(&url)
1190 .header(
1191 "Authorization",
1192 format!("apikey {}", access_token.expose_secret()),
1193 )
1194 .send()
1195 .await?
1196 .json::<serde_json::Value>()
1197 .await?;
1198
1199 let empty_vec = vec![];
1200 let members = response["members"].as_array().unwrap_or(&empty_vec);
1201 if members.is_empty() {
1202 break;
1203 }
1204
1205 for member in members {
1206 if let (Some(status), Some(last_changed), Some(merge_fields)) = (
1208 member["status"].as_str(),
1209 member["last_changed"].as_str(),
1210 member["merge_fields"].as_object(),
1211 ) {
1212 if let (Some(user_id), Some(language_group_id)) = (
1214 merge_fields.get("USERID").and_then(|v| v.as_str()),
1215 merge_fields.get("LANGGRPID").and_then(|v| v.as_str()),
1216 ) {
1217 if !user_id.is_empty() && !language_group_id.is_empty() {
1219 all_data.push((
1220 user_id.to_string(),
1221 last_changed.to_string(),
1222 language_group_id.to_string(),
1223 status.to_string(),
1224 ));
1225 }
1226 }
1227 }
1228 }
1229
1230 let total_items = response["total_items"].as_u64().unwrap_or(0) as usize;
1232 if offset + chunk_size >= total_items {
1233 break;
1234 }
1235
1236 offset += chunk_size;
1237 }
1238
1239 Ok(all_data)
1240}
1241
1242const BATCH_SIZE: usize = 1000;
1243
1244async fn process_unsubscribed_users_from_mailchimp(
1245 conn: &mut PgConnection,
1246 mailchimp_data: Vec<(String, String, String, String)>,
1247) -> anyhow::Result<()> {
1248 let total_records = mailchimp_data.len();
1249 let total_chunks = total_records.div_ceil(BATCH_SIZE);
1250
1251 for (chunk_num, chunk) in mailchimp_data.chunks(BATCH_SIZE).enumerate() {
1252 if chunk.is_empty() {
1253 continue;
1254 }
1255
1256 if let Err(e) = headless_lms_models::marketing_consents::update_unsubscribed_users_from_mailchimp_in_bulk(
1257 conn,
1258 chunk.to_vec(),
1259 )
1260 .await
1261 {
1262 error!(
1263 "Error while processing chunk {}/{}: {}",
1264 chunk_num + 1,
1265 total_chunks,
1266 e
1267 );
1268 }
1269 }
1270
1271 Ok(())
1272}