Refactoring du cache pour une meilleur gestion des metadonnées

This commit is contained in:
2025-10-26 09:18:48 +01:00
parent e3c17bf131
commit fc093743c1
25 changed files with 701 additions and 277 deletions

1
Cargo.lock generated
View File

@@ -2392,6 +2392,7 @@ dependencies = [
"reqwest", "reqwest",
"rusqlite", "rusqlite",
"serde", "serde",
"serde_json",
"serde_yaml", "serde_yaml",
"sha1", "sha1",
"sha2", "sha2",

View File

@@ -6,6 +6,7 @@
use anyhow::Result; use anyhow::Result;
use pmocache::{CacheConfig, StreamTransformer}; use pmocache::{CacheConfig, StreamTransformer};
use serde_json::Value;
use std::sync::Arc; use std::sync::Arc;
/// Configuration pour le cache audio /// Configuration pour le cache audio
@@ -16,10 +17,6 @@ impl CacheConfig for AudioConfig {
"flac" "flac"
} }
fn table_name() -> &'static str {
"audio_tracks"
}
fn cache_type() -> &'static str { fn cache_type() -> &'static str {
"flac" "flac"
} }
@@ -147,17 +144,16 @@ fn create_flac_transformer() -> StreamTransformer {
.channels .channels
.ok_or_else(|| { .ok_or_else(|| {
tracing::error!("Audio file missing channel information"); tracing::error!("Audio file missing channel information");
"Audio file is missing channel information. The file may be corrupted.".to_string() "Audio file is missing channel information. The file may be corrupted."
.to_string()
})? })?
.count(); .count();
let sample_rate = track let sample_rate = track.codec_params.sample_rate.ok_or_else(|| {
.codec_params tracing::error!("Audio file missing sample rate information");
.sample_rate "Audio file is missing sample rate information. The file may be corrupted."
.ok_or_else(|| { .to_string()
tracing::error!("Audio file missing sample rate information"); })?;
"Audio file is missing sample rate information. The file may be corrupted.".to_string()
})?;
let bits_per_sample = track.codec_params.bits_per_sample.unwrap_or(16); let bits_per_sample = track.codec_params.bits_per_sample.unwrap_or(16);
@@ -179,7 +175,10 @@ fn create_flac_transformer() -> StreamTransformer {
} }
Err(e) => { Err(e) => {
tracing::error!("Failed to read audio packet: {}", e); tracing::error!("Failed to read audio packet: {}", e);
return Err(format!("Failed to read audio data: {}. The file may be corrupted.", e)); return Err(format!(
"Failed to read audio data: {}. The file may be corrupted.",
e
));
} }
}; };
@@ -212,7 +211,10 @@ fn create_flac_transformer() -> StreamTransformer {
if samples_i32.is_empty() { if samples_i32.is_empty() {
tracing::error!("No audio samples could be decoded from the file"); tracing::error!("No audio samples could be decoded from the file");
return Err("No audio samples could be decoded. The file may be corrupted or empty.".to_string()); return Err(
"No audio samples could be decoded. The file may be corrupted or empty."
.to_string(),
);
} }
tracing::debug!( tracing::debug!(
@@ -287,12 +289,10 @@ fn create_flac_transformer() -> StreamTransformer {
})?; })?;
let mut sink = ByteSink::new(); let mut sink = ByteSink::new();
flac_stream flac_stream.write(&mut sink).map_err(|e| {
.write(&mut sink) tracing::error!("Failed to write FLAC stream: {:?}", e);
.map_err(|e| { format!("Failed to write FLAC data: {:?}", e)
tracing::error!("Failed to write FLAC stream: {:?}", e); })?;
format!("Failed to write FLAC data: {:?}", e)
})?;
Ok::<Vec<u8>, String>(sink.into_inner()) Ok::<Vec<u8>, String>(sink.into_inner())
}) })
@@ -389,14 +389,12 @@ pub async fn add_with_metadata_extraction(
// Extraire les métadonnées // Extraire les métadonnées
let metadata = crate::metadata::AudioMetadata::from_bytes(&flac_bytes)?; let metadata = crate::metadata::AudioMetadata::from_bytes(&flac_bytes)?;
let metadata_json: Value = serde_json::to_value(&metadata)
// Sérialiser en JSON .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
let metadata_json = serde_json::to_string(&metadata)?;
// Stocker dans la DB // Stocker dans la DB
cache cache
.db .db
.update_metadata(&pk, &metadata_json) .set_metadata(&pk, &metadata_json)
.map_err(|e| anyhow::anyhow!("Database error: {}", e))?; .map_err(|e| anyhow::anyhow!("Database error: {}", e))?;
// Mettre à jour la collection si les métadonnées en fournissent une // Mettre à jour la collection si les métadonnées en fournissent une
@@ -404,8 +402,9 @@ pub async fn add_with_metadata_extraction(
if let Some(auto_collection) = metadata.collection_key() { if let Some(auto_collection) = metadata.collection_key() {
cache cache
.db .db
.add(&pk, url, Some(&auto_collection)) .add(&pk, None, Some(&auto_collection))
.map_err(|e| anyhow::anyhow!("Database error: {}", e))?; .map_err(|e| anyhow::anyhow!("Database error: {}", e))?;
cache.db.set_origin_url(&pk, url)?;
} }
} }

View File

@@ -4,8 +4,8 @@
//! des méthodes de gestion du cache audio à pmoconfig::Config. //! des méthodes de gestion du cache audio à pmoconfig::Config.
use anyhow::Result; use anyhow::Result;
use pmoconfig::Config;
use pmocache::CacheConfigExt; use pmocache::CacheConfigExt;
use pmoconfig::Config;
use std::sync::Arc; use std::sync::Arc;
const DEFAULT_AUDIO_CACHE_DIR: &str = "cache_audio"; const DEFAULT_AUDIO_CACHE_DIR: &str = "cache_audio";

View File

@@ -20,6 +20,7 @@ hex = "0.4"
anyhow = "1.0" anyhow = "1.0"
chrono = "0.4" chrono = "0.4"
serde = { version = "1.0", features = ["derive"] } serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
bytes = "1.6" bytes = "1.6"
# Async # Async

View File

@@ -93,7 +93,7 @@ pub struct ErrorResponse {
/// ///
/// Retourne la liste complète des entrées du cache triées par nombre d'accès décroissant. /// Retourne la liste complète des entrées du cache triées par nombre d'accès décroissant.
pub async fn list_items<C: CacheConfig>(State(cache): State<Arc<Cache<C>>>) -> impl IntoResponse { pub async fn list_items<C: CacheConfig>(State(cache): State<Arc<Cache<C>>>) -> impl IntoResponse {
match cache.db.get_all() { match cache.db.get_all(true) {
Ok(entries) => (StatusCode::OK, Json(entries)).into_response(), Ok(entries) => (StatusCode::OK, Json(entries)).into_response(),
Err(e) => ( Err(e) => (
StatusCode::INTERNAL_SERVER_ERROR, StatusCode::INTERNAL_SERVER_ERROR,
@@ -113,7 +113,7 @@ pub async fn get_item_info<C: CacheConfig>(
State(cache): State<Arc<Cache<C>>>, State(cache): State<Arc<Cache<C>>>,
Path(pk): Path<String>, Path(pk): Path<String>,
) -> impl IntoResponse { ) -> impl IntoResponse {
match cache.db.get(&pk) { match cache.db.get(&pk, true) {
Ok(entry) => (StatusCode::OK, Json(entry)).into_response(), Ok(entry) => (StatusCode::OK, Json(entry)).into_response(),
Err(_) => ( Err(_) => (
StatusCode::NOT_FOUND, StatusCode::NOT_FOUND,
@@ -135,7 +135,7 @@ pub async fn get_download_status<C: CacheConfig>(
Path(pk): Path<String>, Path(pk): Path<String>,
) -> impl IntoResponse { ) -> impl IntoResponse {
// Vérifier que l'item existe dans la DB // Vérifier que l'item existe dans la DB
if cache.db.get(&pk).is_err() { if cache.db.get(&pk, false).is_err() {
return ( return (
StatusCode::NOT_FOUND, StatusCode::NOT_FOUND,
Json(ErrorResponse { Json(ErrorResponse {
@@ -222,7 +222,7 @@ pub async fn delete_item<C: CacheConfig>(
Path(pk): Path<String>, Path(pk): Path<String>,
) -> impl IntoResponse { ) -> impl IntoResponse {
// Vérifier que l'item existe // Vérifier que l'item existe
if cache.db.get(&pk).is_err() { if cache.db.get(&pk, false).is_err() {
return ( return (
StatusCode::NOT_FOUND, StatusCode::NOT_FOUND,
Json(ErrorResponse { Json(ErrorResponse {

View File

@@ -21,9 +21,7 @@ pub trait CacheConfig: Send + Sync {
/// Extension des fichiers (ex: "webp", "flac") /// Extension des fichiers (ex: "webp", "flac")
fn file_extension() -> &'static str; fn file_extension() -> &'static str;
/// Nom de la table dans la base de données (ex: "covers", "audio") /// Nom de la table dans la base de données (ex: "covers", "audio")
fn table_name() -> &'static str {
"cached_items"
}
/// Type de cache (ex: "audio", "image") /// Type de cache (ex: "audio", "image")
fn cache_type() -> &'static str { fn cache_type() -> &'static str {
"file" "file"
@@ -118,7 +116,7 @@ impl<C: CacheConfig> Cache<C> {
) -> Result<Self> { ) -> Result<Self> {
let directory = PathBuf::from(dir); let directory = PathBuf::from(dir);
std::fs::create_dir_all(&directory)?; std::fs::create_dir_all(&directory)?;
let db = DB::init(&directory.join("cache.db"), C::table_name())?; let db = DB::init(&directory.join("cache.db"))?;
Ok(Self { Ok(Self {
dir: directory, dir: directory,
@@ -169,7 +167,7 @@ impl<C: CacheConfig> Cache<C> {
tracing::debug!("Computed pk {} for URL {}", pk, url); tracing::debug!("Computed pk {} for URL {}", pk, url);
// 3. Vérifier si le fichier est déjà en cache // 3. Vérifier si le fichier est déjà en cache
if self.db.get(&pk).is_ok() { if self.db.get(&pk, false).is_ok() {
let file_path = self.file_path(&pk); let file_path = self.file_path(&pk);
if file_path.exists() { if file_path.exists() {
// Déjà en cache, update timestamp et retour rapide // Déjà en cache, update timestamp et retour rapide
@@ -202,8 +200,8 @@ impl<C: CacheConfig> Cache<C> {
} }
// Ajouter immédiatement à la DB // Ajouter immédiatement à la DB
self.db.add(&pk, url, collection)?; self.db.add(&pk, None, collection)?;
self.db.set_origin_url(&pk, url)?;
// Appliquer la politique d'éviction LRU si nécessaire // Appliquer la politique d'éviction LRU si nécessaire
if let Err(e) = self.enforce_limit().await { if let Err(e) = self.enforce_limit().await {
tracing::warn!("Error enforcing cache limit: {}", e); tracing::warn!("Error enforcing cache limit: {}", e);
@@ -264,7 +262,7 @@ impl<C: CacheConfig> Cache<C> {
tracing::debug!("Computed pk {} for source_uri {}", pk, source_uri); tracing::debug!("Computed pk {} for source_uri {}", pk, source_uri);
// 3. Vérifier si le fichier est déjà en cache // 3. Vérifier si le fichier est déjà en cache
if self.db.get(&pk).is_ok() { if self.db.get(&pk, false).is_ok() {
let file_path = self.file_path(&pk); let file_path = self.file_path(&pk);
if file_path.exists() { if file_path.exists() {
// Déjà en cache, update timestamp et retour rapide // Déjà en cache, update timestamp et retour rapide
@@ -301,7 +299,8 @@ impl<C: CacheConfig> Cache<C> {
downloads.insert(pk.clone(), download.clone()); downloads.insert(pk.clone(), download.clone());
} }
self.db.add(&pk, source_uri, collection)?; self.db.add(&pk, None, collection)?;
self.db.set_origin_url(&pk, source_uri);
if let Err(e) = self.enforce_limit().await { if let Err(e) = self.enforce_limit().await {
tracing::warn!("Error enforcing cache limit: {}", e); tracing::warn!("Error enforcing cache limit: {}", e);
@@ -359,13 +358,59 @@ impl<C: CacheConfig> Cache<C> {
.await .await
} }
pub async fn delete_item(&self, pk: &str) -> Result<()> {
// Vérifie l'existence pour signaler une erreur explicite si l'entrée est absente
self.db.get(pk, false)?;
// Oublie un téléchargement en cours pour cette clé
self.downloads.write().await.remove(pk);
// Supprime chaque fichier {pk}.{qualifier}.{ext} (ignorer si déjà absent)
for path in self.get_file_paths(pk)? {
if let Err(err) = tokio::fs::remove_file(&path).await {
if err.kind() != std::io::ErrorKind::NotFound {
return Err(err.into());
}
}
}
// Efface lentrée de la base (les métadonnées partent via ON DELETE CASCADE)
self.db.delete(pk)?;
Ok(())
}
pub async fn delete_collection(&self, collection: &str) -> Result<()> {
let entries = self.db.get_by_collection(collection, false)?;
{
let mut downloads = self.downloads.write().await;
for entry in &entries {
downloads.remove(&entry.pk);
}
}
for entry in &entries {
for path in self.get_file_paths(&entry.pk)? {
if let Err(err) = tokio::fs::remove_file(&path).await {
if err.kind() != std::io::ErrorKind::NotFound {
return Err(err.into());
}
}
}
}
self.db.delete_collection(collection)?;
Ok(())
}
/// Récupère le chemin d'un fichier dans le cache /// Récupère le chemin d'un fichier dans le cache
/// ///
/// # Arguments /// # Arguments
/// ///
/// * `pk` - Clé primaire du fichier /// * `pk` - Clé primaire du fichier
pub async fn get(&self, pk: &str) -> Result<PathBuf> { pub async fn get(&self, pk: &str) -> Result<PathBuf> {
self.db.get(pk)?; self.db.get(pk, false)?;
self.db.update_hit(pk)?; self.db.update_hit(pk)?;
let file_path = self.file_path(pk); let file_path = self.file_path(pk);
@@ -376,13 +421,18 @@ impl<C: CacheConfig> Cache<C> {
} }
} }
pub async fn touch(&self, pk: &str) -> Result<()> {
self.db.update_hit(pk)?;
Ok(())
}
/// Récupère tous les fichiers d'une collection /// Récupère tous les fichiers d'une collection
/// ///
/// # Arguments /// # Arguments
/// ///
/// * `collection` - Identifiant de la collection /// * `collection` - Identifiant de la collection
pub async fn get_collection(&self, collection: &str) -> Result<Vec<PathBuf>> { pub async fn get_collection(&self, collection: &str) -> Result<Vec<PathBuf>> {
let entries = self.db.get_by_collection(collection)?; let entries = self.db.get_by_collection(collection, false)?;
let mut paths = Vec::new(); let mut paths = Vec::new();
for entry in entries { for entry in entries {
@@ -412,20 +462,26 @@ impl<C: CacheConfig> Cache<C> {
/// Consolide le cache en supprimant les orphelins et en re-téléchargeant les fichiers manquants /// Consolide le cache en supprimant les orphelins et en re-téléchargeant les fichiers manquants
pub async fn consolidate(&self) -> Result<()> { pub async fn consolidate(&self) -> Result<()> {
// Récupérer la liste des entrées à traiter // Récupérer la liste des entrées à traiter
let entries = self.db.get_all()?; let entries = self.db.get_all(false)?;
// Supprimer les entrées sans fichiers correspondants // Supprimer les entrées sans fichiers correspondants
for entry in entries { for entry in entries {
let file_path = self.file_path(&entry.pk); let file_path = self.file_path(&entry.pk);
if !file_path.exists() { if !file_path.exists() {
// Re-télécharger le fichier manquant match self.db.get_origin_url(&entry.pk)? {
match self Some(url) => {
.add_from_url(&entry.source_url, entry.collection.as_deref()) if let Err(err) = self.add_from_url(&url, entry.collection.as_deref()).await
.await {
{ tracing::warn!(
Ok(_) => {} "Unable to redownload missing file for {}: {}",
Err(_) => { entry.pk,
// Si le téléchargement échoue, supprimer l'entrée DB err
);
self.db.delete(&entry.pk)?;
}
}
None => {
self.db.delete(&entry.pk)?; self.db.delete(&entry.pk)?;
} }
} }
@@ -441,7 +497,7 @@ impl<C: CacheConfig> Cache<C> {
// Format attendu: {pk}.{qualifier}.{EXT} // Format attendu: {pk}.{qualifier}.{EXT}
// On extrait le pk (première partie avant le premier point) // On extrait le pk (première partie avant le premier point)
if let Some(pk) = file_name.split('.').next() { if let Some(pk) = file_name.split('.').next() {
if self.db.get(pk).is_err() { if self.db.get(pk, false).is_err() {
tokio::fs::remove_file(path).await?; tokio::fs::remove_file(path).await?;
} }
} }
@@ -601,6 +657,42 @@ impl<C: CacheConfig> Cache<C> {
.join(format!("{}.{}.{}", pk, qualifier, C::file_extension())) .join(format!("{}.{}.{}", pk, qualifier, C::file_extension()))
} }
/// Retourne tous les chemins de fichiers stockés pour une clé donnée,
/// quel que soit le qualifier.
///
/// Format: `{pk}.*.{extension}`
pub fn get_file_paths(&self, pk: &str) -> Result<Vec<PathBuf>> {
let mut paths = Vec::new();
let prefix = format!("{pk}.");
let expected_ext = C::file_extension();
for entry in std::fs::read_dir(&self.dir)? {
let entry = entry?;
let path = entry.path();
if !path.is_file() {
continue;
}
let file_name = match entry.file_name().into_string() {
Ok(name) => name,
Err(_) => continue, // nom de fichier non UTF-8 : on lignore
};
if !file_name.starts_with(&prefix) {
continue;
}
if !file_name.ends_with(expected_ext) {
continue;
}
paths.push(path);
}
Ok(paths)
}
/// Valide les données avant de les stocker /// Valide les données avant de les stocker
/// Par défaut, accepte toutes les données /// Par défaut, accepte toutes les données
pub fn validate_data(&self, data: &[u8]) -> Result<Vec<u8>> { pub fn validate_data(&self, data: &[u8]) -> Result<Vec<u8>> {

View File

@@ -51,11 +51,6 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
C::file_extension() C::file_extension()
} }
/// Retourne le nom de la table
fn table_name(&self) -> &'static str {
C::table_name()
}
/// Construit le chemin complet d'un fichier dans le cache /// Construit le chemin complet d'un fichier dans le cache
/// ///
/// Format: `{pk}.{qualificatif}.{extension}` /// Format: `{pk}.{qualificatif}.{extension}`
@@ -161,11 +156,11 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
/// assert_eq!(pk.len(), 32); // 16 bytes = 32 hex chars /// assert_eq!(pk.len(), 32); // 16 bytes = 32 hex chars
/// ``` /// ```
pub fn pk_from_content_header(header: &[u8]) -> String { pub fn pk_from_content_header(header: &[u8]) -> String {
use sha2::{Sha256, Digest}; use sha2::{Digest, Sha256};
let mut hasher = Sha256::new(); let mut hasher = Sha256::new();
hasher.update(header); hasher.update(header);
let result = hasher.finalize(); let result = hasher.finalize();
hex::encode(&result[..16]) // 16 octets = 32 caractères hex hex::encode(&result[..16]) // 16 octets = 32 caractères hex
} }
/// Génère une clé primaire à partir d'une URL (legacy) /// Génère une clé primaire à partir d'une URL (legacy)

View File

@@ -4,9 +4,12 @@
//! des éléments en cache, avec tracking des accès et des statistiques. //! des éléments en cache, avec tracking des accès et des statistiques.
use chrono::Utc; use chrono::Utc;
use rusqlite::{params, Connection}; use rusqlite::{params, Connection, Error, OptionalExtension};
use serde::Serialize; use serde::Serialize;
use serde_json::{Map, Number, Value};
use std::path::Path; use std::path::Path;
use std::str::FromStr;
use std::sync::Mutex; use std::sync::Mutex;
#[cfg(feature = "openapi")] #[cfg(feature = "openapi")]
@@ -21,7 +24,7 @@ pub struct CacheEntry {
pub pk: String, pub pk: String,
/// URL source de l'élément /// URL source de l'élément
#[cfg_attr(feature = "openapi", schema(example = "https://example.com/resource"))] #[cfg_attr(feature = "openapi", schema(example = "https://example.com/resource"))]
pub source_url: String, pub id: String,
/// Collection à laquelle appartient l'élément (optionnel) /// Collection à laquelle appartient l'élément (optionnel)
#[cfg_attr(feature = "openapi", schema(example = "album:123"))] #[cfg_attr(feature = "openapi", schema(example = "album:123"))]
pub collection: Option<String>, pub collection: Option<String>,
@@ -36,7 +39,7 @@ pub struct CacheEntry {
feature = "openapi", feature = "openapi",
schema(example = r#"{"title":"Track","artist":"Artist"}"#) schema(example = r#"{"title":"Track","artist":"Artist"}"#)
)] )]
pub metadata_json: Option<String>, pub metadata: Option<Value>,
} }
/// Base de données SQLite pour le cache /// Base de données SQLite pour le cache
@@ -48,7 +51,6 @@ pub struct CacheEntry {
#[derive(Debug)] #[derive(Debug)]
pub struct DB { pub struct DB {
conn: Mutex<Connection>, conn: Mutex<Connection>,
table_name: String,
} }
impl DB { impl DB {
@@ -67,42 +69,55 @@ impl DB {
/// ///
/// let db = DB::init(Path::new("cache.db"), "my_cache").unwrap(); /// let db = DB::init(Path::new("cache.db"), "my_cache").unwrap();
/// ``` /// ```
pub fn init(path: &Path, table_name: &str) -> Result<Self, rusqlite::Error> { pub fn init(path: &Path) -> Result<Self, rusqlite::Error> {
let conn = Connection::open(path)?; let conn = Connection::open(path)?;
let create_table_sql = format!( conn.execute(
"CREATE TABLE IF NOT EXISTS ASSET ( "CREATE TABLE IF NOT EXISTS asset (
pk TEXT PRIMARY KEY, pk TEXT PRIMARY KEY,
source_url TEXT,
collection TEXT, collection TEXT,
id TEXT,
hits INTEGER DEFAULT 0, hits INTEGER DEFAULT 0,
last_used TEXT, last_used TEXT
metadata_json TEXT
)", )",
table_name [],
); )?;
conn.execute(
conn.execute(&create_table_sql, [])?; "CREATE TABLE IF NOT EXISTS metadata (
pk TEXT,
key TEXT,
value_type TEXT NOT NULL CHECK (value_type IN ('string','number','boolean','null')),
value TEXT,
PRIMARY KEY (pk, key),
FOREIGN KEY (pk) REFERENCES asset (pk) ON DELETE CASCADE
)"
, [])?;
// Créer un index sur la collection pour les requêtes rapides // Créer un index sur la collection pour les requêtes rapides
let create_index_sql = format!( conn.execute(
"CREATE INDEX IF NOT EXISTS idx_{}_collection ON {} (collection)", "CREATE INDEX IF NOT EXISTS idx_asset_collection
table_name, table_name ON ASSET (collection)",
); [],
)?;
conn.execute(&create_index_sql, [])?;
// Créer un index composite pour optimiser la politique LRU (get_oldest) // Créer un index composite pour optimiser la politique LRU (get_oldest)
let create_lru_index_sql = format!( conn.execute(
"CREATE INDEX IF NOT EXISTS idx_{}_lru ON {} (last_used ASC, hits ASC)", "CREATE INDEX IF NOT EXISTS idx_asset_lru
table_name, table_name ON asset (last_used ASC, hits ASC)",
); [],
)?;
conn.execute(&create_lru_index_sql, [])?; // Crée un index composite pour rendre unique les ids si défini dans une collection
conn.execute(
"CREATE UNIQUE INDEX
IF NOT EXISTS asset_collection_id_unique
ON asset (collection, id)
WHERE id IS NOT NULL;",
[],
)?;
Ok(Self { Ok(Self {
conn: Mutex::new(conn), conn: Mutex::new(conn),
table_name: table_name.to_string(),
}) })
} }
@@ -113,8 +128,13 @@ impl DB {
/// * `pk` - Clé primaire de l'élément /// * `pk` - Clé primaire de l'élément
/// * `url` - URL source de l'élément /// * `url` - URL source de l'élément
/// * `collection` - Collection optionnelle à laquelle appartient l'élément /// * `collection` - Collection optionnelle à laquelle appartient l'élément
pub fn add(&self, pk: &str, url: &str, collection: Option<&str>) -> rusqlite::Result<()> { pub fn add(
self.add_with_metadata(pk, url, collection, None) &self,
pk: &str,
id: Option<&str>,
collection: Option<&str>,
) -> rusqlite::Result<()> {
self.add_with_metadata(pk, id, collection, None)
} }
/// Ajoute ou met à jour une entrée avec métadonnées JSON optionnelles /// Ajoute ou met à jour une entrée avec métadonnées JSON optionnelles
@@ -128,52 +148,317 @@ impl DB {
pub fn add_with_metadata( pub fn add_with_metadata(
&self, &self,
pk: &str, pk: &str,
url: &str, id: Option<&str>,
collection: Option<&str>, collection: Option<&str>,
metadata_json: Option<&str>, metadata: Option<&Value>,
) -> rusqlite::Result<()> { ) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!(
"INSERT INTO {} (pk, source_url, collection, hits, last_used, metadata_json)
VALUES (?1, ?2, ?3, 0, ?4, ?5)
ON CONFLICT(pk) DO UPDATE SET
source_url = excluded.source_url,
collection = excluded.collection,
last_used = excluded.last_used,
metadata_json = excluded.metadata_json",
self.table_name
);
conn.execute( conn.execute(
&sql, "INSERT INTO asset (pk, id, collection, hits, last_used)
params![pk, url, collection, Utc::now().to_rfc3339(), metadata_json], VALUES (?1, ?2, ?3, 0, ?4)
ON CONFLICT(pk) DO UPDATE SET
id = excluded.id,
collection = excluded.collection,
last_used = excluded.last_used",
params![pk, id, collection, Utc::now().to_rfc3339()],
)?;
if metadata.is_some() {
self.set_metadata(pk, metadata.unwrap())?
}
Ok(())
}
/// Remplace toutes les métadonnées associées à une entrée.
///
/// # Arguments
///
/// * `pk` - Clé primaire de l'élément ciblé.
/// * `metadata` - Objet JSON complet décrivant les nouvelles métadonnées.
///
/// # Errors
///
/// Retourne une erreur si `metadata` n'est pas un objet JSON ou si l'écriture
/// SQLite échoue.
pub fn set_metadata(&self, pk: &str, metadata: &Value) -> rusqlite::Result<()> {
let metadata_obj = metadata.as_object().ok_or_else(|| {
Error::InvalidParameterName("metadata must be a JSON object".to_owned())
})?;
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
tx.execute("DELETE FROM metadata WHERE pk = ?1", params![pk])?;
for (key, value) in metadata_obj.iter() {
let (value_type, value_text): (&str, Option<String>) = match value {
Value::Null => ("null", None),
Value::Bool(b) => ("boolean", Some(b.to_string())),
Value::Number(n) => ("number", Some(n.to_string())),
Value::String(s) => ("string", Some(s.clone())),
Value::Array(_) | Value::Object(_) => ("string", Some(value.to_string())),
};
tx.execute(
"INSERT INTO metadata (pk, key, value_type, value) VALUES (?1, ?2, ?3, ?4)",
params![pk, key, value_type, value_text.as_deref()],
)?;
}
tx.commit()
}
/// Insère ou met à jour une métadonnée individuelle.
///
/// # Arguments
///
/// * `pk` - Clé primaire de l'élément concerné.
/// * `key` - Nom de la métadonnée à enregistrer.
/// * `value` - Valeur JSON à stocker pour cette clé.
pub fn set_a_metadata(&self, pk: &str, key: &str, value: Value) -> rusqlite::Result<()> {
let (value_type, value_text): (&str, Option<String>) = match value {
Value::Null => ("null", None),
Value::Bool(b) => ("boolean", Some(b.to_string())),
Value::Number(n) => ("number", Some(n.to_string())),
Value::String(s) => ("string", Some(s)),
Value::Array(arr) => ("string", Some(Value::Array(arr).to_string())),
Value::Object(map) => ("string", Some(Value::Object(map).to_string())),
};
let conn = self.conn.lock().unwrap();
conn.execute(
"INSERT INTO metadata (pk, key, value_type, value)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(pk, key) DO UPDATE SET
value_type = excluded.value_type,
value = excluded.value",
params![pk, key, value_type, value_text.as_deref()],
)?; )?;
Ok(()) Ok(())
} }
/// Récupère une entrée de la base de données par sa clé /// Alias interne pour récupérer une métadonnée individuelle.
///
/// Préférer `get_metadata_value` pour les appels externes.
pub fn get_a_metadata(&self, pk: &str, key: &str) -> rusqlite::Result<Option<Value>> {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT value_type, value FROM metadata WHERE pk = ?1 AND key = ?2",
params![pk, key],
|row| {
let value_type: String = row.get(0)?;
let raw: Option<String> = row.get(1)?;
decode_metadata_value(key, &value_type, raw)
},
)
.optional()
}
/// Récupère toutes les métadonnées d'une entrée sous forme d'objet JSON.
///
/// Retourne `Ok(None)` si aucune métadonnée n'est présente.
pub fn get_metadata(&self, pk: &str) -> rusqlite::Result<Option<Value>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn.prepare("SELECT key, value_type, value FROM metadata WHERE pk = ?1")?;
let rows = stmt.query_map([pk], |row| {
let key: String = row.get(0)?;
let value_type: String = row.get(1)?;
let value: Option<String> = row.get(2)?;
Ok((key, value_type, value))
})?;
let mut metadata = Map::new();
let mut found = false;
for row in rows {
let (key, value_type, raw) = row?;
found = true;
let value = match value_type.as_str() {
"null" => Value::Null,
"boolean" => {
let raw = raw.as_deref().ok_or_else(|| {
Error::InvalidParameterName(format!(
"missing boolean metadata for key '{key}'"
))
})?;
let parsed = raw.parse::<bool>().map_err(|_| {
Error::InvalidParameterName(format!(
"invalid boolean metadata for key '{key}'"
))
})?;
Value::Bool(parsed)
}
"number" => {
let raw = raw.as_deref().ok_or_else(|| {
Error::InvalidParameterName(format!(
"missing number metadata for key '{key}'"
))
})?;
let number = Number::from_str(raw).map_err(|_| {
Error::InvalidParameterName(format!(
"invalid number metadata for key '{key}'"
))
})?;
Value::Number(number)
}
"string" => Value::String(raw.unwrap_or_default()),
other => {
return Err(Error::InvalidParameterName(format!(
"unknown metadata type '{other}' for key '{key}'"
)))
}
};
metadata.insert(key, value);
}
if found {
Ok(Some(Value::Object(metadata)))
} else {
Ok(None)
}
}
/// Enregistre l'URL d'origine liée à un élément du cache.
pub fn set_origin_url(&self, pk: &str, origin_url: &str) -> rusqlite::Result<()> {
self.set_a_metadata(pk, "origin_url", Value::String(origin_url.to_owned()))
}
/// Récupère l'URL d'origine précédemment stockée pour un élément.
///
/// Retourne `Ok(None)` si aucune URL n'est définie.
pub fn get_origin_url(&self, pk: &str) -> rusqlite::Result<Option<String>> {
match self.get_metadata_value(pk, "origin_url")? {
Some(Value::String(url)) => Ok(Some(url)),
Some(Value::Null) | None => Ok(None),
Some(other) => Err(Error::InvalidParameterName(format!(
"metadata 'origin_url' must be a string, got {other}"
))),
}
}
/// Récupère uniquement les métadonnées JSON d'une entrée
/// ///
/// # Arguments /// # Arguments
/// ///
/// * `pk` - Clé primaire de l'élément à récupérer /// * `pk` - Clé primaire de l'élément
pub fn get(&self, pk: &str) -> rusqlite::Result<CacheEntry> { ///
let conn = self.conn.lock().unwrap(); /// # Returns
let sql = format!( ///
"SELECT pk, source_url, collection, hits, last_used, metadata_json FROM {} WHERE pk = ?1", /// Les métadonnées JSON si présentes, None sinon
self.table_name pub fn get_metadata_json(&self, pk: &str) -> rusqlite::Result<Option<String>> {
); Ok(self.get_metadata(pk)?.map(|value| value.to_string()))
}
conn.query_row(&sql, [pk], |row| { /// Récupère une métadonnée individuelle, si elle existe.
Ok(CacheEntry { pub fn get_metadata_value(&self, pk: &str, key: &str) -> rusqlite::Result<Option<Value>> {
pk: row.get(0)?, let conn = self.conn.lock().unwrap();
source_url: row.get(1)?,
collection: row.get(2)?, conn.query_row(
hits: row.get(3)?, "SELECT value_type, value FROM metadata WHERE pk = ?1 AND key = ?2",
last_used: row.get(4)?, params![pk, key],
metadata_json: row.get(5)?, |row| {
}) let value_type: String = row.get(0)?;
}) let raw: Option<String> = row.get(1)?;
decode_metadata_value(key, &value_type, raw)
},
)
.optional()
}
/// Récupère une entrée de la base de données par sa clé
///
/// # Arguments
/// * `pk` - Clé primaire de l'élément à récupérer.
/// * `with_metadata` - Charge les métadonnées associées si `true`.
pub fn get(&self, pk: &str, with_metadata: bool) -> rusqlite::Result<CacheEntry> {
let mut entry = {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT pk, id, collection, hits, last_used \
FROM asset \
WHERE pk = ?1",
[pk],
|row| {
Ok(CacheEntry {
pk: row.get(0)?,
id: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata: None,
})
},
)?
};
if with_metadata {
entry.metadata = self.get_metadata(&entry.pk)?;
}
Ok(entry)
}
/// Récupère une entrée en utilisant la paire `(collection, id)`.
///
/// # Arguments
///
/// * `collection` - Collection dans laquelle chercher.
/// * `id` - Identifiant logique de l'élément.
/// * `with_metadata` - Charge les métadonnées associées si `true`.
pub fn get_from_id(
&self,
collection: &str,
id: &str,
with_metadata: bool,
) -> rusqlite::Result<CacheEntry> {
let mut entry = {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT pk, id, collection, hits, last_used \
FROM asset \
WHERE collection = ?1 AND id = ?2",
params![collection, id],
|row| {
Ok(CacheEntry {
pk: row.get(0)?,
id: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata: None,
})
},
)?
};
if with_metadata {
entry.metadata = self.get_metadata(&entry.pk)?;
}
Ok(entry)
}
/// Définit ou remplace l'identifiant logique (`id`) d'une entrée.
///
/// Retourne `QueryReturnedNoRows` si la clé primaire est inconnue.
pub fn set_id(&self, pk: &str, id: &str) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap();
let updated = conn.execute("UPDATE asset SET id = ?2 WHERE pk = ?1", params![pk, id])?;
if updated == 0 {
return Err(Error::QueryReturnedNoRows);
}
Ok(())
} }
/// Met à jour le compteur d'accès et la date du dernier accès /// Met à jour le compteur d'accès et la date du dernier accès
@@ -183,47 +468,57 @@ impl DB {
/// * `pk` - Clé primaire de l'élément /// * `pk` - Clé primaire de l'élément
pub fn update_hit(&self, pk: &str) -> rusqlite::Result<()> { pub fn update_hit(&self, pk: &str) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!(
"UPDATE {} SET hits = hits + 1, last_used = ?1 WHERE pk = ?2",
self.table_name
);
conn.execute(&sql, params![Utc::now().to_rfc3339(), pk])?; conn.execute(
&"UPDATE asset
SET hits = hits + 1, last_used = ?1
WHERE pk = ?2",
params![Utc::now().to_rfc3339(), pk],
)?;
Ok(()) Ok(())
} }
/// Purge toutes les entrées de la base de données /// Purge toutes les entrées de la base de données.
pub fn purge(&self) -> rusqlite::Result<()> { pub fn purge(&self) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!("DELETE FROM {}", self.table_name); conn.execute("DELETE FROM asset", [])?;
conn.execute(&sql, [])?;
Ok(()) Ok(())
} }
/// Récupère toutes les entrées, triées par nombre d'accès décroissant /// Récupère toutes les entrées, triées par nombre d'accès décroissant.
pub fn get_all(&self) -> rusqlite::Result<Vec<CacheEntry>> { ///
/// # Arguments
///
/// * `include_metadata` - Ajoute les métadonnées à chaque entrée si `true`.
pub fn get_all(&self, include_metadata: bool) -> rusqlite::Result<Vec<CacheEntry>> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!(
"SELECT pk, source_url, collection, hits, last_used, metadata_json FROM {} ORDER BY hits DESC",
self.table_name
);
let mut stmt = conn.prepare(&sql)?; let mut stmt = conn.prepare(
"SELECT pk, id, collection, hits, last_used
FROM asset
ORDER BY hits DESC",
)?;
let entries = stmt let mut entries = stmt
.query_map([], |row| { .query_map([], |row| {
Ok(CacheEntry { Ok(CacheEntry {
pk: row.get(0)?, pk: row.get(0)?,
source_url: row.get(1)?, id: row.get(1)?,
collection: row.get(2)?, collection: row.get(2)?,
hits: row.get(3)?, hits: row.get(3)?,
last_used: row.get(4)?, last_used: row.get(4)?,
metadata_json: row.get(5)?, metadata: None,
}) })
})? })?
.collect::<rusqlite::Result<Vec<_>>>()?; .collect::<rusqlite::Result<Vec<_>>>()?;
if include_metadata {
for entry in entries.iter_mut() {
entry.metadata = self.get_metadata(&entry.pk)?;
}
}
Ok(entries) Ok(entries)
} }
@@ -231,53 +526,57 @@ impl DB {
/// ///
/// # Arguments /// # Arguments
/// ///
/// * `collection` - Identifiant de la collection /// * `collection` - Identifiant de la collection.
pub fn get_by_collection(&self, collection: &str) -> rusqlite::Result<Vec<CacheEntry>> { /// * `include_metadata` - Ajoute les métadonnées à chaque entrée si `true`.
pub fn get_by_collection(
&self,
collection: &str,
include_metadata: bool,
) -> rusqlite::Result<Vec<CacheEntry>> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!(
"SELECT pk, source_url, collection, hits, last_used, metadata_json FROM {} WHERE collection = ?1 ORDER BY hits DESC",
self.table_name
);
let mut stmt = conn.prepare(&sql)?; let mut stmt = conn.prepare(
"SELECT pk, id, collection, hits, last_used
FROM asset
WHERE collection = ?1 ORDER BY hits DESC",
)?;
let entries = stmt let mut entries = stmt
.query_map([collection], |row| { .query_map([collection], |row| {
Ok(CacheEntry { Ok(CacheEntry {
pk: row.get(0)?, pk: row.get(0)?,
source_url: row.get(1)?, id: row.get(1)?,
collection: row.get(2)?, collection: row.get(2)?,
hits: row.get(3)?, hits: row.get(3)?,
last_used: row.get(4)?, last_used: row.get(4)?,
metadata_json: row.get(5)?, metadata: None,
}) })
})? })?
.collect::<rusqlite::Result<Vec<_>>>()?; .collect::<rusqlite::Result<Vec<_>>>()?;
if include_metadata {
for entry in entries.iter_mut() {
entry.metadata = self.get_metadata(&entry.pk)?;
}
}
Ok(entries) Ok(entries)
} }
/// Supprime toutes les entrées d'une collection /// Supprime toutes les entrées d'une collection.
/// ///
/// # Arguments /// Les métadonnées associées sont supprimées automatiquement grâce à la
/// /// contrainte `ON DELETE CASCADE`.
/// * `collection` - Identifiant de la collection à supprimer
pub fn delete_collection(&self, collection: &str) -> rusqlite::Result<()> { pub fn delete_collection(&self, collection: &str) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!("DELETE FROM {} WHERE collection = ?1", self.table_name); conn.execute("DELETE FROM asset WHERE collection = ?1", [collection])?;
conn.execute(&sql, [collection])?;
Ok(()) Ok(())
} }
/// Supprime une entrée de la base de données /// Supprime une entrée de la base de données ainsi que ses métadonnées.
///
/// # Arguments
///
/// * `pk` - Clé primaire de l'élément à supprimer
pub fn delete(&self, pk: &str) -> rusqlite::Result<()> { pub fn delete(&self, pk: &str) -> rusqlite::Result<()> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!("DELETE FROM {} WHERE pk = ?1", self.table_name); conn.execute("DELETE FROM asset WHERE pk = ?1", [pk])?;
conn.execute(&sql, [pk])?;
Ok(()) Ok(())
} }
@@ -288,8 +587,7 @@ impl DB {
/// Le nombre total d'entrées /// Le nombre total d'entrées
pub fn count(&self) -> rusqlite::Result<usize> { pub fn count(&self) -> rusqlite::Result<usize> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!("SELECT COUNT(*) FROM {}", self.table_name); let count: i64 = conn.query_row("SELECT COUNT(*) FROM asset", [], |row| row.get(0))?;
let count: i64 = conn.query_row(&sql, [], |row| row.get(0))?;
Ok(count as usize) Ok(count as usize)
} }
@@ -307,65 +605,67 @@ impl DB {
/// Liste des entrées les plus anciennes, triées par last_used ASC /// Liste des entrées les plus anciennes, triées par last_used ASC
pub fn get_oldest(&self, limit: usize) -> rusqlite::Result<Vec<CacheEntry>> { pub fn get_oldest(&self, limit: usize) -> rusqlite::Result<Vec<CacheEntry>> {
let conn = self.conn.lock().unwrap(); let conn = self.conn.lock().unwrap();
let sql = format!(
let mut stmt = conn.prepare(
"SELECT pk, source_url, collection, hits, last_used, metadata_json "SELECT pk, source_url, collection, hits, last_used, metadata_json
FROM {} FROM asset
ORDER BY last_used ASC, hits ASC ORDER BY last_used ASC, hits ASC
LIMIT ?1", LIMIT ?1",
self.table_name )?;
);
let mut stmt = conn.prepare(&sql)?;
let entries = stmt let entries = stmt
.query_map([limit], |row| { .query_map([limit], |row| {
Ok(CacheEntry { Ok(CacheEntry {
pk: row.get(0)?, pk: row.get(0)?,
source_url: row.get(1)?, id: row.get(1)?,
collection: row.get(2)?, collection: row.get(2)?,
hits: row.get(3)?, hits: row.get(3)?,
last_used: row.get(4)?, last_used: row.get(4)?,
metadata_json: row.get(5)?, metadata: None,
}) })
})? })?
.collect::<rusqlite::Result<Vec<_>>>()?; .collect::<rusqlite::Result<Vec<_>>>()?;
Ok(entries) Ok(entries)
} }
}
/// Récupère uniquement les métadonnées JSON d'une entrée /// Convertit une ligne de la table `metadata` en valeur JSON.
/// fn decode_metadata_value(
/// # Arguments key: &str,
/// value_type: &str,
/// * `pk` - Clé primaire de l'élément raw: Option<String>,
/// ) -> rusqlite::Result<Value> {
/// # Returns match value_type {
/// "null" => Ok(Value::Null),
/// Les métadonnées JSON si présentes, None sinon "boolean" => {
pub fn get_metadata_json(&self, pk: &str) -> rusqlite::Result<Option<String>> { let raw = raw.as_deref().ok_or_else(|| {
let conn = self.conn.lock().unwrap(); Error::InvalidParameterName(format!("missing boolean metadata for '{key}'"))
let sql = format!( })?;
"SELECT metadata_json FROM {} WHERE pk = ?1", raw.parse::<bool>().map(Value::Bool).map_err(|_| {
self.table_name Error::InvalidParameterName(format!("invalid boolean metadata for '{key}'"))
); })
}
conn.query_row(&sql, [pk], |row| row.get(0)) "number" => {
} let raw = raw.as_deref().ok_or_else(|| {
Error::InvalidParameterName(format!("missing number metadata for '{key}'"))
/// Met à jour uniquement les métadonnées JSON d'une entrée existante })?;
/// Number::from_str(raw).map(Value::Number).map_err(|_| {
/// # Arguments Error::InvalidParameterName(format!("invalid number metadata for '{key}'"))
/// })
/// * `pk` - Clé primaire de l'élément }
/// * `metadata_json` - Métadonnées JSON à stocker "string" => {
pub fn update_metadata(&self, pk: &str, metadata_json: &str) -> rusqlite::Result<()> { let raw = raw.unwrap_or_default();
let conn = self.conn.lock().unwrap(); let trimmed = raw.trim_start();
let sql = format!( if trimmed.starts_with('{') || trimmed.starts_with('[') {
"UPDATE {} SET metadata_json = ?1 WHERE pk = ?2", if let Ok(json) = serde_json::from_str::<Value>(&raw) {
self.table_name return Ok(json);
); }
}
conn.execute(&sql, params![metadata_json, pk])?; Ok(Value::String(raw))
Ok(()) }
other => Err(Error::InvalidParameterName(format!(
"unknown metadata type '{other}' for key '{key}'"
))),
} }
} }

View File

@@ -501,7 +501,8 @@ pub async fn peek_header(url: &str, max_bytes: usize) -> Result<Vec<u8>, String>
// Si le serveur ne supporte pas Range (status 200 au lieu de 206), // Si le serveur ne supporte pas Range (status 200 au lieu de 206),
// on lit quand même mais on limite la lecture // on lit quand même mais on limite la lecture
if !response.status().is_success() && response.status() != reqwest::StatusCode::PARTIAL_CONTENT { if !response.status().is_success() && response.status() != reqwest::StatusCode::PARTIAL_CONTENT
{
return Err(format!("HTTP error: {}", response.status())); return Err(format!("HTTP error: {}", response.status()));
} }

View File

@@ -143,8 +143,8 @@ pub use cache::{Cache, CacheConfig};
pub use cache_trait::{pk_from_content_header, pk_from_url, FileCache}; pub use cache_trait::{pk_from_content_header, pk_from_url, FileCache};
pub use db::{CacheEntry, DB}; pub use db::{CacheEntry, DB};
pub use download::{ pub use download::{
download, download_with_transformer, ingest_with_transformer, peek_header, download, download_with_transformer, ingest_with_transformer, peek_header, peek_reader_header,
peek_reader_header, Download, StreamTransformer, Download, StreamTransformer,
}; };
#[cfg(feature = "pmoserver")] #[cfg(feature = "pmoserver")]

View File

@@ -70,9 +70,7 @@ where
(status = 200, description = "Configuration complète", body = serde_json::Value) (status = 200, description = "Configuration complète", body = serde_json::Value)
) )
)] )]
async fn get_full_config( async fn get_full_config(State(config): State<Arc<Config>>) -> Result<Json<JsonValue>, ApiError> {
State(config): State<Arc<Config>>,
) -> Result<Json<JsonValue>, ApiError> {
let value = config.get_value(&[])?; let value = config.get_value(&[])?;
let json_value = yaml_to_json(&value)?; let json_value = yaml_to_json(&value)?;
Ok(Json(json_value)) Ok(Json(json_value))

View File

@@ -520,16 +520,27 @@ impl Config {
Ok(Value::String(s)) => match s.parse::<u16>() { Ok(Value::String(s)) => match s.parse::<u16>() {
Ok(port) => port, Ok(port) => port,
Err(_) => { Err(_) => {
tracing::warn!("Invalid HTTP port '{}', using default {}", s, DEFAULT_HTTP_PORT); tracing::warn!(
"Invalid HTTP port '{}', using default {}",
s,
DEFAULT_HTTP_PORT
);
DEFAULT_HTTP_PORT DEFAULT_HTTP_PORT
} }
}, },
Ok(_) => { Ok(_) => {
tracing::warn!("HTTP port not a number or string, using default {}", DEFAULT_HTTP_PORT); tracing::warn!(
"HTTP port not a number or string, using default {}",
DEFAULT_HTTP_PORT
);
DEFAULT_HTTP_PORT DEFAULT_HTTP_PORT
} }
Err(err) => { Err(err) => {
tracing::warn!("Failed to get HTTP port: {}, using default {}", err, DEFAULT_HTTP_PORT); tracing::warn!(
"Failed to get HTTP port: {}, using default {}",
err,
DEFAULT_HTTP_PORT
);
DEFAULT_HTTP_PORT DEFAULT_HTTP_PORT
} }
} }
@@ -586,7 +597,6 @@ impl Config {
self.set_value(&["devices", devtype, name, "udn"], Value::String(udn)) self.set_value(&["devices", devtype, name, "udn"], Value::String(udn))
} }
impl_string_config!( impl_string_config!(
/// Gets the Qobuz username from configuration /// Gets the Qobuz username from configuration
get_qobuz_username, get_qobuz_username,
@@ -618,9 +628,19 @@ impl Config {
Ok((username, password)) Ok((username, password))
} }
impl_usize_config!(get_log_cache_size, set_log_cache_size, &["host", "logger", "buffer_capacity"], DEFAULT_LOG_BUFFER_CAPACITY); impl_usize_config!(
get_log_cache_size,
set_log_cache_size,
&["host", "logger", "buffer_capacity"],
DEFAULT_LOG_BUFFER_CAPACITY
);
impl_bool_config!(get_log_enable_console, set_log_enable_console, &["host", "logger", "enable_console"], DEFAULT_LOG_ENABLE_CONSOLE); impl_bool_config!(
get_log_enable_console,
set_log_enable_console,
&["host", "logger", "enable_console"],
DEFAULT_LOG_ENABLE_CONSOLE
);
/// Récupère le niveau de log minimum depuis la configuration /// Récupère le niveau de log minimum depuis la configuration
pub fn get_log_min_level(&self) -> Result<String> { pub fn get_log_min_level(&self) -> Result<String> {

View File

@@ -15,10 +15,6 @@ impl CacheConfig for CoversConfig {
"webp" "webp"
} }
fn table_name() -> &'static str {
"covers"
}
fn cache_type() -> &'static str { fn cache_type() -> &'static str {
"image" "image"
} }

View File

@@ -4,8 +4,8 @@
//! des méthodes de gestion du cache de couvertures à pmoconfig::Config. //! des méthodes de gestion du cache de couvertures à pmoconfig::Config.
use anyhow::Result; use anyhow::Result;
use pmoconfig::Config;
use pmocache::CacheConfigExt; use pmocache::CacheConfigExt;
use pmoconfig::Config;
use std::sync::Arc; use std::sync::Arc;
const DEFAULT_COVER_CACHE_DIR: &str = "cache_covers"; const DEFAULT_COVER_CACHE_DIR: &str = "cache_covers";

View File

@@ -1,12 +1,11 @@
use std::sync::Arc; use std::sync::Arc;
use bevy_reflect::Reflect; use bevy_reflect::Reflect;
use htmlescape::decode_html;
use once_cell::sync::Lazy; use once_cell::sync::Lazy;
use pmodidl::{DIDLLite, MediaMetadataParser}; use pmodidl::{DIDLLite, MediaMetadataParser};
use pmoupnp::state_variables::{StateVariable, StateVariableError}; use pmoupnp::state_variables::{StateVariable, StateVariableError};
use pmoupnp::variable_types::StateVarType; use pmoupnp::variable_types::StateVarType;
use htmlescape::decode_html;
fn avtransporturimetadataparser(value: &str) -> Result<Box<dyn Reflect>, StateVariableError> { fn avtransporturimetadataparser(value: &str) -> Result<Box<dyn Reflect>, StateVariableError> {
// Nettoyage de base // Nettoyage de base
@@ -32,7 +31,8 @@ fn avtransporturimetadataparser(value: &str) -> Result<Box<dyn Reflect>, StateVa
} }
fn avtransporturimetadatamarshal(value: &dyn Reflect) -> Result<String, StateVariableError> { fn avtransporturimetadatamarshal(value: &dyn Reflect) -> Result<String, StateVariableError> {
let didl = value.downcast_ref::<DIDLLite>() let didl = value
.downcast_ref::<DIDLLite>()
.ok_or_else(|| StateVariableError::ConversionError("DIDLLite".into()))?; .ok_or_else(|| StateVariableError::ConversionError("DIDLLite".into()))?;
let xml = quick_xml::se::to_string(didl) let xml = quick_xml::se::to_string(didl)
.map_err(|e| StateVariableError::ConversionError(format!("serialize error: {}", e)))?; .map_err(|e| StateVariableError::ConversionError(format!("serialize error: {}", e)))?;

View File

@@ -25,9 +25,7 @@ async fn main() -> Result<()> {
eprintln!("======================================\n"); eprintln!("======================================\n");
// Create client // Create client
let mut client = RadioParadiseClient::builder() let mut client = RadioParadiseClient::builder().build().await?;
.build()
.await?;
eprintln!("Client configured for FLAC streaming\n"); eprintln!("Client configured for FLAC streaming\n");

View File

@@ -65,10 +65,12 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
let decode_task = tokio::task::spawn_blocking(move || -> anyhow::Result<Vec<(u64, usize)>> { let decode_task = tokio::task::spawn_blocking(move || -> anyhow::Result<Vec<(u64, usize)>> {
let mut decoder = StreamingPCMDecoder::new(http_stream)?; let mut decoder = StreamingPCMDecoder::new(http_stream)?;
println!(" 🎼 Stream info: {}Hz, {} channels, {} bits", println!(
decoder.sample_rate(), " 🎼 Stream info: {}Hz, {} channels, {} bits",
decoder.channels(), decoder.sample_rate(),
decoder.bits_per_sample()); decoder.channels(),
decoder.bits_per_sample()
);
let mut chunk_times = Vec::new(); let mut chunk_times = Vec::new();
let mut chunk_count = 0; let mut chunk_count = 0;
@@ -78,17 +80,21 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
chunk_times.push((chunk.position_ms, chunk.samples.len())); chunk_times.push((chunk.position_ms, chunk.samples.len()));
if chunk_count % 50 == 0 { if chunk_count % 50 == 0 {
println!(" 📦 Chunk {} at {}ms ({} samples)", println!(
chunk_count, " 📦 Chunk {} at {}ms ({} samples)",
chunk.position_ms, chunk_count,
chunk.samples.len()); chunk.position_ms,
chunk.samples.len()
);
} }
} }
Ok(chunk_times) Ok(chunk_times)
}); });
let chunk_times = decode_task.await.map_err(|e| anyhow::anyhow!("Join error: {}", e))??; let chunk_times = decode_task
.await
.map_err(|e| anyhow::anyhow!("Join error: {}", e))??;
let total_time = start_time.elapsed(); let total_time = start_time.elapsed();
println!("\n✅ Streaming Complete!"); println!("\n✅ Streaming Complete!");
@@ -101,7 +107,11 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
} }
if let Some((last_pos, _)) = chunk_times.last() { if let Some((last_pos, _)) = chunk_times.last() {
println!(" Last chunk at: {}ms (~{:.1}s)", last_pos, last_pos / 1000); println!(
" Last chunk at: {}ms (~{:.1}s)",
last_pos,
last_pos / 1000
);
} }
println!("\n💡 Analysis:"); println!("\n💡 Analysis:");

View File

@@ -24,7 +24,6 @@ pub const DEFAULT_BLOCK_TIMEOUT_SECS: u64 = 180;
/// Default User-Agent /// Default User-Agent
pub const DEFAULT_USER_AGENT: &str = "pmoparadise/0.1.0"; pub const DEFAULT_USER_AGENT: &str = "pmoparadise/0.1.0";
/// Radio Paradise HTTP client /// Radio Paradise HTTP client
/// ///
/// This client provides access to Radio Paradise's streaming API, /// This client provides access to Radio Paradise's streaming API,
@@ -102,7 +101,6 @@ impl RadioParadiseClient {
cloned cloned
} }
/// Get a block by event ID /// Get a block by event ID
/// ///
/// If `event` is None, returns the current block. /// If `event` is None, returns the current block.
@@ -283,7 +281,6 @@ impl ClientBuilder {
self self
} }
/// Set the channel (0 = main mix, 1 = mellow, 2 = rock, 3 = world/etc) /// Set the channel (0 = main mix, 1 = mellow, 2 = rock, 3 = world/etc)
pub fn channel(mut self, channel: u8) -> Self { pub fn channel(mut self, channel: u8) -> Self {
self.channel = channel; self.channel = channel;
@@ -360,5 +357,4 @@ mod tests {
assert_eq!(builder.api_base, DEFAULT_API_BASE); assert_eq!(builder.api_base, DEFAULT_API_BASE);
assert_eq!(builder.channel, 0); assert_eq!(builder.channel, 0);
} }
} }

View File

@@ -387,8 +387,12 @@ impl WorkerState {
let mut decoder = StreamingPCMDecoder::new(http_stream) let mut decoder = StreamingPCMDecoder::new(http_stream)
.context("Failed to create streaming decoder")?; .context("Failed to create streaming decoder")?;
info!("Streaming decoder initialized: {}Hz, {} channels, {} bits", info!(
decoder.sample_rate(), decoder.channels(), decoder.bits_per_sample()); "Streaming decoder initialized: {}Hz, {} channels, {} bits",
decoder.sample_rate(),
decoder.channels(),
decoder.bits_per_sample()
);
// Decode chunks and send them // Decode chunks and send them
while let Some(chunk) = decoder.decode_chunk()? { while let Some(chunk) = decoder.decode_chunk()? {
@@ -649,14 +653,10 @@ impl WorkerState {
let duration_ms = song.duration; let duration_ms = song.duration;
// Encode PCM to FLAC // Encode PCM to FLAC
let flac_bytes = encode_samples_to_flac( let flac_bytes =
track_samples, encode_samples_to_flac(track_samples, channels, sample_rate, bits_per_sample)
channels, .await
sample_rate, .context("Failed to encode song to FLAC")?;
bits_per_sample,
)
.await
.context("Failed to encode song to FLAC")?;
let track_id = self.compute_track_id(block, song_index); let track_id = self.compute_track_id(block, song_index);
let placeholder_uri = format!("{}#{}", block.url, song_index); let placeholder_uri = format!("{}#{}", block.url, song_index);
@@ -715,7 +715,10 @@ impl WorkerState {
} }
} }
} else { } else {
warn!(channel = self.descriptor.slug, "Unable to resolve cover URL for {}", cover_path); warn!(
channel = self.descriptor.slug,
"Unable to resolve cover URL for {}", cover_path
);
} }
} }
Ok(None) Ok(None)
@@ -724,7 +727,10 @@ impl WorkerState {
fn compute_track_id(&self, block: &Block, song_index: usize) -> String { fn compute_track_id(&self, block: &Block, song_index: usize) -> String {
// Use deterministic ID based on block event and song index // Use deterministic ID based on block event and song index
// This allows checking if a song is cached before downloading the block // This allows checking if a song is cached before downloading the block
format!("rp:{}:event_{}_song_{}", self.descriptor.id, block.event, song_index) format!(
"rp:{}:event_{}_song_{}",
self.descriptor.id, block.event, song_index
)
} }
async fn maybe_schedule_poll(&mut self) { async fn maybe_schedule_poll(&mut self) {
@@ -811,7 +817,12 @@ impl WorkerState {
}; };
// Check if file exists // Check if file exists
if self.cache_manager.audio_file_path(&audio_pk).await.is_none() { if self
.cache_manager
.audio_file_path(&audio_pk)
.await
.is_none()
{
debug!( debug!(
channel = self.descriptor.slug, channel = self.descriptor.slug,
event = block.event, event = block.event,
@@ -845,10 +856,15 @@ impl WorkerState {
let track_id = self.compute_track_id(block, *song_index); let track_id = self.compute_track_id(block, *song_index);
// Get metadata (we already checked it exists in check_all_songs_cached) // Get metadata (we already checked it exists in check_all_songs_cached)
let metadata = self.cache_manager.get_metadata(&track_id).await let metadata = self
.cache_manager
.get_metadata(&track_id)
.await
.ok_or_else(|| anyhow!("Metadata disappeared for track_id: {}", track_id))?; .ok_or_else(|| anyhow!("Metadata disappeared for track_id: {}", track_id))?;
let audio_pk = metadata.cached_audio_pk.clone() let audio_pk = metadata
.cached_audio_pk
.clone()
.ok_or_else(|| anyhow!("Audio PK disappeared for track_id: {}", track_id))?; .ok_or_else(|| anyhow!("Audio PK disappeared for track_id: {}", track_id))?;
// Get cover PK if available // Get cover PK if available
@@ -874,10 +890,15 @@ impl WorkerState {
cached_cover_pk: cover_pk, cached_cover_pk: cover_pk,
..metadata.clone() ..metadata.clone()
}; };
self.cache_manager.update_metadata(track_id.clone(), updated_metadata).await; self.cache_manager
.update_metadata(track_id.clone(), updated_metadata)
.await;
} }
let file_path = self.cache_manager.audio_file_path(&audio_pk).await let file_path = self
.cache_manager
.audio_file_path(&audio_pk)
.await
.ok_or_else(|| anyhow!("File disappeared for audio_pk: {}", audio_pk))?; .ok_or_else(|| anyhow!("File disappeared for audio_pk: {}", audio_pk))?;
let duration_ms = song.duration; let duration_ms = song.duration;
@@ -1120,4 +1141,3 @@ async fn encode_samples_to_flac(
}) })
.await? .await?
} }

View File

@@ -36,7 +36,6 @@ struct ParadiseQuery {
channel: Option<u8>, channel: Option<u8>,
} }
#[derive(Debug, Default, Deserialize, IntoParams)] #[derive(Debug, Default, Deserialize, IntoParams)]
#[serde(default)] #[serde(default)]
#[into_params(parameter_in = Query)] #[into_params(parameter_in = Query)]
@@ -801,7 +800,10 @@ pub fn create_api_router(state: RadioParadiseState) -> Router {
.route("/channels/{channel_id}/status", get(get_channel_status)) .route("/channels/{channel_id}/status", get(get_channel_status))
.route("/channels/{channel_id}/playlist", get(get_channel_playlist)) .route("/channels/{channel_id}/playlist", get(get_channel_playlist))
.route("/channels/{channel_id}/history", get(get_channel_history)) .route("/channels/{channel_id}/history", get(get_channel_history))
.route("/channels/{channel_id}/stream/{connection_id}", get(stream_channel_by_connection)) .route(
"/channels/{channel_id}/stream/{connection_id}",
get(stream_channel_by_connection),
)
.with_state(state) .with_state(state)
} }

View File

@@ -164,7 +164,6 @@ impl RadioParadiseClient {
Ok(Bytes::from(data)) Ok(Bytes::from(data))
} }
} }
#[cfg(test)] #[cfg(test)]

View File

@@ -125,7 +125,9 @@ impl StreamingPCMDecoder<ChannelReader> {
} }
pub fn decode_chunk(&mut self) -> anyhow::Result<Option<PCMChunk>> { pub fn decode_chunk(&mut self) -> anyhow::Result<Option<PCMChunk>> {
if self.done { return Ok(None); } if self.done {
return Ok(None);
}
// Crée le FrameReader à la volée (emprunt de self.reader) // Crée le FrameReader à la volée (emprunt de self.reader)
let mut frames = self.reader.blocks(); let mut frames = self.reader.blocks();
@@ -133,7 +135,10 @@ impl StreamingPCMDecoder<ChannelReader> {
// API claxon 0.6.x : il FAUT fournir un Vec<i32> par valeur // API claxon 0.6.x : il FAUT fournir un Vec<i32> par valeur
let buf: Vec<i32> = Vec::new(); let buf: Vec<i32> = Vec::new();
let frame = match frames.read_next_or_eof(buf) { let frame = match frames.read_next_or_eof(buf) {
Ok(None) => { self.done = true; return Ok(None); } Ok(None) => {
self.done = true;
return Ok(None);
}
Ok(Some(f)) => f, Ok(Some(f)) => f,
Err(e) => return Err(anyhow::anyhow!("FLAC decode error: {}", e)), Err(e) => return Err(anyhow::anyhow!("FLAC decode error: {}", e)),
}; };

View File

@@ -127,9 +127,6 @@ async fn test_now_playing() {
} }
} }
#[tokio::test] #[tokio::test]
async fn test_prefetch_next() { async fn test_prefetch_next() {
let mock_server = MockServer::start().await; let mock_server = MockServer::start().await;

View File

@@ -5,7 +5,7 @@
use crate::Server; use crate::Server;
use anyhow::Result; use anyhow::Result;
use pmoconfig::{api, get_config, ApiDoc}; use pmoconfig::{ApiDoc, api, get_config};
use utoipa::OpenApi; use utoipa::OpenApi;
/// Trait d'extension pour ajouter l'API de configuration à pmoserver /// Trait d'extension pour ajouter l'API de configuration à pmoserver

View File

@@ -367,23 +367,17 @@ impl UpnpServerExt for Server {
} }
async fn init_caches(&mut self) -> Result<(Arc<CoverCache>, Arc<AudioCache>), anyhow::Error> { async fn init_caches(&mut self) -> Result<(Arc<CoverCache>, Arc<AudioCache>), anyhow::Error> {
use pmocovers::CoverCacheConfigExt;
use pmoaudiocache::AudioCacheConfigExt; use pmoaudiocache::AudioCacheConfigExt;
use pmocovers::CoverCacheConfigExt;
let config = pmoconfig::get_config(); let config = pmoconfig::get_config();
let cover_cache = self let cover_cache = self
.init_cover_cache( .init_cover_cache(&config.get_covers_dir()?, config.get_covers_size()?)
&config.get_covers_dir()?,
config.get_covers_size()?,
)
.await?; .await?;
let audio_cache = self let audio_cache = self
.init_audio_cache( .init_audio_cache(&config.get_audiocache_dir()?, config.get_audiocache_size()?)
&config.get_audiocache_dir()?,
config.get_audiocache_size()?,
)
.await?; .await?;
Ok((cover_cache, audio_cache)) Ok((cover_cache, audio_cache))