From 2e8dc489203e96e5e74d00c797af4c24645cfdb3 Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Mon, 5 Jan 2026 20:54:18 +0300 Subject: [PATCH] spotify: Add resilient user threads with recovery and backoff Introduce start_user_thread helper to spawn per-user threads with automatic retries, exponential backoff, and NATS publish after max retries for external restart. Replace direct thread spawns with the helper. Improve progress tracker by adding an inner recovery loop, fixing cache update logic, and handling Redis errors gracefully. --- crates/spotify/src/lib.rs | 323 ++++++++++++++++++++------------------ 1 file changed, 174 insertions(+), 149 deletions(-) diff --git a/crates/spotify/src/lib.rs b/crates/spotify/src/lib.rs index 1e23b2b4..2c942de3 100644 --- a/crates/spotify/src/lib.rs +++ b/crates/spotify/src/lib.rs @@ -53,6 +53,97 @@ pub async fn run() -> Result<(), Error> { let thread_map: Arc>>> = Arc::new(Mutex::new(HashMap::new())); + // Helper function to start a user thread with auto-recovery + let start_user_thread = |email: String, + token: String, + did: String, + client_id: String, + client_secret: String, + stop_flag: Arc, + cache: Cache, + nc: async_nats::Client| { + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + let mut retry_count = 0; + let max_retries = 5; + + loop { + if stop_flag.load(std::sync::atomic::Ordering::Relaxed) { + println!( + "{} Stop flag set, exiting recovery loop", + format!("[{}]", email).bright_green() + ); + break; + } + + match rt.block_on(async { + watch_currently_playing( + email.clone(), + token.clone(), + did.clone(), + stop_flag.clone(), + cache.clone(), + client_id.clone(), + client_secret.clone(), + ) + .await + }) { + Ok(_) => { + println!( + "{} Thread completed normally", + format!("[{}]", email).bright_green() + ); + break; + } + Err(e) => { + retry_count += 1; + println!( + "{} Thread crashed (attempt {}/{}): {}", + format!("[{}]", email).bright_green(), + retry_count, + max_retries, + e.to_string().bright_red() + ); + + if retry_count >= max_retries { + println!( + "{} Max retries reached, publishing to NATS for external restart", + format!("[{}]", email).bright_green() + ); + match rt + .block_on(nc.publish("rocksky.spotify.user", email.clone().into())) + { + Ok(_) => { + println!( + "{} Published message to restart thread", + format!("[{}]", email).bright_green() + ); + } + Err(e) => { + println!( + "{} Error publishing message to restart thread: {}", + format!("[{}]", email).bright_green(), + e.to_string().bright_red() + ); + } + } + break; + } + + // Exponential backoff: 2^retry_count seconds, max 60 seconds + let backoff_seconds = std::cmp::min(2_u64.pow(retry_count as u32), 60); + println!( + "{} Retrying in {} seconds...", + format!("[{}]", email).bright_green(), + backoff_seconds + ); + std::thread::sleep(std::time::Duration::from_secs(backoff_seconds)); + } + } + } + }) + }; + // Start threads for all users for user in users { let email = user.0.clone(); @@ -70,50 +161,16 @@ pub async fn run() -> Result<(), Error> { .unwrap() .insert(email.clone(), Arc::clone(&stop_flag)); - thread::spawn(move || { - let rt = tokio::runtime::Runtime::new().unwrap(); - match rt.block_on(async { - watch_currently_playing( - email.clone(), - token, - did, - stop_flag, - cache.clone(), - client_id, - client_secret, - ) - .await?; - Ok::<(), Error>(()) - }) { - Ok(_) => {} - Err(e) => { - println!( - "{} Error starting thread for user: {} - {}", - format!("[{}]", email).bright_green(), - email.bright_green(), - e.to_string().bright_red() - ); - - // If there's an error, publish a message to restart the thread - match rt.block_on(nc.publish("rocksky.spotify.user", email.clone().into())) { - Ok(_) => { - println!( - "{} Published message to restart thread for user: {}", - format!("[{}]", email).bright_green(), - email.bright_green() - ); - } - Err(e) => { - println!( - "{} Error publishing message to restart thread: {}", - format!("[{}]", email).bright_green(), - e.to_string().bright_red() - ); - } - } - } - } - }); + start_user_thread( + email, + token, + did, + client_id, + client_secret, + stop_flag, + cache, + nc, + ); } // Handle subscription messages @@ -153,33 +210,18 @@ pub async fn run() -> Result<(), Error> { let client_id = user.3.clone(); let client_secret = user.4.clone(); let cache = cache.clone(); + let nc = nc.clone(); - thread::spawn(move || { - let rt = tokio::runtime::Runtime::new().unwrap(); - match rt.block_on(async { - watch_currently_playing( - email.clone(), - token, - did, - new_stop_flag, - cache.clone(), - client_id, - client_secret, - ) - .await?; - Ok::<(), Error>(()) - }) { - Ok(_) => {} - Err(e) => { - println!( - "{} Error restarting thread for user: {} - {}", - format!("[{}]", email).bright_green(), - email.bright_green(), - e.to_string().bright_red() - ); - } - } - }); + start_user_thread( + email, + token, + did, + client_id, + client_secret, + new_stop_flag, + cache, + nc, + ); println!("Restarted thread for user: {}", user_id.bright_green()); } else { @@ -200,44 +242,16 @@ pub async fn run() -> Result<(), Error> { thread_map.insert(email.clone(), Arc::clone(&stop_flag)); - thread::spawn(move || { - let rt = tokio::runtime::Runtime::new().unwrap(); - match rt.block_on(async { - watch_currently_playing( - email.clone(), - token, - did, - stop_flag, - cache.clone(), - client_id, - client_secret, - ) - .await?; - Ok::<(), Error>(()) - }) { - Ok(_) => {} - Err(e) => { - println!( - "{} Error starting thread for user: {} - {}", - format!("[{}]", email).bright_green(), - email.bright_green(), - e.to_string().bright_red() - ); - match rt - .block_on(nc.publish("rocksky.spotify.user", email.clone().into())) - { - Ok(_) => {} - Err(e) => { - println!( - "{} Error publishing message to restart thread: {}", - format!("[{}]", email).bright_green(), - e.to_string().bright_red() - ); - } - } - } - } - }); + start_user_thread( + email, + token, + did, + client_id, + client_secret, + stop_flag, + cache, + nc, + ); } } } @@ -810,31 +824,54 @@ pub async fn watch_currently_playing( let spotify_email_clone = spotify_email.clone(); let cache_clone = cache.clone(); thread::spawn(move || { - loop { - if stop_flag_clone.load(std::sync::atomic::Ordering::Relaxed) { - println!( - "{} Stopping Thread", - format!("[{}]", spotify_email_clone).bright_green() - ); - break; - } - if let Ok(Some(cached)) = cache_clone.get(&format!("{}:current", spotify_email_clone)) { - if serde_json::from_str::(&cached).is_err() { - thread::sleep(std::time::Duration::from_millis(800)); - continue; + // Inner thread with error recovery + let result: Result<(), Error> = (|| { + loop { + if stop_flag_clone.load(std::sync::atomic::Ordering::Relaxed) { + println!( + "{} Stopping progress tracker thread", + format!("[{}]", spotify_email_clone).bright_green() + ); + break; } - let mut current_song = serde_json::from_str::(&cached)?; - - if let Some(item) = current_song.item.clone() { - if current_song.is_playing - && current_song.progress_ms.unwrap_or(0) < item.duration_ms.into() + if let Ok(Some(cached)) = + cache_clone.get(&format!("{}:current", spotify_email_clone)) + { + if let Ok(mut current_song) = serde_json::from_str::(&cached) { - current_song.progress_ms = - Some(current_song.progress_ms.unwrap_or(0) + 800); + if let Some(item) = current_song.item.clone() { + if current_song.is_playing + && current_song.progress_ms.unwrap_or(0) < item.duration_ms.into() + { + current_song.progress_ms = + Some(current_song.progress_ms.unwrap_or(0) + 800); + match cache_clone.setex( + &format!("{}:current", spotify_email_clone), + &serde_json::to_string(¤t_song).unwrap_or_default(), + 16, + ) { + Ok(_) => {} + Err(e) => { + println!( + "{} redis error: {}", + format!("[{}]", spotify_email_clone).bright_green(), + e.to_string().bright_red() + ); + } + } + thread::sleep(std::time::Duration::from_millis(800)); + continue; + } + } + } + } + + if let Ok(Some(cached)) = cache_clone.get(&spotify_email_clone) { + if cached != "No content" { match cache_clone.setex( &format!("{}:current", spotify_email_clone), - &serde_json::to_string(¤t_song)?, + &cached, 16, ) { Ok(_) => {} @@ -846,33 +883,21 @@ pub async fn watch_currently_playing( ); } } - thread::sleep(std::time::Duration::from_millis(800)); - continue; } } - continue; - } - if let Ok(Some(cached)) = cache_clone.get(&spotify_email_clone) { - if cached == "No content" { - thread::sleep(std::time::Duration::from_millis(800)); - continue; - } - match cache_clone.setex(&format!("{}:current", spotify_email_clone), &cached, 16) { - Ok(_) => {} - Err(e) => { - println!( - "{} redis error: {}", - format!("[{}]", spotify_email_clone).bright_green(), - e.to_string().bright_red() - ); - } - } + thread::sleep(std::time::Duration::from_millis(800)); } + Ok(()) + })(); - thread::sleep(std::time::Duration::from_millis(800)); + if let Err(e) = result { + println!( + "{} Progress tracker thread error: {}", + format!("[{}]", spotify_email_clone).bright_green(), + e.to_string().bright_red() + ); } - Ok::<(), Error>(()) }); loop { -- 2.51.2