diff --git a/package.json b/package.json index 8de05db..9799019 100644 --- a/package.json +++ b/package.json @@ -8,6 +8,7 @@ "build": "pnpm turbo run build --filter='./packages/*' --filter='./apps/*'", "build:rust": "turbo run build:rust", "backfill": "./scripts/backfill-tap.sh", + "backfill:lightrail": "SQLX_OFFLINE=true cargo run -p cadet --bin lightrail-backfill", "tunnel:up": "./scripts/dev-tunnel.sh up", "tunnel:down": "./scripts/dev-tunnel.sh down", "tunnel:status": "./scripts/dev-tunnel.sh status", diff --git a/services/cadet/Cargo.toml b/services/cadet/Cargo.toml index d442d66..5219cf6 100644 --- a/services/cadet/Cargo.toml +++ b/services/cadet/Cargo.toml @@ -7,6 +7,10 @@ edition = "2021" name = "tap-backfill" path = "src/bin/tap_backfill.rs" +[[bin]] +name = "lightrail-backfill" +path = "src/bin/lightrail_backfill.rs" + [dependencies] tokio.workspace = true tokio-tungstenite.workspace = true diff --git a/services/cadet/src/bin/lightrail_backfill.rs b/services/cadet/src/bin/lightrail_backfill.rs new file mode 100644 index 0000000..e3e6826 --- /dev/null +++ b/services/cadet/src/bin/lightrail_backfill.rs @@ -0,0 +1,230 @@ +use std::{ + collections::HashSet, + fs::{self, OpenOptions}, + io::Write, + path::PathBuf, + time::Duration, +}; + +use cadet::{db, ingestors::car::CarImportIngestor}; +use serde::Deserialize; +use sqlx::PgPool; +use tokio::time::sleep; + +const DEFAULT_LIGHTRAIL_URL: &str = "https://lightrail.microcosm.blue"; +const DEFAULT_COLLECTION: &str = "fm.teal.alpha.feed.play"; +const DEFAULT_STATE_FILE: &str = ".teal-lightrail-backfill-done.txt"; + +#[derive(Debug, Deserialize)] +struct LightrailRepo { + did: String, +} + +#[derive(Debug, Deserialize)] +struct ListReposByCollectionResponse { + cursor: Option, + repos: Vec, +} + +fn setup_tracing() { + tracing_subscriber::fmt() + .with_max_level(tracing::Level::ERROR) + .init(); +} + +fn env_u64(name: &str, default: u64) -> u64 { + std::env::var(name) + .ok() + .and_then(|value| value.parse().ok()) + .unwrap_or(default) +} + +fn env_usize(name: &str, default: usize) -> usize { + std::env::var(name) + .ok() + .and_then(|value| value.parse().ok()) + .unwrap_or(default) +} + +fn load_completed(path: &PathBuf) -> anyhow::Result> { + if std::env::var("LIGHTRAIL_BACKFILL_RESET").as_deref() == Ok("1") { + let _ = fs::remove_file(path); + } + + match fs::read_to_string(path) { + Ok(contents) => Ok(contents + .lines() + .map(str::trim) + .filter(|line| !line.is_empty()) + .map(ToString::to_string) + .collect()), + Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(HashSet::new()), + Err(err) => Err(err.into()), + } +} + +fn record_completed(path: &PathBuf, did: &str) -> anyhow::Result<()> { + let mut file = OpenOptions::new().create(true).append(true).open(path)?; + writeln!(file, "{did}")?; + file.flush()?; + Ok(()) +} + +async fn discover_repos( + client: &reqwest::Client, + lightrail_url: &str, + collection: &str, + page_limit: usize, +) -> anyhow::Result> { + let mut cursor = None::; + let mut repos = Vec::new(); + + loop { + let mut request = client + .get(format!( + "{}/xrpc/com.atproto.sync.listReposByCollection", + lightrail_url.trim_end_matches('/') + )) + .query(&[ + ("collection", collection), + ("limit", &page_limit.to_string()), + ]); + + if let Some(cursor) = cursor.as_deref() { + request = request.query(&[("cursor", cursor)]); + } + + let response = request.send().await?.error_for_status()?; + let page = response.json::().await?; + let page_len = page.repos.len(); + repos.extend(page.repos.into_iter().map(|repo| repo.did)); + + match page.cursor { + Some(next_cursor) + if page_len > 0 && cursor.as_deref() != Some(next_cursor.as_str()) => + { + cursor = Some(next_cursor); + } + _ => break, + } + } + + repos.sort(); + repos.dedup(); + Ok(repos) +} + +async fn refresh_materialized_views(pool: &PgPool) -> anyhow::Result<()> { + sqlx::query("REFRESH MATERIALIZED VIEW mv_artist_play_counts") + .execute(pool) + .await?; + sqlx::query("REFRESH MATERIALIZED VIEW mv_release_play_counts") + .execute(pool) + .await?; + sqlx::query("REFRESH MATERIALIZED VIEW mv_recording_play_counts") + .execute(pool) + .await?; + sqlx::query("REFRESH MATERIALIZED VIEW mv_global_play_count") + .execute(pool) + .await?; + Ok(()) +} + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + dotenvy::dotenv().ok(); + setup_tracing(); + + std::env::set_var("CADET_DEFER_MATERIALIZED_VIEW_REFRESH", "1"); + + let lightrail_url = + std::env::var("LIGHTRAIL_URL").unwrap_or_else(|_| DEFAULT_LIGHTRAIL_URL.to_string()); + let collection = + std::env::var("LIGHTRAIL_COLLECTION").unwrap_or_else(|_| DEFAULT_COLLECTION.to_string()); + let page_limit = env_usize("LIGHTRAIL_PAGE_LIMIT", 1000); + let max_repos = env_usize("LIGHTRAIL_BACKFILL_MAX_REPOS", usize::MAX); + let retry_count = env_u64("LIGHTRAIL_BACKFILL_RETRIES", 2); + let retry_sleep_ms = env_u64("LIGHTRAIL_BACKFILL_RETRY_SLEEP_MS", 2500); + let state_file = PathBuf::from( + std::env::var("LIGHTRAIL_BACKFILL_STATE_FILE") + .unwrap_or_else(|_| DEFAULT_STATE_FILE.to_string()), + ); + + let client = reqwest::Client::builder() + .timeout(Duration::from_secs(120)) + .user_agent("teal-cadet-lightrail-backfill/0.1") + .build()?; + + eprintln!("Discovering repos with {collection} via {lightrail_url}"); + let mut repos = discover_repos(&client, &lightrail_url, &collection, page_limit).await?; + if repos.len() > max_repos { + repos.truncate(max_repos); + } + + let mut completed = load_completed(&state_file)?; + let pool = db::init_pool().await?; + let ingestor = CarImportIngestor::new(pool.clone()); + let total = repos.len(); + let mut imported = 0_usize; + let mut skipped = 0_usize; + let mut failed = Vec::new(); + + eprintln!( + "Starting Lightrail CAR backfill for {total} repos ({already_done} already done)", + already_done = completed.len() + ); + + for (index, did) in repos.into_iter().enumerate() { + let ordinal = index + 1; + if completed.contains(&did) { + skipped += 1; + continue; + } + + eprintln!("[{ordinal}/{total}] importing {did}"); + let mut last_error = None; + for attempt in 0..=retry_count { + match ingestor.fetch_and_process_identity_car(&did).await { + Ok(import_id) => { + record_completed(&state_file, &did)?; + completed.insert(did.clone()); + imported += 1; + eprintln!("[{ordinal}/{total}] imported {did} ({import_id})"); + last_error = None; + break; + } + Err(err) => { + last_error = Some(err); + if attempt < retry_count { + eprintln!( + "[{ordinal}/{total}] retrying {did} after failure: {last_error}", + last_error = last_error.as_ref().unwrap() + ); + sleep(Duration::from_millis(retry_sleep_ms)).await; + } + } + } + } + + if let Some(err) = last_error { + eprintln!("[{ordinal}/{total}] failed {did}: {err}"); + failed.push((did, err.to_string())); + } + } + + eprintln!("Refreshing materialized play views"); + refresh_materialized_views(&pool).await?; + + eprintln!( + "Lightrail backfill complete: imported {imported}, skipped {skipped}, failed {}", + failed.len() + ); + if !failed.is_empty() { + for (did, error) in &failed { + eprintln!("failed {did}: {error}"); + } + anyhow::bail!("{} repos failed during Lightrail backfill", failed.len()); + } + + Ok(()) +} diff --git a/services/cadet/src/ingestors/car/car_import.rs b/services/cadet/src/ingestors/car/car_import.rs index 27b8305..03911fa 100644 --- a/services/cadet/src/ingestors/car/car_import.rs +++ b/services/cadet/src/ingestors/car/car_import.rs @@ -49,6 +49,7 @@ use tracing::{info, warn}; pub struct ExtractedRecord { pub collection: String, pub rkey: String, + pub cid: String, pub data: serde_json::Value, } @@ -170,6 +171,7 @@ impl CarImportIngestor { records.push(ExtractedRecord { collection, rkey, + cid: record_cid.to_string(), data, }); } @@ -251,7 +253,7 @@ impl CarImportIngestor { "fm.teal.alpha.feed.play" => { info!(" 📀 Processing play record..."); let result = self - .process_play_record(&record.data, did, &record.rkey) + .process_play_record(&record.data, did, &record.rkey, &record.cid) .await; if result.is_ok() { info!(" ✅ Successfully processed play record"); @@ -275,7 +277,7 @@ impl CarImportIngestor { "fm.teal.alpha.actor.status" => { info!(" 📢 Processing status record..."); let result = self - .process_status_record(&record.data, did, &record.rkey) + .process_status_record(&record.data, did, &record.rkey, &record.cid) .await; if result.is_ok() { info!(" ✅ Successfully processed status record"); @@ -287,7 +289,7 @@ impl CarImportIngestor { "fm.teal.alpha.actor.profileStatus" => { info!(" 🧭 Processing profile status record..."); let result = self - .process_profile_status_record(&record.data, did, &record.rkey) + .process_profile_status_record(&record.data, did, &record.rkey, &record.cid) .await; if result.is_ok() { info!(" ✅ Successfully processed profile status record"); @@ -309,7 +311,13 @@ impl CarImportIngestor { | "fm.teal.alpha.feed.social.badgeAssignment" => { info!(" 💬 Processing social record..."); let result = self - .process_social_record(&record.collection, &record.data, did, &record.rkey) + .process_social_record( + &record.collection, + &record.data, + did, + &record.rkey, + &record.cid, + ) .await; if result.is_ok() { info!(" ✅ Successfully processed social record"); @@ -342,21 +350,22 @@ impl CarImportIngestor { } /// Process a play record using the existing PlayIngestor - async fn process_play_record(&self, data: &Value, did: &str, rkey: &str) -> Result<()> { + async fn process_play_record( + &self, + data: &Value, + did: &str, + rkey: &str, + cid: &str, + ) -> Result<()> { + let data = Self::normalize_play_record_json(data.clone()); let play_record: types::fm_teal::alpha::feed::play::Play = - value::from_json_value::(data.clone())?; + value::from_json_value::(data)?; let play_ingestor = super::super::teal::feed_play::PlayIngestor::new(self.sql.clone()); let uri = super::super::teal::assemble_at_uri(did, "fm.teal.alpha.feed.play", rkey); play_ingestor - .insert_play( - &play_record, - &uri, - &format!("car-import-{}", uuid::Uuid::new_v4()), - did, - rkey, - ) + .insert_play(&play_record, &uri, cid, did, rkey) .await?; info!( @@ -366,6 +375,22 @@ impl CarImportIngestor { Ok(()) } + fn normalize_play_record_json(mut data: Value) -> Value { + let Some(duration) = data.get_mut("duration") else { + return data; + }; + + let Some(float_duration) = duration.as_f64() else { + return data; + }; + + if duration.as_i64().is_none() && float_duration.is_finite() { + *duration = Value::Number((float_duration.round() as i64).into()); + } + + data + } + /// Process a profile record using the existing ActorProfileIngestor async fn process_profile_record(&self, data: &Value, did: &str, _rkey: &str) -> Result<()> { let profile_record = super::super::teal::actor_profile::deserialize_profile(data)?; @@ -385,7 +410,13 @@ impl CarImportIngestor { } /// Process a status record using the existing ActorStatusIngestor - async fn process_status_record(&self, data: &Value, did: &str, rkey: &str) -> Result<()> { + async fn process_status_record( + &self, + data: &Value, + did: &str, + rkey: &str, + cid: &str, + ) -> Result<()> { let status_record: types::fm_teal::alpha::actor::status::Status = value::from_json_value::(data.clone())?; @@ -393,12 +424,7 @@ impl CarImportIngestor { super::super::teal::actor_status::ActorStatusIngestor::new(self.sql.clone()); status_ingestor - .insert_status( - did, - rkey, - &format!("car-import-{}", uuid::Uuid::new_v4()), - &status_record, - ) + .insert_status(did, rkey, cid, &status_record) .await?; info!("Successfully stored status record from CAR import"); @@ -411,6 +437,7 @@ impl CarImportIngestor { data: &Value, did: &str, rkey: &str, + cid: &str, ) -> Result<()> { let profile_status_record: types::fm_teal::alpha::actor::profile_status::ProfileStatus = value::from_json_value::( @@ -423,12 +450,7 @@ impl CarImportIngestor { ); profile_status_ingestor - .insert_profile_status( - did, - rkey, - &format!("car-import-{}", uuid::Uuid::new_v4()), - &profile_status_record, - ) + .insert_profile_status(did, rkey, cid, &profile_status_record) .await?; info!("Successfully stored profile status record from CAR import"); @@ -442,6 +464,7 @@ impl CarImportIngestor { data: &Value, did: &str, rkey: &str, + cid: &str, ) -> Result<()> { let kind = match collection { "fm.teal.alpha.feed.social.post" => super::super::teal::social::SocialCollection::Post, @@ -466,14 +489,7 @@ impl CarImportIngestor { }; let social_ingestor = super::super::teal::social::SocialRecordIngestor::new(self.sql.clone(), kind); - social_ingestor - .insert_record( - did, - rkey, - &format!("car-import-{}", uuid::Uuid::new_v4()), - data, - ) - .await + social_ingestor.insert_record(did, rkey, cid, data).await } /// Fetch and process a CAR file from a PDS for a given identity diff --git a/services/cadet/src/ingestors/teal/feed_play.rs b/services/cadet/src/ingestors/teal/feed_play.rs index ef62c60..11dc3ec 100644 --- a/services/cadet/src/ingestors/teal/feed_play.rs +++ b/services/cadet/src/ingestors/teal/feed_play.rs @@ -1463,19 +1463,20 @@ impl PlayIngestor { .await?; } - // Refresh materialized views concurrently (if needed, consider if this should be done less frequently) - sqlx::query!("REFRESH MATERIALIZED VIEW mv_artist_play_counts;") - .execute(&self.sql) - .await?; - sqlx::query!("REFRESH MATERIALIZED VIEW mv_release_play_counts;") - .execute(&self.sql) - .await?; - sqlx::query!("REFRESH MATERIALIZED VIEW mv_recording_play_counts;") - .execute(&self.sql) - .await?; - sqlx::query!("REFRESH MATERIALIZED VIEW mv_global_play_count;") - .execute(&self.sql) - .await?; + if std::env::var("CADET_DEFER_MATERIALIZED_VIEW_REFRESH").as_deref() != Ok("1") { + sqlx::query!("REFRESH MATERIALIZED VIEW mv_artist_play_counts;") + .execute(&self.sql) + .await?; + sqlx::query!("REFRESH MATERIALIZED VIEW mv_release_play_counts;") + .execute(&self.sql) + .await?; + sqlx::query!("REFRESH MATERIALIZED VIEW mv_recording_play_counts;") + .execute(&self.sql) + .await?; + sqlx::query!("REFRESH MATERIALIZED VIEW mv_global_play_count;") + .execute(&self.sql) + .await?; + } // // Optionally check materialised views (consider removing in production for performance) // // For debugging purposes, can keep for now