Skip to main content

headless_lms_server/programs/
chatbot_syncer.rs

1use std::{
2    collections::{HashMap, HashSet},
3    time::Duration,
4};
5
6use chrono::Utc;
7use dotenvy::dotenv;
8use sqlx::{PgConnection, PgPool};
9use url::Url;
10use uuid::Uuid;
11
12use crate::config::program_config::ProgramConfig;
13use crate::setup_tracing;
14use headless_lms_utils::periodic_worker::{
15    PeriodicWorkerConfig, StillRunningLog, run_periodic_worker,
16};
17
18use headless_lms_base::config::ApplicationConfiguration;
19use headless_lms_chatbot::{
20    azure_blob_storage::AzureBlobClient,
21    azure_datasources::{create_azure_datasource, does_azure_datasource_exist},
22    azure_search_index::{create_search_index, does_search_index_exist},
23    azure_search_indexer::{
24        check_search_indexer_status, create_search_indexer, does_search_indexer_exist,
25        run_search_indexer_now,
26    },
27    azure_skillset::{create_skillset, does_skillset_exist},
28    content_cleaner::convert_material_blocks_to_markdown_with_llm,
29};
30use headless_lms_models::{
31    application_task_default_language_models::ApplicationTask,
32    chapters::DatabaseChapter,
33    course_page_markdown_content::CoursePageMarkdownContent,
34    pages::{Page, PageVisibility},
35};
36use headless_lms_utils::{
37    document_schema_processor::{GutenbergBlock, remove_sensitive_attributes},
38    url_encoding::url_encode,
39};
40
41const SYNC_INTERVAL_SECS: u64 = 10;
42const PRINT_STILL_RUNNING_MESSAGE_TICKS_THRESHOLD: u32 = 60;
43const FAILURE_COOLDOWN_SECS: i64 = 300;
44const MAX_CONSECUTIVE_FAILURES: i32 = 5;
45
46pub async fn main() -> anyhow::Result<()> {
47    initialize_environment()?;
48    let config = initialize_configuration().await?;
49    if config.app_configuration.azure_configuration.is_none() {
50        warn!("Azure configuration not provided. Not running chatbot syncer.");
51        // Sleep indefinitely to prevent the program from exiting. This only happens in development.
52        loop {
53            tokio::time::sleep(Duration::from_secs(u64::MAX)).await;
54        }
55    }
56    if config.app_configuration.test_chatbot {
57        warn!(
58            "Using mock azure configuration, this must be a test/dev environment. Not running chatbot syncer."
59        );
60        // Sleep indefinitely to prevent the program from exiting. This only happens in development.
61        loop {
62            tokio::time::sleep(Duration::from_secs(u64::MAX)).await;
63        }
64    }
65
66    let db_pool = initialize_database_pool(&config.database_url).await?;
67    let blob_client = initialize_blob_client(&config).await?;
68
69    let mut reported_permanently_failing_page_ids: HashSet<Uuid> = HashSet::new();
70
71    info!("Starting chatbot syncer.");
72
73    run_periodic_worker(
74        PeriodicWorkerConfig {
75            tick_interval: Duration::from_secs(SYNC_INTERVAL_SECS),
76            still_running: Some(StillRunningLog {
77                every: PRINT_STILL_RUNNING_MESSAGE_TICKS_THRESHOLD,
78                message: "Still syncing for chatbot.",
79                initial_ticks: 0,
80            }),
81            delay_missed_ticks: false,
82        },
83        async || {
84            // Acquired per tick so that the pool's liveness check replaces a connection that died
85            // while we were idle; a checked out connection is never healed.
86            match db_pool.acquire().await {
87                Ok(mut conn) => {
88                    if let Err(e) = sync_pages(
89                        &mut conn,
90                        &config,
91                        &blob_client,
92                        &mut reported_permanently_failing_page_ids,
93                    )
94                    .await
95                    {
96                        error!("Error during synchronization: {:?}", e);
97                    }
98                }
99                Err(e) => error!("Failed to acquire a database connection: {:?}", e),
100            }
101            Ok(())
102        },
103    )
104    .await
105}
106
107fn initialize_environment() -> anyhow::Result<()> {
108    dotenv().ok();
109    ProgramConfig::ensure_default_rust_log_for_workers();
110    setup_tracing()?;
111    Ok(())
112}
113
114struct SyncerConfig {
115    database_url: String,
116    name: String,
117    app_configuration: ApplicationConfiguration,
118}
119
120async fn initialize_configuration() -> anyhow::Result<SyncerConfig> {
121    let database_url = ProgramConfig::database_url_with_default();
122    let base_url_raw = ProgramConfig::required("BASE_URL")?;
123    let base_url =
124        Url::parse(&base_url_raw).map_err(|e| anyhow::anyhow!("invalid BASE_URL: {}", e))?;
125
126    let name = base_url
127        .host_str()
128        .ok_or_else(|| anyhow::anyhow!("BASE_URL must have a host"))?
129        .replace(".", "-");
130
131    let app_configuration = ApplicationConfiguration::try_from_env()?;
132
133    Ok(SyncerConfig {
134        database_url,
135        name,
136        app_configuration,
137    })
138}
139
140/// Initializes the PostgreSQL connection pool.
141async fn initialize_database_pool(database_url: &str) -> anyhow::Result<PgPool> {
142    PgPool::connect(database_url).await.map_err(|e| {
143        anyhow::anyhow!(
144            "Failed to connect to the database at {}: {:?}",
145            database_url,
146            e
147        )
148    })
149}
150
151/// Initializes the Azure Blob Storage client.
152async fn initialize_blob_client(config: &SyncerConfig) -> anyhow::Result<AzureBlobClient> {
153    let blob_client = AzureBlobClient::new(&config.app_configuration, &config.name).await?;
154    blob_client.ensure_container_exists().await?;
155    Ok(blob_client)
156}
157
158/// Synchronizes pages to the chatbot backend.
159/// `reported_permanently_failing_page_ids` carries the previously warned-about set across ticks so
160/// the warning is emitted on change rather than on every sweep.
161async fn sync_pages(
162    conn: &mut PgConnection,
163    config: &SyncerConfig,
164    blob_client: &AzureBlobClient,
165    reported_permanently_failing_page_ids: &mut HashSet<Uuid>,
166) -> anyhow::Result<()> {
167    let base_url = Url::parse(&config.app_configuration.base_url)?;
168    let chatbot_configs =
169        headless_lms_models::chatbot_configurations::get_for_azure_search_maintenance(conn).await?;
170
171    let course_ids: Vec<Uuid> = chatbot_configs
172        .iter()
173        .filter_map(|config| config.course_id)
174        .collect::<HashSet<_>>()
175        .into_iter()
176        .collect();
177
178    let sync_statuses =
179        headless_lms_models::chatbot_page_sync_statuses::ensure_sync_statuses_exist(
180            conn,
181            &course_ids,
182        )
183        .await?;
184
185    // (page_id, page_history_id)
186    let latest_history_ids =
187        headless_lms_models::page_history::get_latest_page_history_ids_by_course_ids(
188            conn,
189            &course_ids,
190        )
191        .await?;
192
193    let shared_index_name = config.name.clone();
194    ensure_search_index_exists(
195        &shared_index_name,
196        &config.app_configuration,
197        &blob_client.container_name,
198    )
199    .await?;
200
201    if !check_search_indexer_status(&shared_index_name, &config.app_configuration).await? {
202        warn!("Search indexer is not ready to index. Skipping synchronization.");
203        return Ok(());
204    }
205
206    let mut any_changes = false;
207    let mut permanently_failing_page_ids: HashSet<Uuid> = HashSet::new();
208
209    for (course_id, statuses) in sync_statuses.iter() {
210        let page_ids: Vec<Uuid> = statuses.iter().map(|s| s.page_id).collect();
211        let public_pages_set: HashSet<Uuid> =
212            headless_lms_models::pages::get_by_ids_and_visibility(
213                conn,
214                &page_ids,
215                PageVisibility::Public,
216            )
217            .await?
218            .into_iter()
219            .map(|p| p.id)
220            .collect();
221
222        let outdated_statuses: Vec<_> = statuses
223            .iter()
224            .filter(|status| {
225                if !public_pages_set.contains(&status.page_id) {
226                    return false;
227                }
228
229                let is_outdated = latest_history_ids
230                    .get(&status.page_id)
231                    .is_some_and(|history_id| {
232                        status.synced_page_revision_id != Some(*history_id)
233                    });
234
235                if !is_outdated {
236                    return false;
237                }
238
239                if status.consecutive_failures >= MAX_CONSECUTIVE_FAILURES {
240                    permanently_failing_page_ids.insert(status.page_id);
241                    return false;
242                }
243
244                if let Some(error_msg) = &status.error_message
245                    && !error_msg.is_empty() {
246                        let error_age_seconds = (Utc::now() - status.updated_at).num_seconds();
247                        if error_age_seconds < FAILURE_COOLDOWN_SECS {
248                            debug!(
249                                "Skipping page {} due to recent failure ({} seconds ago, {} consecutive failures): {}",
250                                status.page_id, error_age_seconds, status.consecutive_failures, error_msg
251                            );
252                            return false;
253                        }
254                    }
255
256                true
257            })
258            .collect();
259
260        if outdated_statuses.is_empty() {
261            continue;
262        }
263
264        any_changes = true;
265        info!(
266            "Syncing {} pages for course id: {}.",
267            outdated_statuses.len(),
268            course_id
269        );
270        for status in &outdated_statuses {
271            info!(
272                "Page id: {}, synced page revision id: {:?}.",
273                status.page_id, status.synced_page_revision_id
274            );
275        }
276
277        let page_ids: Vec<Uuid> = outdated_statuses.iter().map(|s| s.page_id).collect();
278        let md_ids: Vec<Uuid> = outdated_statuses
279            .iter()
280            .filter_map(|s| s.converted_markdown_content_id)
281            .collect();
282        let pages = headless_lms_models::pages::get_by_ids_and_visibility(
283            conn,
284            &page_ids,
285            PageVisibility::Public,
286        )
287        .await?;
288
289        if !pages.is_empty() {
290            sync_pages_batch(
291                conn,
292                &pages,
293                &md_ids,
294                blob_client,
295                &base_url,
296                &config.app_configuration,
297                &latest_history_ids,
298            )
299            .await?;
300        } else {
301            info!("No pages to sync for course id: {}.", course_id);
302        }
303
304        let hidden_page_ids: Vec<Uuid> = statuses
305            .iter()
306            .filter(|status| {
307                !public_pages_set.contains(&status.page_id)
308                    && status.synced_page_revision_id.is_some()
309            })
310            .map(|s| s.page_id)
311            .collect();
312
313        if !hidden_page_ids.is_empty() {
314            info!(
315                "Clearing sync statuses for {} hidden pages: {:?}",
316                hidden_page_ids.len(),
317                hidden_page_ids
318            );
319            headless_lms_models::chatbot_page_sync_statuses::clear_sync_statuses(
320                conn,
321                &hidden_page_ids,
322            )
323            .await?;
324        }
325
326        delete_old_files(conn, *course_id, blob_client).await?;
327    }
328
329    // Only on change: a page stays permanently failing until someone intervenes, and this runs
330    // every SYNC_INTERVAL_SECS, so warning unconditionally would repeat the same line all day.
331    if permanently_failing_page_ids != *reported_permanently_failing_page_ids {
332        if !permanently_failing_page_ids.is_empty() {
333            warn!(
334                "Skipping {} pages that have failed to sync at least {} times in a row. Manual intervention required: {:?}",
335                permanently_failing_page_ids.len(),
336                MAX_CONSECUTIVE_FAILURES,
337                permanently_failing_page_ids
338            );
339        }
340        *reported_permanently_failing_page_ids = permanently_failing_page_ids;
341    }
342
343    if any_changes {
344        run_search_indexer_now(&shared_index_name, &config.app_configuration).await?;
345        info!("New files have been synced and the search indexer has been started.");
346    }
347
348    Ok(())
349}
350
351/// Ensures that the specified search index exists, creating it if necessary.
352async fn ensure_search_index_exists(
353    name: &str,
354    app_config: &ApplicationConfiguration,
355    container_name: &str,
356) -> anyhow::Result<()> {
357    if !does_search_index_exist(name, app_config).await? {
358        create_search_index(name.to_owned(), app_config).await?;
359    }
360    if !does_skillset_exist(name, app_config).await? {
361        create_skillset(name, name, app_config).await?;
362    }
363    if !does_azure_datasource_exist(name, app_config).await? {
364        create_azure_datasource(name, container_name, app_config).await?;
365    }
366    if !does_search_indexer_exist(name, app_config).await? {
367        create_search_indexer(name, name, name, name, app_config).await?;
368    }
369
370    Ok(())
371}
372
373/// Processes and synchronizes a batch of pages.
374async fn sync_pages_batch(
375    conn: &mut PgConnection,
376    pages: &[Page],
377    // map from page id to course page markdown content id
378    md_ids: &[Uuid],
379    blob_client: &AzureBlobClient,
380    base_url: &Url,
381    app_config: &ApplicationConfiguration,
382    latest_history_ids: &HashMap<Uuid, Uuid>,
383) -> anyhow::Result<()> {
384    let course_id = pages
385        .first()
386        .ok_or_else(|| anyhow::anyhow!("No pages to sync."))?
387        .course_id
388        .ok_or_else(|| anyhow::anyhow!("The first page does not belong to any course."))?;
389
390    let course = headless_lms_models::courses::get_course(conn, course_id).await?;
391    let chapters = headless_lms_models::chapters::get_course_chapters(conn, course_id).await?;
392    let md_contents =
393        headless_lms_models::course_page_markdown_content::get_many(conn, md_ids).await?;
394    let organization =
395        headless_lms_models::organizations::get_organization(conn, course.organization_id).await?;
396    let task_lm = headless_lms_models::application_task_default_language_models::get_for_task(
397        conn,
398        ApplicationTask::ContentCleaning,
399    )
400    .await?;
401
402    let mut base_url = base_url.clone();
403    base_url.set_path(&format!(
404        "/org/{}/courses/{}",
405        organization.slug, course.slug
406    ));
407
408    let mut allowed_file_paths = Vec::new();
409    let mut page_revision_map = HashMap::new();
410    // newly created md. map page_id to page_history_id and md content
411    let mut new_markdown_contents_map = HashMap::new();
412
413    for page in pages {
414        info!("Syncing page id: {}.", page.id);
415
416        let mut page_url = base_url.clone();
417        page_url.set_path(&format!("{}{}", base_url.path(), page.url_path));
418
419        let parsed_content: Vec<GutenbergBlock> = serde_json::from_value(page.content.clone())?;
420        let sanitized_blocks = remove_sensitive_attributes(parsed_content);
421
422        let page_md_content: Option<&CoursePageMarkdownContent> =
423            md_contents.iter().find(|x| x.page_id == page.id);
424        let latest_page_history_id: Option<&Uuid> = latest_history_ids.get(&page.id);
425
426        let up_to_date_md_content = page_md_content.and_then(|c| {
427            latest_page_history_id.and_then(|id| {
428                if id == &c.page_history_id {
429                    Some(c.markdown_content.to_string())
430                } else {
431                    None
432                }
433            })
434        });
435
436        let content_as_markdown = if let Some(content) = up_to_date_md_content.to_owned() {
437            info!("Using previously generated Markdown for page {}", page.id);
438            content
439        } else {
440            match convert_material_blocks_to_markdown_with_llm(
441                &sanitized_blocks,
442                app_config,
443                &task_lm,
444            )
445            .await
446            {
447                Ok(markdown) => {
448                    info!("Successfully cleaned content for page {}", page.id);
449                    // Check if the markdown is empty, or if it just contains all spaces or newlines
450                    if markdown.trim().is_empty() {
451                        warn!(
452                            "Markdown is empty for page {}. Generating fallback content with a fake heading.",
453                            page.id
454                        );
455                        format!("# {}", page.title)
456                    } else {
457                        markdown
458                    }
459                }
460                Err(e) => {
461                    let error_msg = format!("Sync failed: LLM processing error: {}", e);
462                    warn!(
463                        "Failed to clean content with LLM for page {}: {}. Using serialized sanitized content instead.",
464                        page.id, error_msg
465                    );
466                    if let Err(db_err) =
467                        headless_lms_models::chatbot_page_sync_statuses::set_page_sync_error(
468                            conn, page.id, &error_msg,
469                        )
470                        .await
471                    {
472                        warn!(
473                            "Failed to record sync error for page {}: {:?}",
474                            page.id, db_err
475                        );
476                    }
477                    // Fallback to original content
478                    serde_json::to_string(&sanitized_blocks)?
479                }
480            }
481        };
482
483        // save markdown content if new markdown was generated
484        // if there is an error saving it to blobs, we can try uploading the same content
485        // if the page hasn't been changed between tries.
486        if let Some(history_id) = latest_page_history_id
487            && up_to_date_md_content.is_none()
488        {
489            new_markdown_contents_map.insert(
490                page.id,
491                (history_id.to_owned(), content_as_markdown.to_owned()),
492            );
493        }
494
495        let blob_path = generate_blob_path(page)?;
496        let chapter: Option<&DatabaseChapter> = chapters
497            .iter()
498            .find(|c| page.chapter_id.is_some_and(|c_id| c_id == c.id));
499
500        allowed_file_paths.push(blob_path.clone());
501        let mut metadata = HashMap::new();
502        // Azure Blob Storage metadata values must be ASCII-only. URL-encode values that may
503        // contain non-ASCII characters (e.g., Finnish characters like ä, ö) to ensure they
504        // are ASCII-compatible. We decode the url and the title before we save them in our database.
505        metadata.insert("url".to_string(), url_encode(page_url.as_ref()));
506        metadata.insert("title".to_string(), url_encode(&page.title));
507        metadata.insert(
508            "course_id".to_string(),
509            page.course_id.unwrap_or(Uuid::nil()).to_string().into(),
510        );
511        metadata.insert(
512            "language".to_string(),
513            course.language_code.to_string().into(),
514        );
515        metadata.insert("filepath".to_string(), blob_path.clone().into());
516        if let Some(c) = chapter {
517            metadata.insert(
518                "chunk_context".to_string(),
519                url_encode(&format!(
520                    "This chunk is a snippet from page {} from chapter {}: {} of the course {}.",
521                    page.title, c.chapter_number, c.name, course.name,
522                )),
523            );
524        } else {
525            metadata.insert(
526                "chunk_context".to_string(),
527                url_encode(&format!(
528                    "This chunk is a snippet from page {} of the course {}.",
529                    page.title, course.name,
530                )),
531            );
532        }
533
534        if let Err(e) = blob_client
535            .upload_file(&blob_path, content_as_markdown.as_bytes(), Some(metadata))
536            .await
537        {
538            let error_msg = format!("Sync failed: Upload error: {}", e);
539            warn!("Failed to upload file {}: {:?}", blob_path, e);
540            if let Err(db_err) =
541                headless_lms_models::chatbot_page_sync_statuses::set_page_sync_error(
542                    conn, page.id, &error_msg,
543                )
544                .await
545            {
546                warn!(
547                    "Failed to record upload error for page {}: {:?}",
548                    page.id, db_err
549                );
550            }
551        } else if let Some(history_id) = latest_page_history_id {
552            page_revision_map.insert(page.id, *history_id);
553        }
554    }
555
556    if let Err(e) = headless_lms_models::chatbot_page_sync_statuses::save_markdown_content(
557        conn,
558        new_markdown_contents_map,
559    )
560    .await
561    {
562        warn!("Failed to save converted page content in DB: {}", e);
563    };
564
565    // update revision ids for all pages
566    headless_lms_models::chatbot_page_sync_statuses::update_page_revision_ids(
567        conn,
568        page_revision_map,
569    )
570    .await?;
571
572    Ok(())
573}
574
575/// Generates the blob storage path for a given page.
576fn generate_blob_path(page: &Page) -> anyhow::Result<String> {
577    let course_id = page
578        .course_id
579        .ok_or_else(|| anyhow::anyhow!("Page {} does not belong to any course.", page.id))?;
580
581    Ok(format!("courses/{}/pages/{}.md", course_id, page.id))
582}
583
584/// Deletes files from blob storage that are no longer associated with any public page.
585/// This includes files for deleted pages, hidden pages, and any other pages that are no longer public.
586async fn delete_old_files(
587    conn: &mut PgConnection,
588    course_id: Uuid,
589    blob_client: &AzureBlobClient,
590) -> anyhow::Result<()> {
591    let mut courses_prefix = "courses/".to_string();
592    courses_prefix.push_str(&course_id.to_string());
593    let existing_files = blob_client.list_files_with_prefix(&courses_prefix).await?;
594
595    let pages = headless_lms_models::pages::get_all_by_course_id_and_visibility(
596        conn,
597        course_id,
598        PageVisibility::Public,
599    )
600    .await?;
601
602    let allowed_paths: HashSet<String> = pages
603        .iter()
604        .filter_map(|page| generate_blob_path(page).ok())
605        .collect();
606
607    for file in existing_files {
608        if !allowed_paths.contains(&file) {
609            info!("Deleting obsolete file: {}", file);
610            blob_client.delete_file(&file).await?;
611        }
612    }
613
614    Ok(())
615}