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 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 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 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
140async 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
151async 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
158async 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 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 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
351async 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
373async fn sync_pages_batch(
375 conn: &mut PgConnection,
376 pages: &[Page],
377 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 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 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 serde_json::to_string(&sanitized_blocks)?
479 }
480 }
481 };
482
483 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 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 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
575fn 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
584async 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}