From d6920703a110d0214129e3efd428e746c4954fdb Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Wed, 17 Dec 2025 07:25:47 +0100 Subject: [PATCH] =?UTF-8?q?Am=C3=A9lioration=20des=20pmoplaylist?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pmoplaylist/src/manager.rs | 95 +++++++++++++++++++++++++++++++++++--- pmoserver/src/server.rs | 68 ++++++++++++++------------- 2 files changed, 124 insertions(+), 39 deletions(-) diff --git a/pmoplaylist/src/manager.rs b/pmoplaylist/src/manager.rs index 0db151d5..cffe0ff8 100644 --- a/pmoplaylist/src/manager.rs +++ b/pmoplaylist/src/manager.rs @@ -6,12 +6,12 @@ use crate::playlist::core::PlaylistConfig; use crate::playlist::Playlist; use crate::Result; use once_cell::sync::OnceCell; -use pmocache::{CacheBroadcastEvent, CacheSubscription}; +use pmocache::{CacheBroadcastEvent, CacheEvent, CacheSubscription}; use std::collections::HashMap; use std::path::PathBuf; use std::sync::RwLock as StdRwLock; use std::sync::{ - atomic::{AtomicU64, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, Arc, }; use std::time::Duration; @@ -33,6 +33,7 @@ struct ManagerInner { track_index: StdRwLock>>, // cache_pk -> playlists cache_subscriptions: StdRwLock>, event_tx: broadcast::Sender, + lazy_listener_started: AtomicBool, } /// Type d'évènement émis par le PlaylistManager. @@ -87,14 +88,24 @@ impl PlaylistManager { track_index: StdRwLock::new(HashMap::new()), cache_subscriptions: StdRwLock::new(HashMap::new()), event_tx: broadcast::channel(256).0, + lazy_listener_started: AtomicBool::new(false), }), }; // Lancer la task d'�viction en background - let manager_clone = manager.clone(); - tokio::spawn(async move { - manager_clone.eviction_task().await; - }); + { + let manager_clone = manager.clone(); + tokio::spawn(async move { + manager_clone.eviction_task().await; + }); + } + + { + let manager_clone = manager.clone(); + tokio::spawn(async move { + manager_clone.ensure_lazy_listener().await; + }); + } Ok(manager) } @@ -667,6 +678,77 @@ impl PlaylistManager { } } + async fn ensure_lazy_listener(&self) { + if self.inner.lazy_listener_started.load(Ordering::SeqCst) { + return; + } + + let cache = match audio_cache() { + Ok(cache) => cache, + Err(_) => return, + }; + + if self + .inner + .lazy_listener_started + .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) + .is_err() + { + return; + } + + let manager = self.clone(); + let inner = self.inner.clone(); + tokio::spawn(async move { + let mut rx = cache.subscribe_events(); + while let Ok(event) = rx.recv().await { + if let CacheEvent::LazyDownloaded { lazy_pk, real_pk } = event { + manager + .handle_lazy_download_event(&lazy_pk, &real_pk) + .await; + } + } + inner.lazy_listener_started.store(false, Ordering::SeqCst); + }); + } + + async fn handle_lazy_download_event(&self, lazy_pk: &str, real_pk: &str) { + let playlists = { + let index = self.inner.track_index.read().unwrap(); + index.get(lazy_pk).cloned() + }; + + let Some(playlists) = playlists else { + tracing::debug!( + "Lazy download {} converted to {} but no playlists referenced it", + lazy_pk, + real_pk + ); + return; + }; + + for playlist_id in playlists { + match self.get_write_handle(playlist_id.clone()).await { + Ok(writer) => { + if let Err(e) = writer.update_cache_pk(lazy_pk, real_pk).await { + tracing::error!( + "Failed to update playlist {} from {} to {}: {}", + playlist_id, + lazy_pk, + real_pk, + e + ); + } + } + Err(e) => tracing::debug!( + "Failed to acquire write handle for playlist {} during lazy swap: {}", + playlist_id, + e + ), + } + } + } + /// Task d'�viction en background async fn eviction_task(&self) { loop { @@ -739,6 +821,7 @@ pub fn register_audio_cache(cache: Arc) { let manager = manager.clone(); tokio::spawn(async move { manager.sync_cache_subscriptions().await; + manager.ensure_lazy_listener().await; }); } } diff --git a/pmoserver/src/server.rs b/pmoserver/src/server.rs index 859d2cbc..f7fc0d6c 100644 --- a/pmoserver/src/server.rs +++ b/pmoserver/src/server.rs @@ -28,7 +28,7 @@ use std::net::SocketAddr; use std::sync::Arc; use tokio::{signal, sync::RwLock, task::JoinHandle}; use tokio_util::sync::CancellationToken; -use tracing::info; +use tracing::{error, info, warn}; use utoipa::OpenApi; use utoipa_swagger_ui::SwaggerUi; @@ -543,43 +543,45 @@ impl Server { ); let router = self.router.clone(); + let shutdown_token = self.shutdown_token.clone(); // Créer un channel pour signaler l'arrêt gracieux let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>(); - // Tâche qui écoute Ctrl+C et signale l'arrêt - let shutdown_token_for_signal = self.shutdown_token.clone(); - let shutdown_signal = tokio::spawn(async move { - signal::ctrl_c().await.expect("failed to listen for ctrl_c"); - info!("Ctrl+C reçu, arrêt gracieux"); - // Déclencher le token d'arrêt pour tous les composants - shutdown_token_for_signal.cancel(); - // Envoyer le signal d'arrêt au serveur HTTP (ignorer l'erreur si le receiver a déjà été drop) - let _ = shutdown_tx.send(()); - }); - - // Tâche serveur avec arrêt gracieux - let server_task = tokio::spawn(async move { - let r = router.read().await.clone(); - let listener = tokio::net::TcpListener::bind(addr).await.unwrap(); - - // Serveur avec arrêt gracieux - axum::serve(listener, r.into_make_service()) - .with_graceful_shutdown(async move { - // Attendre le signal d'arrêt - let _ = shutdown_rx.await; - }) - .await - .unwrap(); - - info!("Serveur HTTP arrêté proprement"); - }); - self.join_handle = Some(tokio::spawn(async move { - // Attendre que le serveur se termine (après réception du signal) - let _ = server_task.await; - // Nettoyer la tâche de signal - let _ = shutdown_signal.await; + let server_future = async { + let r = router.read().await.clone(); + let listener = tokio::net::TcpListener::bind(addr).await.unwrap(); + + axum::serve(listener, r.into_make_service()) + .with_graceful_shutdown(async move { + let _ = shutdown_rx.await; + }) + .await + }; + + tokio::pin!(server_future); + let ctrl_c = signal::ctrl_c(); + tokio::pin!(ctrl_c); + + tokio::select! { + result = &mut server_future => { + if let Err(err) = result { + error!("Serveur HTTP arrêté avec une erreur: {}", err); + } else { + info!("Serveur HTTP arrêté proprement"); + } + } + _ = &mut ctrl_c => { + info!("Ctrl+C reçu, arrêt gracieux"); + shutdown_token.cancel(); + let _ = shutdown_tx.send(()); + + if tokio::time::timeout(std::time::Duration::from_secs(5), &mut server_future).await.is_err() { + warn!("Arrêt gracieux trop long, fermeture forcée du serveur HTTP"); + } + } + } })); }