Skip to main content

headless_lms_server/programs/mailchimp_syncer/
mod.rs

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
85/// These fields are excluded from removing all fields that are not in the schema
86const 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;
95/// Bulk archive of every operation's result, so it can be far larger and slower than the JSON
96/// calls the shared client's default timeout is sized for.
97const 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
111/// The main function that initializes environment variables, config, and sync process.
112pub 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    // Iterate through access tokens and ensure Mailchimp schema is set up
130    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        // Check and sync tags for this access token once every hour
168        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                // this usually happens if the database is reset while running bin/dev etc.
194                info!("syncer may have lost its connection to the db, trying to reconnect");
195                conn = db_pool.acquire().await?;
196            }
197        }
198    }
199}
200
201/// Initializes environment variables, logging, and tracing setup.
202fn initialize_environment() -> anyhow::Result<()> {
203    dotenv().ok();
204    ProgramConfig::ensure_default_rust_log_for_workers();
205    setup_tracing()?;
206    Ok(())
207}
208
209/// Structure to hold the configuration settings, such as the database URL.
210struct SyncerConfig {
211    database_url: String,
212}
213
214/// Initializes and returns configuration settings (database URL).
215async fn initialize_configuration() -> anyhow::Result<SyncerConfig> {
216    let database_url = ProgramConfig::database_url_with_default();
217
218    Ok(SyncerConfig { database_url })
219}
220
221/// Initializes the PostgreSQL connection pool from the provided database URL.
222async 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
232/// Ensures the Mailchimp schema is up to date, adding required fields and removing any extra ones.
233async 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        // Remove extra fields not in REQUIRED_FIELDS or FIELDS_EXCLUDED_FROM_REMOVING
243        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    // Add any required fields that are missing
269    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
300/// Fetches the current merge fields from the Mailchimp list schema.
301async 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
350/// Adds a new merge field to the Mailchimp list.
351async 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
395/// Removes a merge field from the Mailchimp list by with a field ID.
396async 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
432/// Fetch tags from mailchimp and sync the changes to the database
433pub 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    // Extract tags from the response
456    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    // Fetch the current tags from the database
482    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    // Check if any tags from Mailchimp are renamed
490    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    // Check if any tags in the database have been removed via Mailchimp
512    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        // Check if this tag exists in the Mailchimp tags
520        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
540/// Synchronizes the user contacts with Mailchimp.
541/// Added a boolean flag to determine whether to process unsubscribes.
542async 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    // Iterate through tokens and fetch and send user details to Mailchimp
556    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        // Fetch all users from Mailchimp and sync possible changes locally
577        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        // Fetch unsynced emails and update them in Mailchimp
598        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        // Fetch unsynced user consents and update them in Mailchimp
683        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            // Store the successfully synced user IDs from syncing user consents
727            successfully_synced_user_ids.extend(consent_synced_user_ids);
728        }
729    }
730
731    // If there are any successfully synced users, update the database to mark them as synced
732    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
758/// Sends a batch of users to Mailchimp for synchronization.
759pub 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    // Prepare each user's data for Mailchimp
771    for user in users_details {
772        // Check user has given permission to send data to mailchimp
773        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
1021/// Sets or removes the "{slug}-completed" tag on Mailchimp members based on course completion. Skips users without user_mailchimp_id.
1022async 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
1109/// Updates the email addresses of multiple users in a Mailchimp mailing list.
1110async 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; // Skip this user if they are not subscribed because Mailchimp only updates emails that are subscribed
1125            }
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            // Prepare the body for the PUT request
1133            let body = serde_json::json!({
1134                "email_address": &user.email,
1135                "status": &user.email_subscription_in_mailchimp,
1136            });
1137
1138            // Update the email
1139            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
1172/// Fetches data from Mailchimp in chunks.
1173async 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            // Process the member, but only if necessary fields are present and valid
1207            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                // Ensure both USERID and LANGGRPID are present and valid
1213                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                    // Avoid adding data if any field is missing or empty
1218                    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        // Check the pagination info from the response
1231        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}