headless_lms_server/programs/
service_info_fetcher.rs1use crate::config::program_config::ProgramConfig;
2use crate::{domain::models_requests, setup_tracing};
3use anyhow::Result;
4use dotenvy::dotenv;
5use futures::stream::{self, StreamExt};
6use headless_lms_models::{
7 exercise_service_info::{ExerciseServiceInfo, fetch_and_upsert_service_info},
8 exercise_services::ExerciseService,
9};
10use sqlx::PgPool;
11use tokio::time::{Duration, sleep};
12use tracing::info;
13
14const N: usize = 10;
15
16pub async fn main() -> anyhow::Result<()> {
17 dotenv().ok();
18 ProgramConfig::ensure_default_rust_log_for_workers();
19 setup_tracing()?;
20
21 let database_url = ProgramConfig::database_url_with_default();
22 let db_pool = PgPool::connect(&database_url).await?;
23
24 let mut conn = db_pool.acquire().await?;
25
26 loop {
27 let exercise_services =
28 headless_lms_models::exercise_services::get_exercise_services(&mut conn).await?;
29 debug!(
30 "Fetching and updating statuses from {} services",
31 exercise_services.len()
32 );
33 let iter_stream = stream::iter(exercise_services.iter().map(|exercise_service| {
34 do_fetch_and_upsert_service_info(db_pool.clone(), exercise_service)
35 }));
36 let buffer_unordered = iter_stream.buffer_unordered(N);
38 let results = buffer_unordered.collect::<Vec<_>>().await;
39 let (succeeded, failed) = results.into_iter().partition::<Vec<_>, _>(|o| o.is_ok());
40 info!(
41 "Fetching and updating statuses complete. Succeeded: {}, failed: {}",
42 succeeded.len(),
43 failed.len()
44 );
45 sleep(Duration::from_secs(60)).await;
46 }
47}
48
49pub async fn do_fetch_and_upsert_service_info(
50 pool: PgPool,
51 exercise_service: &ExerciseService,
52) -> Result<ExerciseServiceInfo> {
53 let mut conn = pool.acquire().await?;
54 Ok(fetch_and_upsert_service_info(
55 &mut conn,
56 exercise_service,
57 models_requests::fetch_service_info,
58 )
59 .await?)
60}