debugage lazy cache

This commit is contained in:
2025-12-15 15:05:35 +01:00
parent d2d8111668
commit 68e6f528e5
18 changed files with 641 additions and 194 deletions

View File

@@ -8,17 +8,23 @@ use crate::db::DB;
use crate::download::{
download_with_transformer, ingest_with_transformer, Download, StreamTransformer,
};
use crate::lazy::{lazy_prefix_from_pk, LazyEntryRemoteData, LazyProvider};
use anyhow::{anyhow, bail, Result};
use serde_json::{Number, Value};
use sha2::{Digest, Sha256};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::sync::{Arc, RwLock as StdRwLock};
use tokio::io::{AsyncRead, AsyncReadExt};
use tokio::sync::{broadcast, RwLock};
use tracing;
enum FinalizeMode<'a> {
InsertNew,
ConvertLazy { lazy_pk: &'a str },
}
// ============================================================================
// LAZY PK SUPPORT
// ============================================================================
@@ -41,7 +47,11 @@ pub fn generate_lazy_pk(url: &str) -> String {
/// Vérifie si un PK est en mode lazy
pub fn is_lazy_pk(pk: &str) -> bool {
pk.starts_with(LAZY_PK_PREFIX)
if pk.starts_with(LAZY_PK_PREFIX) {
return true;
}
lazy_prefix_from_pk(pk).is_some()
}
/// Events émis par le cache pour notifier les changements d'état
@@ -136,6 +146,8 @@ pub struct Cache<C: CacheConfig> {
min_prebuffer_size: u64,
/// LAZY PK SUPPORT: Channel pour broadcaster les events (lazy downloads, etc.)
served_tx: Option<broadcast::Sender<CacheEvent>>,
/// Providers responsables de préfixes lazy spécifiques
lazy_providers: StdRwLock<HashMap<String, Arc<dyn LazyProvider>>>,
/// Phantom data pour le type de configuration
_phantom: std::marker::PhantomData<C>,
}
@@ -242,6 +254,7 @@ impl<C: CacheConfig> Cache<C> {
download: Arc<Download>,
collection: Option<&str>,
origin_url: Option<&str>,
mode: FinalizeMode<'_>,
) -> Result<String> {
// Attendre le prébuffering (pour le cache progressif)
if self.min_prebuffer_size > 0 {
@@ -256,10 +269,17 @@ impl<C: CacheConfig> Cache<C> {
);
}
// Ajouter à la DB une fois le prébuffer terminé
self.db.add(pk, None, collection)?;
if let Some(url) = origin_url {
self.db.set_origin_url(pk, url)?;
// Ajouter ou commuter la DB selon le mode
match mode {
FinalizeMode::InsertNew => {
self.db.add(pk, None, collection)?;
if let Some(url) = origin_url {
self.db.set_origin_url(pk, url)?;
}
}
FinalizeMode::ConvertLazy { lazy_pk } => {
self.db.update_lazy_to_downloaded(lazy_pk, pk)?;
}
}
// Sauvegarder les métadonnées techniques du transformer
@@ -400,6 +420,7 @@ impl<C: CacheConfig> Cache<C> {
transformer_factory,
min_prebuffer_size: DEFAULT_PREBUFFER_SIZE,
served_tx: Some(served_tx),
lazy_providers: StdRwLock::new(HashMap::new()),
_phantom: std::marker::PhantomData,
})
}
@@ -432,6 +453,34 @@ impl<C: CacheConfig> Cache<C> {
self.min_prebuffer_size
}
/// Enregistre un provider responsable d'un préfixe de lazy PK.
pub fn register_lazy_provider(&self, provider: Arc<dyn LazyProvider>) {
let prefix = provider.lazy_prefix().to_string();
let mut guard = self
.lazy_providers
.write()
.expect("lazy provider registry poisoned");
guard.insert(prefix, provider);
}
/// Désenregistre un provider à partir de son préfixe.
pub fn unregister_lazy_provider(&self, prefix: &str) {
let mut guard = self
.lazy_providers
.write()
.expect("lazy provider registry poisoned");
guard.remove(prefix);
}
fn provider_for_lazy_pk(&self, lazy_pk: &str) -> Option<Arc<dyn LazyProvider>> {
let prefix = lazy_prefix_from_pk(lazy_pk)?;
let guard = self
.lazy_providers
.read()
.expect("lazy provider registry poisoned");
guard.get(prefix).cloned()
}
/// S'abonne aux diffusions HTTP pour un `pk` donné.
///
/// La callback est appelée à chaque fois qu'un élément est servi avec succès via les routes
@@ -611,10 +660,63 @@ impl<C: CacheConfig> Cache<C> {
}
// Finaliser avec prébuffering et nettoyage
self.finalize_download(&pk, download, collection, Some(url))
self.finalize_download(&pk, download, collection, Some(url), FinalizeMode::InsertNew)
.await
}
/// Télécharge un fichier lazy et commute l'entrée existante
pub async fn download_lazy_from_url(
&self,
lazy_pk: &str,
url: &str,
collection: Option<&str>,
) -> Result<String> {
// Si déjà converti, retourner directement
if let Ok(Some(real_pk)) = self.db.get_pk_by_lazy_pk(lazy_pk) {
return Ok(real_pk);
}
// Si l'URL pointe déjà vers un fichier complet, commuter sans re-télécharger
if let Ok(Some(existing_pk)) = self.db.get_pk_by_origin_url(url) {
if existing_pk != lazy_pk && self.check_cached_and_complete(&existing_pk).await? {
self.db.update_lazy_to_downloaded(lazy_pk, &existing_pk)?;
return Ok(existing_pk);
}
}
// 1. Télécharger les 2048 premiers octets pour calculer le pk
let header = crate::download::peek_header(url, 2048)
.await
.map_err(|e| anyhow!("Failed to peek header: {}", e))?;
let pk = crate::cache_trait::pk_from_content_header(&header);
// 2. Si déjà en cache (complet), commuter directement
if self.check_cached_and_complete(&pk).await? {
self.db.update_lazy_to_downloaded(lazy_pk, &pk)?;
return Ok(pk);
}
// 3. Lancer le téléchargement complet
tracing::debug!("Starting lazy download for pk {} (lazy {})", pk, lazy_pk);
let file_path = self.get_file_path(&pk);
let transformer = self.transformer_factory.as_ref().map(|f| f());
let download = download_with_transformer(&file_path, url, transformer);
{
let mut downloads = self.downloads.write().await;
downloads.insert(pk.clone(), download.clone());
}
self.finalize_download(
&pk,
download,
collection,
Some(url),
FinalizeMode::ConvertLazy { lazy_pk },
)
.await
}
/// Ajoute un fichier à partir d'un flux asynchrone.
///
/// Cette méthode utilise le même système d'identifiants basé sur le contenu que `add_from_url`.
@@ -725,7 +827,7 @@ impl<C: CacheConfig> Cache<C> {
}
// Finaliser avec prébuffering et nettoyage
self.finalize_download(&pk, download, collection, source_uri)
self.finalize_download(&pk, download, collection, source_uri, FinalizeMode::InsertNew)
.await
}
@@ -1297,71 +1399,95 @@ impl<C: CacheConfig> Cache<C> {
}
}
/// Ajoute une URL en mode deferred (pas de download immédiat)
///
/// Vérifie d'abord si l'URL existe déjà en DB :
/// - Si eager (déjà téléchargé) : retourne le lazy_pk si existe, sinon le real pk
/// - Si lazy (pas encore téléchargé) : retourne le lazy pk existant
/// - Sinon : crée nouvelle entry lazy
///
/// # Arguments
///
/// * `url` - URL à cacher
/// * `collection` - Collection optionnelle
///
/// # Returns
///
/// PK (lazy ou real selon l'état)
///
/// # Example
///
/// ```rust,no_run
/// let cache = Cache::<AudioConfig>::new("./cache", 1000)?;
/// let lazy_pk = cache.add_from_url_deferred("https://example.com/track.mp3", Some("qobuz")).await?;
/// // → Returns "L:abc123..." (lazy PK)
/// // Fichier pas encore téléchargé, juste métadonnées en DB
/// ```
/// Garantit l'existence d'une entrée lazy spécifique.
pub async fn ensure_lazy_entry(
&self,
lazy_pk: &str,
collection: Option<&str>,
origin_url: Option<&str>,
) -> Result<()> {
if let Ok(true) = self.db.has_lazy_entry(lazy_pk) {
self.db.update_hit_by_lazy_pk(lazy_pk)?;
} else {
self.db.add_lazy(lazy_pk, None, collection)?;
}
if let Some(url) = origin_url {
self.db.set_origin_url_for_lazy(lazy_pk, url)?;
}
Ok(())
}
/// Récupère auprès du provider les métadonnées/couvertures associées.
pub async fn fetch_lazy_provider_data(&self, lazy_pk: &str) -> Result<LazyEntryRemoteData> {
if let Some(provider) = self.provider_for_lazy_pk(lazy_pk) {
let metadata = provider.metadata(lazy_pk).await?;
let cover_url = provider.cover_url(lazy_pk).await?;
Ok(LazyEntryRemoteData { metadata, cover_url })
} else {
Ok(LazyEntryRemoteData::default())
}
}
/// Résout l'URL d'origine pour un lazy PK, via la DB ou un provider.
pub async fn resolve_lazy_url(&self, lazy_pk: &str) -> Result<String> {
if let Ok(Some(url)) = self.db.get_origin_url(lazy_pk) {
return Ok(url);
}
if let Some(provider) = self.provider_for_lazy_pk(lazy_pk) {
return provider.get_url(lazy_pk).await;
}
bail!("No origin URL or lazy provider registered for {}", lazy_pk);
}
/// Télécharge un lazy PK en résolvant automatiquement son URL.
pub async fn download_lazy(&self, lazy_pk: &str, collection: Option<&str>) -> Result<String> {
let (existing_pk, _lazy_sec, existing_collection) = self
.db
.get_entry_by_pk_or_lazy_pk(lazy_pk)?
.ok_or_else(|| anyhow!("Lazy pk {} not found in DB", lazy_pk))?;
let url = self.resolve_lazy_url(lazy_pk).await?;
let collection = collection.or(existing_collection.as_deref());
if let Some(real_pk) = existing_pk {
if real_pk != lazy_pk && self.check_cached_and_complete(&real_pk).await? {
return Ok(real_pk);
}
}
self.download_lazy_from_url(lazy_pk, &url, collection).await
}
/// Ajoute une URL sans lancer immédiatement le téléchargement.
pub async fn add_from_url_deferred(
&self,
url: &str,
collection: Option<&str>,
) -> Result<String> {
// 1. Vérifier si URL déjà en cache
if let Ok(Some((pk_opt, lazy_pk_opt))) = self.db.get_entry_by_url(url) {
// URL existe déjà
if let Some(pk) = pk_opt {
// Fichier déjà téléchargé (eager ou lazy→eager)
tracing::debug!("URL {} already downloaded with pk {}", url, pk);
self.db.update_hit(&pk)?;
// Retourner lazy_pk si existe (pour compatibilité Control Point)
// sinon retourner pk
if let Some(lpk) = lazy_pk_opt {
return Ok(lpk);
}
return Ok(pk);
} else if let Some(lpk) = lazy_pk_opt {
// Entry lazy existante (pas encore téléchargé)
tracing::debug!("URL {} already in lazy mode with pk {}", url, lpk);
self.db.update_hit_by_lazy_pk(&lpk)?;
return Ok(lpk);
}
}
// 2. URL inconnue → créer nouvelle entry lazy
let lazy_pk = generate_lazy_pk(url);
// Vérifier si ce lazy_pk existe déjà (collision improbable mais...)
let lazy_pk = format!("L:{}", generate_lazy_pk(url));
if let Ok(true) = self.db.has_lazy_entry(&lazy_pk) {
bail!("Lazy PK collision for URL: {}", url);
}
// 3. Ajouter en DB
self.db.add_lazy(&lazy_pk, None, collection)?;
self.db.set_origin_url_for_lazy(&lazy_pk, url)?;
self.ensure_lazy_entry(&lazy_pk, collection, Some(url)).await?;
tracing::debug!("Created new lazy pk {} for URL {}", lazy_pk, url);
Ok(lazy_pk)
}
}