amélioration de la webapp

This commit is contained in:
2025-10-19 13:42:29 +02:00
parent 90a7b8b6fa
commit 776c535862
179 changed files with 2426 additions and 1685 deletions

Binary file not shown.

View File

@@ -1,8 +1,8 @@
use pmoapp::{WebAppExt, Webapp};
use pmomediarenderer::MEDIA_RENDERER;
use pmomediaserver::{MEDIA_SERVER, sources::SourcesExt};
use pmosource::MusicSourceExt;
use pmoserver::Server;
use pmosource::MusicSourceExt;
use pmoupnp::{UpnpServerExt, upnp_api::UpnpApiExt};
use tracing::info;

View File

@@ -297,7 +297,7 @@ pub trait WebAppExt {
/// # Type Parameter
///
/// * `W` - Type RustEmbed contenant les fichiers de la webapp
async fn add_webapp<W>(&mut self, path: &str)
async fn add_webapp<W>(&mut self, path: &str)
where
W: RustEmbed + Clone + Send + Sync + 'static;
@@ -310,7 +310,7 @@ pub trait WebAppExt {
/// # Type Parameter
///
/// * `W` - Type RustEmbed contenant les fichiers de la webapp
async fn add_webapp_with_redirect<W>(&mut self, path: &str)
async fn add_webapp_with_redirect<W>(&mut self, path: &str)
where
W: RustEmbed + Clone + Send + Sync + 'static;
}

View File

@@ -8,6 +8,13 @@
<span v-else>🔄</span>
Refresh
</button>
<button
@click="togglePlayback"
:disabled="!nowPlaying || !nowPlaying.stream_url"
class="btn-play"
>
{{ isPlaying ? '⏹️ Stop' : '▶️ Play' }}
</button>
<select v-model="selectedChannel" @change="changeChannel" class="channel-select">
<option v-for="channel in channels" :key="channel.id" :value="channel.id">
{{ channel.name }}
@@ -16,6 +23,21 @@
</div>
</div>
<!-- Audio Player -->
<div v-if="isPlaying || audioError" class="audio-player-container">
<audio
v-if="!audioError"
ref="audioPlayer"
controls
@ended="handleAudioEnded"
@error="handleAudioError"
></audio>
<p v-if="audioError" class="audio-error">{{ audioError }}</p>
<button @click="stopPlayback" class="btn-stop">
{{ audioError ? '✕ Close' : '⏹️ Stop' }}
</button>
</div>
<div v-if="error" class="error-message">
{{ error }}
</div>
@@ -114,7 +136,7 @@
</template>
<script setup>
import { ref, onMounted } from 'vue'
import { ref, onMounted, nextTick } from 'vue'
const API_BASE = '/api/radioparadise'
@@ -123,6 +145,9 @@ const error = ref(null)
const nowPlaying = ref(null)
const channels = ref([])
const selectedChannel = ref(0)
const audioPlayer = ref(null)
const isPlaying = ref(false)
const audioError = ref('')
// Format duration from milliseconds to MM:SS
function formatDuration(ms) {
@@ -178,6 +203,61 @@ function changeChannel() {
// TODO: Implement channel switching in the API
}
function playStream() {
if (!nowPlaying.value?.stream_url) {
audioError.value = 'No stream URL available'
return
}
audioError.value = ''
isPlaying.value = true
nextTick(() => {
const player = audioPlayer.value
if (!player) {
return
}
player.src = nowPlaying.value.stream_url
player.play().catch((e) => {
console.error('Failed to start playback:', e)
audioError.value = `Cannot play stream: ${e.message}`
isPlaying.value = false
})
})
}
function stopPlayback() {
if (audioPlayer.value) {
audioPlayer.value.pause()
audioPlayer.value.currentTime = 0
audioPlayer.value.src = ''
}
isPlaying.value = false
audioError.value = ''
}
function togglePlayback() {
if (isPlaying.value) {
stopPlayback()
} else {
playStream()
}
}
function handleAudioEnded() {
isPlaying.value = false
}
function handleAudioError() {
const audio = audioPlayer.value
if (audio?.error) {
audioError.value = `Audio playback error (code ${audio.error.code})`
} else {
audioError.value = 'Unknown audio playback error'
}
isPlaying.value = false
}
// Initialize on mount
onMounted(async () => {
await fetchChannels()
@@ -239,6 +319,26 @@ onMounted(async () => {
cursor: not-allowed;
}
.btn-play {
background: rgba(46, 204, 113, 0.2);
color: #2ecc71;
border: 1px solid rgba(46, 204, 113, 0.4);
padding: 8px 16px;
border-radius: 4px;
cursor: pointer;
font-weight: bold;
transition: background 0.3s;
}
.btn-play:hover:not(:disabled) {
background: rgba(46, 204, 113, 0.35);
}
.btn-play:disabled {
opacity: 0.5;
cursor: not-allowed;
}
.channel-select {
padding: 8px 12px;
border-radius: 4px;
@@ -513,4 +613,31 @@ onMounted(async () => {
color: #999;
font-size: 0.9em;
}
.audio-player-container {
margin: 12px 0 24px;
padding: 16px;
border-radius: 8px;
background: rgba(0, 0, 0, 0.2);
border: 1px solid #333;
display: flex;
gap: 12px;
align-items: center;
}
.audio-error {
margin: 0;
color: #ff6b6b;
flex: 1;
}
.btn-stop {
padding: 8px 16px;
border-radius: 4px;
border: none;
cursor: pointer;
background: rgba(231, 76, 60, 0.2);
color: #e74c3c;
font-weight: bold;
}
</style>

View File

@@ -55,13 +55,37 @@
</div>
<!-- Services -->
<div v-else class="services-list">
<ServicePanel
v-for="service in device.services"
:key="service.name"
:service="service"
:device-udn="device.udn"
/>
<div v-else class="device-content">
<div class="device-summary">
<div class="meta-row">
<span class="meta-label">UDN:</span>
<code class="meta-value">{{ device.udn }}</code>
</div>
<div class="meta-row" v-if="device.description_url">
<span class="meta-label">Description:</span>
<a
:href="device.description_url"
target="_blank"
rel="noopener"
class="meta-link"
>
View XML
</a>
</div>
<div class="meta-row" v-if="device.base_url">
<span class="meta-label">Base URL:</span>
<code class="meta-value">{{ device.base_url }}</code>
</div>
</div>
<div class="services-list">
<ServicePanel
v-for="service in device.services"
:key="service.name"
:service="service"
:device-udn="device.udn"
/>
</div>
</div>
</div>
</transition>
@@ -353,6 +377,51 @@ onUnmounted(() => {
color: #95a5a6;
}
.device-content {
display: flex;
flex-direction: column;
gap: 1rem;
}
.device-summary {
display: flex;
flex-wrap: wrap;
gap: 0.75rem 1.5rem;
padding: 0.75rem 1rem;
border-radius: 6px;
background: rgba(0, 0, 0, 0.25);
border: 1px solid rgba(52, 152, 219, 0.25);
}
.meta-row {
display: flex;
align-items: center;
gap: 0.5rem;
font-size: 0.9rem;
}
.meta-label {
font-weight: 600;
color: #95a5a6;
}
.meta-value {
background: rgba(44, 62, 80, 0.6);
padding: 0.25rem 0.5rem;
border-radius: 4px;
color: #ecf0f1;
}
.meta-link {
color: #1abc9c;
text-decoration: none;
font-weight: 600;
}
.meta-link:hover {
text-decoration: underline;
}
.services-list {
display: flex;
flex-direction: column;

View File

@@ -30,6 +30,13 @@
<span v-if="action.out_arguments.length > 0" class="badge out-badge" title="Output arguments">
{{ action.out_arguments.length }}
</span>
<span
v-if="action.stateless"
class="badge stateless-badge"
title="Does not mutate state variables"
>
🧊 Stateless
</span>
<span class="expand-indicator">
{{ expandedAction === action.name ? '▼' : '▶' }}
</span>
@@ -38,6 +45,12 @@
<transition name="expand-args">
<div v-if="expandedAction === action.name" class="action-details">
<div v-if="action.stateless" class="action-flags">
<span class="stateless-pill">
Stateless action no state variables updated
</span>
</div>
<!-- Input arguments -->
<div v-if="action.in_arguments.length > 0" class="arguments-section">
<h5 class="section-title">
@@ -299,6 +312,12 @@ onMounted(() => {
border: 1px solid rgba(46, 204, 113, 0.3);
}
.stateless-badge {
background: rgba(155, 89, 182, 0.2);
color: #9b59b6;
border: 1px solid rgba(155, 89, 182, 0.3);
}
.expand-indicator {
color: #3498db;
font-size: 0.9rem;
@@ -316,6 +335,23 @@ onMounted(() => {
border-top: 1px solid rgba(52, 152, 219, 0.2);
}
.action-flags {
margin-top: 0.75rem;
}
.stateless-pill {
display: inline-block;
padding: 0.3rem 0.6rem;
border-radius: 999px;
background: rgba(155, 89, 182, 0.15);
border: 1px solid rgba(155, 89, 182, 0.25);
color: #d2a6e6;
font-size: 0.8rem;
font-weight: 600;
text-transform: uppercase;
letter-spacing: 0.6px;
}
.arguments-section {
margin-top: 1rem;
}

View File

@@ -49,7 +49,10 @@ async fn main() {
tokio::spawn(async move {
let mut source = SourceNode::new();
source.add_subscriber(buffer_tx);
source.generate_chunks(30, 4800, 48000, 440.0).await.unwrap();
source
.generate_chunks(30, 4800, 48000, 440.0)
.await
.unwrap();
});
println!("Waiting for all rooms to finish...\n");

View File

@@ -7,8 +7,7 @@
//! - Système d'événements pour la communication entre nodes
use pmoaudio::{
ChromecastConfig, ChromecastSink, DiskSink, DiskSinkConfig, SourceNode,
VolumeNode,
ChromecastConfig, ChromecastSink, DiskSink, DiskSinkConfig, SourceNode, VolumeNode,
};
use tokio::sync::mpsc;
@@ -146,8 +145,14 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("\n=== Demo completed successfully! ===");
println!("\nSummary:");
println!("- Generated {} chunks of {} samples each", num_chunks, chunk_size);
println!("- Total duration: {:.2} seconds", (num_chunks as usize * chunk_size) as f32 / sample_rate as f32);
println!(
"- Generated {} chunks of {} samples each",
num_chunks, chunk_size
);
println!(
"- Total duration: {:.2} seconds",
(num_chunks as usize * chunk_size) as f32 / sample_rate as f32
);
println!("- Output to Chromecast: Living Room (192.168.1.100)");
println!(
"- Output to file: {}",

View File

@@ -73,10 +73,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// 7. Générer l'audio (10 chunks de 4800 samples à 48kHz = ~1 seconde)
source
.generate_chunks(
10, // nombre de chunks
4800, // samples par chunk (100ms @ 48kHz)
48000, // sample rate
440.0, // fréquence (La 440 Hz)
10, // nombre de chunks
4800, // samples par chunk (100ms @ 48kHz)
48000, // sample rate
440.0, // fréquence (La 440 Hz)
)
.await?;

View File

@@ -90,7 +90,13 @@ impl AudioChunk {
}
/// Crée un nouveau chunk audio avec un gain spécifique
pub fn with_gain(order: u64, left: Vec<f32>, right: Vec<f32>, sample_rate: u32, gain: f32) -> Self {
pub fn with_gain(
order: u64,
left: Vec<f32>,
right: Vec<f32>,
sample_rate: u32,
gain: f32,
) -> Self {
Self {
order,
left: Arc::new(left),

View File

@@ -76,13 +76,13 @@
//! - **RwLock** : Pour partage concurrent du compteur [`TimerNode`]
mod audio_chunk;
mod nodes;
pub mod events;
mod nodes;
pub use audio_chunk::AudioChunk;
pub use events::{
AudioDataEvent, EventPublisher, EventReceiver, NodeEvent, NodeListener,
SourceNameUpdateEvent, VolumeChangeEvent,
AudioDataEvent, EventPublisher, EventReceiver, NodeEvent, NodeListener, SourceNameUpdateEvent,
VolumeChangeEvent,
};
pub use nodes::{
buffer_node::BufferNode,

View File

@@ -1,4 +1,7 @@
use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}};
use crate::{
nodes::{AudioError, MultiSubscriberNode},
AudioChunk,
};
use std::collections::VecDeque;
use std::sync::Arc;
use tokio::sync::{mpsc, RwLock};

View File

@@ -189,10 +189,7 @@ impl ChromecastSink {
// Établir la connexion
self.connect().await?;
let mut stats = ChromecastStats::new(
self.node_id.clone(),
self.config.device_name.clone(),
);
let mut stats = ChromecastStats::new(self.node_id.clone(), self.config.device_name.clone());
// Boucle principale
while let Some(chunk) = self.rx.recv().await {

View File

@@ -1,4 +1,7 @@
use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}};
use crate::{
nodes::{AudioError, MultiSubscriberNode},
AudioChunk,
};
use std::sync::Arc;
use tokio::sync::mpsc;
@@ -56,10 +59,10 @@ impl DecoderNode {
if src_idx < left_data.len() - 1 {
let frac = src_pos - src_idx as f64;
let left_sample =
left_data[src_idx] * (1.0 - frac as f32) + left_data[src_idx + 1] * frac as f32;
let right_sample =
right_data[src_idx] * (1.0 - frac as f32) + right_data[src_idx + 1] * frac as f32;
let left_sample = left_data[src_idx] * (1.0 - frac as f32)
+ left_data[src_idx + 1] * frac as f32;
let right_sample = right_data[src_idx] * (1.0 - frac as f32)
+ right_data[src_idx + 1] * frac as f32;
new_left.push(left_sample);
new_right.push(right_sample);
@@ -69,7 +72,8 @@ impl DecoderNode {
}
}
let new_chunk = AudioChunk::new(chunk.order, new_left, new_right, target_sample_rate);
let new_chunk =
AudioChunk::new(chunk.order, new_left, new_right, target_sample_rate);
self.subscribers.push(Arc::new(new_chunk)).await?;
}
}

View File

@@ -3,11 +3,7 @@
//! Ce module fournit un sink qui écrit les chunks audio sur disque,
//! avec support de la dérivation automatique du nom de fichier depuis la source.
use crate::{
events::SourceNameUpdateEvent,
nodes::AudioError,
AudioChunk,
};
use crate::{events::SourceNameUpdateEvent, nodes::AudioError, AudioChunk};
use std::path::PathBuf;
use std::sync::Arc;
use tokio::fs::File;
@@ -161,7 +157,13 @@ impl DiskSink {
// Nettoyer le nom de la source pour en faire un nom de fichier valide
let clean_name = name
.chars()
.map(|c| if c.is_alphanumeric() || c == '_' || c == '-' { c } else { '_' })
.map(|c| {
if c.is_alphanumeric() || c == '_' || c == '-' {
c
} else {
'_'
}
})
.collect::<String>();
format!("{}.{}", clean_name, self.config.format.extension())
@@ -180,9 +182,9 @@ impl DiskSink {
// Créer le répertoire parent si nécessaire
if let Some(parent) = path.parent() {
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| AudioError::ProcessingError(format!("Failed to create directory: {}", e)))?;
tokio::fs::create_dir_all(parent).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to create directory: {}", e))
})?;
}
// Créer le writer approprié selon le format
@@ -328,10 +330,9 @@ impl AudioFileWriter {
bytes.extend_from_slice(&sample_i16.to_le_bytes());
}
self.file
.write_all(&bytes)
.await
.map_err(|e| AudioError::ProcessingError(format!("Failed to write audio data: {}", e)))?;
self.file.write_all(&bytes).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to write audio data: {}", e))
})?;
self.total_samples += chunk.len();
Ok(())
@@ -361,10 +362,9 @@ impl AudioFileWriter {
header.extend_from_slice(b"data");
header.extend_from_slice(&0u32.to_le_bytes()); // Taille des données (à mettre à jour)
self.file
.write_all(&header)
.await
.map_err(|e| AudioError::ProcessingError(format!("Failed to write WAV header: {}", e)))?;
self.file.write_all(&header).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to write WAV header: {}", e))
})?;
Ok(())
}
@@ -378,24 +378,37 @@ impl AudioFileWriter {
// Positionner au début et réécrire les tailles
use tokio::io::AsyncSeekExt;
self.file.seek(std::io::SeekFrom::Start(4)).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to seek in file: {}", e))
})?;
self.file.write_all(&file_size.to_le_bytes()).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to update file size: {}", e))
})?;
self.file
.seek(std::io::SeekFrom::Start(4))
.await
.map_err(|e| {
AudioError::ProcessingError(format!("Failed to seek in file: {}", e))
})?;
self.file
.write_all(&file_size.to_le_bytes())
.await
.map_err(|e| {
AudioError::ProcessingError(format!("Failed to update file size: {}", e))
})?;
self.file.seek(std::io::SeekFrom::Start(40)).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to seek in file: {}", e))
})?;
self.file.write_all(&data_size.to_le_bytes()).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to update data size: {}", e))
})?;
self.file
.seek(std::io::SeekFrom::Start(40))
.await
.map_err(|e| {
AudioError::ProcessingError(format!("Failed to seek in file: {}", e))
})?;
self.file
.write_all(&data_size.to_le_bytes())
.await
.map_err(|e| {
AudioError::ProcessingError(format!("Failed to update data size: {}", e))
})?;
}
self.file.flush().await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to flush file: {}", e))
})?;
self.file
.flush()
.await
.map_err(|e| AudioError::ProcessingError(format!("Failed to flush file: {}", e)))?;
Ok(())
}

View File

@@ -1,4 +1,7 @@
use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}};
use crate::{
nodes::{AudioError, MultiSubscriberNode},
AudioChunk,
};
use std::sync::Arc;
use tokio::sync::mpsc;
@@ -46,12 +49,8 @@ impl DspNode {
*sample *= self.gain;
}
let new_chunk = AudioChunk::new(
chunk.order,
left_data,
right_data,
chunk.sample_rate,
);
let new_chunk =
AudioChunk::new(chunk.order, left_data, right_data, chunk.sample_rate);
self.subscribers.push(Arc::new(new_chunk)).await?;
}
@@ -118,12 +117,7 @@ impl LowPassDspNode {
new_right.push(self.prev_right);
}
let new_chunk = AudioChunk::new(
chunk.order,
new_left,
new_right,
chunk.sample_rate,
);
let new_chunk = AudioChunk::new(chunk.order, new_left, new_right, chunk.sample_rate);
self.subscribers.push(Arc::new(new_chunk)).await?;
}

View File

@@ -59,10 +59,7 @@ impl SingleSubscriberNode {
}
pub async fn push(&self, chunk: Arc<AudioChunk>) -> Result<(), AudioError> {
self.tx
.send(chunk)
.await
.map_err(|_| AudioError::SendError)
self.tx.send(chunk).await.map_err(|_| AudioError::SendError)
}
}

View File

@@ -197,7 +197,9 @@ impl MpdSink {
/// Envoie un chunk au serveur MPD (mock)
async fn send_chunk(&self, _chunk: &AudioChunk) -> Result<(), AudioError> {
if !self.connected {
return Err(AudioError::ProcessingError("Not connected to MPD".to_string()));
return Err(AudioError::ProcessingError(
"Not connected to MPD".to_string(),
));
}
// Dans une vraie implémentation:

View File

@@ -1,4 +1,4 @@
use crate::{AudioChunk, nodes::AudioError};
use crate::{nodes::AudioError, AudioChunk};
use std::sync::Arc;
use tokio::sync::mpsc;
@@ -115,11 +115,13 @@ impl SinkStats {
let sum_squares_left: f64 = chunk.left.iter().map(|&x| (x * x) as f64).sum();
let sum_squares_right: f64 = chunk.right.iter().map(|&x| (x * x) as f64).sum();
self.rms_left = ((self.rms_left.powi(2) * (self.total_samples - chunk.len() as u64) as f64
self.rms_left = ((self.rms_left.powi(2)
* (self.total_samples - chunk.len() as u64) as f64
+ sum_squares_left)
/ self.total_samples as f64)
.sqrt();
self.rms_right = ((self.rms_right.powi(2) * (self.total_samples - chunk.len() as u64) as f64
self.rms_right = ((self.rms_right.powi(2)
* (self.total_samples - chunk.len() as u64) as f64
+ sum_squares_right)
/ self.total_samples as f64)
.sqrt();

View File

@@ -1,4 +1,7 @@
use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}};
use crate::{
nodes::{AudioError, MultiSubscriberNode},
AudioChunk,
};
use std::sync::Arc;
use tokio::sync::mpsc;
@@ -63,12 +66,8 @@ impl SourceNode {
sample_rate: u32,
) -> Result<(), AudioError> {
for i in 0..count {
let chunk = AudioChunk::new(
i,
vec![0.0; chunk_size],
vec![0.0; chunk_size],
sample_rate,
);
let chunk =
AudioChunk::new(i, vec![0.0; chunk_size], vec![0.0; chunk_size], sample_rate);
self.subscribers.push(Arc::new(chunk)).await?;
}
Ok(())

View File

@@ -1,4 +1,7 @@
use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}};
use crate::{
nodes::{AudioError, MultiSubscriberNode},
AudioChunk,
};
use std::sync::Arc;
use tokio::sync::{mpsc, RwLock};

View File

@@ -32,7 +32,10 @@ async fn test_complete_pipeline() {
tokio::spawn(async move {
let mut source = SourceNode::new();
source.add_subscriber(decoder_tx);
source.generate_chunks(10, 4800, 48000, 440.0).await.unwrap();
source
.generate_chunks(10, 4800, 48000, 440.0)
.await
.unwrap();
});
// Attendre la fin
@@ -73,7 +76,10 @@ async fn test_multiroom_buffering() {
tokio::spawn(async move {
let mut source = SourceNode::new();
source.add_subscriber(buffer_tx);
source.generate_chunks(20, 1000, 48000, 440.0).await.unwrap();
source
.generate_chunks(20, 1000, 48000, 440.0)
.await
.unwrap();
});
let stats1 = sink1_handle.await.unwrap();
@@ -108,7 +114,10 @@ async fn test_timer_accuracy() {
source.add_subscriber(timer_tx);
// 48000 samples à 48kHz = 1 seconde
source.generate_chunks(1, 48000, 48000, 440.0).await.unwrap();
source
.generate_chunks(1, 48000, 48000, 440.0)
.await
.unwrap();
});
sink_handle.await.unwrap();

View File

@@ -30,12 +30,20 @@ async fn main() -> anyhow::Result<()> {
let header = &data[0..4];
if header == b"fLaC" {
println!("✓ File is FLAC!");
} else if header[0..3] == *b"ID3" || (header.len() >= 2 && header[0] == 0xFF && (header[1] & 0xE0) == 0xE0) {
} else if header[0..3] == *b"ID3"
|| (header.len() >= 2 && header[0] == 0xFF && (header[1] & 0xE0) == 0xE0)
{
println!("✗ File is still MP3!");
println!(" Header: {:02X} {:02X} {:02X} {:02X}", header[0], header[1], header[2], header[3]);
println!(
" Header: {:02X} {:02X} {:02X} {:02X}",
header[0], header[1], header[2], header[3]
);
} else {
println!("? Unknown format");
println!(" Header: {:02X} {:02X} {:02X} {:02X}", header[0], header[1], header[2], header[3]);
println!(
" Header: {:02X} {:02X} {:02X} {:02X}",
header[0], header[1], header[2], header[3]
);
}
}

View File

@@ -70,7 +70,10 @@ fn create_flac_transformer() -> StreamTransformer {
buffer.extend_from_slice(&chunk);
}
tracing::debug!("Downloaded {} bytes total, starting FLAC conversion", buffer.len());
tracing::debug!(
"Downloaded {} bytes total, starting FLAC conversion",
buffer.len()
);
// 2. Si c'est déjà du FLAC, on l'écrit directement
if buffer.len() >= 4 && &buffer[0..4] == b"fLaC" {
@@ -85,21 +88,26 @@ fn create_flac_transformer() -> StreamTransformer {
// 3. Décoder l'audio avec Symphonia
let (samples, channels, sample_rate, bits_per_sample) = {
use std::io::Cursor;
use symphonia::core::audio::SampleBuffer;
use symphonia::core::codecs::{DecoderOptions, CODEC_TYPE_NULL};
use symphonia::core::errors::Error as SymphoniaError;
use symphonia::core::formats::FormatOptions;
use symphonia::core::io::MediaSourceStream;
use symphonia::core::meta::MetadataOptions;
use symphonia::core::probe::Hint;
use symphonia::core::errors::Error as SymphoniaError;
use std::io::Cursor;
let cursor = Cursor::new(buffer);
let mss = MediaSourceStream::new(Box::new(cursor), Default::default());
let hint = Hint::new();
let probed = symphonia::default::get_probe()
.format(&hint, mss, &FormatOptions::default(), &MetadataOptions::default())
.format(
&hint,
mss,
&FormatOptions::default(),
&MetadataOptions::default(),
)
.map_err(|e| format!("Failed to probe format: {}", e))?;
let mut format = probed.format;
@@ -114,15 +122,18 @@ fn create_flac_transformer() -> StreamTransformer {
.make(&track.codec_params, &DecoderOptions::default())
.map_err(|e| format!("Failed to create decoder: {}", e))?;
let channels = track.codec_params.channels
let channels = track
.codec_params
.channels
.ok_or_else(|| "No channel info".to_string())?
.count();
let sample_rate = track.codec_params.sample_rate
let sample_rate = track
.codec_params
.sample_rate
.ok_or_else(|| "No sample rate info".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);
let mut samples_i32 = Vec::new();
let track_id = track.id;
@@ -135,7 +146,9 @@ fn create_flac_transformer() -> StreamTransformer {
decoder.reset();
continue;
}
Err(SymphoniaError::IoError(e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => {
Err(SymphoniaError::IoError(e))
if e.kind() == std::io::ErrorKind::UnexpectedEof =>
{
break;
}
Err(e) => return Err(format!("Decode error: {}", e)),
@@ -183,13 +196,13 @@ fn create_flac_transformer() -> StreamTransformer {
tracing::debug!("Normalizing to 16-bit");
let samples = samples_i32.iter().map(|&s| (s >> 16) as i32).collect();
(samples, 16)
},
}
17..=24 => {
// Pour 17-24 bits, normaliser vers la plage 24-bit
tracing::debug!("Normalizing to 24-bit");
let samples = samples_i32.iter().map(|&s| (s >> 8) as i32).collect();
(samples, 24)
},
}
_ => {
// Pour 25-32 bits, garder la pleine échelle i32
tracing::debug!("Keeping 32-bit");
@@ -200,15 +213,20 @@ fn create_flac_transformer() -> StreamTransformer {
(normalized_samples, channels, sample_rate, target_bits)
};
tracing::debug!("Encoding to FLAC: {} samples, {} channels, {} Hz, {} bits",
samples.len(), channels, sample_rate, bits_per_sample);
tracing::debug!(
"Encoding to FLAC: {} samples, {} channels, {} Hz, {} bits",
samples.len(),
channels,
sample_rate,
bits_per_sample
);
// 4. Encoder en FLAC avec flacenc
// Note: L'encodage FLAC est une opération bloquante/CPU-intensive,
// donc nous l'exécutons dans un thread bloquant pour ne pas bloquer le runtime Tokio
let flac_data = tokio::task::spawn_blocking(move || {
use flacenc::component::BitRepr;
use flacenc::bitsink::ByteSink;
use flacenc::component::BitRepr;
use flacenc::error::Verify;
let config = flacenc::config::Encoder::default()
@@ -222,15 +240,13 @@ fn create_flac_transformer() -> StreamTransformer {
sample_rate as usize,
);
let flac_stream = flacenc::encode_with_fixed_block_size(
&config,
source,
config.block_size,
)
.map_err(|e| format!("FLAC encode error: {:?}", e))?;
let flac_stream =
flacenc::encode_with_fixed_block_size(&config, source, config.block_size)
.map_err(|e| format!("FLAC encode error: {:?}", e))?;
let mut sink = ByteSink::new();
flac_stream.write(&mut sink)
flac_stream
.write(&mut sink)
.map_err(|e| format!("FLAC write error: {:?}", e))?;
Ok::<Vec<u8>, String>(sink.into_inner())
@@ -241,7 +257,9 @@ fn create_flac_transformer() -> StreamTransformer {
tracing::debug!("FLAC encoding complete: {} bytes", flac_data.len());
// 5. Écrire le fichier FLAC
file.write_all(&flac_data).await.map_err(|e| e.to_string())?;
file.write_all(&flac_data)
.await
.map_err(|e| e.to_string())?;
file.flush().await.map_err(|e| e.to_string())?;
// 6. Mettre à jour la progression finale
@@ -328,13 +346,17 @@ pub async fn add_with_metadata_extraction(
let metadata_json = serde_json::to_string(&metadata)?;
// Stocker dans la DB
cache.db.update_metadata(&pk, &metadata_json)
cache
.db
.update_metadata(&pk, &metadata_json)
.map_err(|e| anyhow::anyhow!("Database error: {}", e))?;
// Mettre à jour la collection si les métadonnées en fournissent une
if collection.is_none() {
if let Some(auto_collection) = metadata.collection_key() {
cache.db.add(&pk, url, Some(&auto_collection))
cache
.db
.add(&pk, url, Some(&auto_collection))
.map_err(|e| anyhow::anyhow!("Database error: {}", e))?;
}
}
@@ -366,7 +388,9 @@ pub async fn add_with_metadata_extraction(
/// # }
/// ```
pub fn get_metadata(cache: &Cache, pk: &str) -> Result<crate::metadata::AudioMetadata> {
let metadata_json = cache.db.get_metadata_json(pk)
let metadata_json = cache
.db
.get_metadata_json(pk)
.map_err(|e| anyhow::anyhow!("Database error: {}", e))?
.ok_or_else(|| anyhow::anyhow!("No metadata found for pk: {}", pk))?;
@@ -375,4 +399,3 @@ pub fn get_metadata(cache: &Cache, pk: &str) -> Result<crate::metadata::AudioMet
Ok(metadata)
}

View File

@@ -4,6 +4,7 @@
//! pour standardiser le stockage dans le cache.
use anyhow::{anyhow, Result};
use std::io::Cursor;
use symphonia::core::audio::SampleBuffer;
use symphonia::core::codecs::{DecoderOptions, CODEC_TYPE_NULL};
use symphonia::core::errors::Error as SymphoniaError;
@@ -11,7 +12,6 @@ use symphonia::core::formats::FormatOptions;
use symphonia::core::io::MediaSourceStream;
use symphonia::core::meta::MetadataOptions;
use symphonia::core::probe::Hint;
use std::io::Cursor;
/// Convertit des données audio en FLAC
///
@@ -54,7 +54,12 @@ pub fn convert_to_flac(data: &[u8], extension: Option<&str>) -> Result<Vec<u8>>
// Prober le format
let probed = symphonia::default::get_probe()
.format(&hint, mss, &FormatOptions::default(), &MetadataOptions::default())
.format(
&hint,
mss,
&FormatOptions::default(),
&MetadataOptions::default(),
)
.map_err(|e| anyhow!("Impossible de détecter le format audio: {}", e))?;
let mut format = probed.format;

View File

@@ -131,14 +131,14 @@
//! - [`pmoserver`] : Serveur HTTP
pub mod cache;
pub mod metadata;
pub mod flac;
pub mod metadata;
#[cfg(feature = "pmoserver")]
pub mod openapi;
// Re-exports principaux
pub use cache::{Cache, AudioConfig, new_cache, add_with_metadata_extraction, get_metadata};
pub use cache::{add_with_metadata_extraction, get_metadata, new_cache, AudioConfig, Cache};
pub use metadata::AudioMetadata;
#[cfg(feature = "pmoserver")]
@@ -161,7 +161,11 @@ pub trait AudioCacheExt {
/// # Returns
///
/// * `Arc<Cache>` - Instance partagée du cache
async fn init_audio_cache(&mut self, cache_dir: &str, limit: usize) -> anyhow::Result<std::sync::Arc<Cache>>;
async fn init_audio_cache(
&mut self,
cache_dir: &str,
limit: usize,
) -> anyhow::Result<std::sync::Arc<Cache>>;
/// Initialise le cache audio avec la configuration par défaut.
///
@@ -170,7 +174,7 @@ pub trait AudioCacheExt {
}
#[cfg(feature = "pmoserver")]
use pmocache::pmoserver_ext::{create_file_router, create_api_router};
use pmocache::pmoserver_ext::{create_api_router, create_file_router};
#[cfg(feature = "pmoserver")]
use std::sync::Arc;
#[cfg(feature = "pmoserver")]
@@ -178,14 +182,18 @@ use utoipa::OpenApi;
#[cfg(feature = "pmoserver")]
impl AudioCacheExt for pmoserver::Server {
async fn init_audio_cache(&mut self, cache_dir: &str, limit: usize) -> anyhow::Result<Arc<Cache>> {
async fn init_audio_cache(
&mut self,
cache_dir: &str,
limit: usize,
) -> anyhow::Result<Arc<Cache>> {
let cache = Arc::new(crate::cache::new_cache(cache_dir, limit)?);
// Router de fichiers pour servir les pistes FLAC
// Routes: GET /audio/tracks/{pk} et GET /audio/tracks/{pk}/{param}
let file_router = create_file_router(
cache.clone(),
"audio/flac" // Content-Type
"audio/flac", // Content-Type
);
self.add_router("/", file_router).await;

View File

@@ -87,12 +87,12 @@ impl AudioMetadata {
/// println!("Titre: {:?}", metadata.title);
/// ```
pub fn from_file(path: &Path) -> Result<Self> {
let tagged_file = Probe::open(path)?
.options(ParseOptions::new())
.read()?;
let tagged_file = Probe::open(path)?.options(ParseOptions::new()).read()?;
let properties = tagged_file.properties();
let tag = tagged_file.primary_tag().or_else(|| tagged_file.first_tag());
let tag = tagged_file
.primary_tag()
.or_else(|| tagged_file.first_tag());
let mut metadata = Self {
title: None,
@@ -138,7 +138,9 @@ impl AudioMetadata {
.read()?;
let properties = tagged_file.properties();
let tag = tagged_file.primary_tag().or_else(|| tagged_file.first_tag());
let tag = tagged_file
.primary_tag()
.or_else(|| tagged_file.first_tag());
let mut metadata = Self {
title: None,

View File

@@ -11,7 +11,9 @@ fn main() {
println!(" dl.wait_until_finished().await?;\n");
println!("2. Téléchargement avec transformation:");
println!(" let transformer: StreamTransformer = Box::new(|response, mut file, update_progress| {{");
println!(
" let transformer: StreamTransformer = Box::new(|response, mut file, update_progress| {{"
);
println!(" Box::pin(async move {{");
println!(" let mut stream = response.bytes_stream();");
println!(" let mut total = 0u64;");

View File

@@ -1,7 +1,7 @@
// Exemple d'utilisation du module download avec transformations
use pmocache::download::{download_with_transformer, StreamTransformer};
use futures_util::StreamExt;
use pmocache::download::{download_with_transformer, StreamTransformer};
use tokio::io::AsyncWriteExt;
/// Exemple de transformer qui compresse les données en gzip
@@ -63,7 +63,13 @@ fn create_uppercase_transformer() -> StreamTransformer {
// Transformer en majuscules (seulement pour texte ASCII)
let transformed: Vec<u8> = chunk
.iter()
.map(|&b| if b.is_ascii_lowercase() { b.to_ascii_uppercase() } else { b })
.map(|&b| {
if b.is_ascii_lowercase() {
b.to_ascii_uppercase()
} else {
b
}
})
.collect();
file.write_all(&transformed)
@@ -210,7 +216,10 @@ async fn main() {
Ok(_) => {
println!(" ✓ Téléchargement terminé!");
println!(" - Taille source: {} bytes", dl.current_size().await);
println!(" - Taille transformée: {} bytes", dl.transformed_size().await);
println!(
" - Taille transformée: {} bytes",
dl.transformed_size().await
);
}
Err(e) => {
eprintln!(" ✗ Erreur: {}", e);
@@ -223,16 +232,15 @@ async fn main() {
let _ = std::fs::remove_file(&skip_file);
let transformer = create_skip_header_transformer(100);
let dl = download_with_transformer(
&skip_file,
"https://www.rust-lang.org/",
Some(transformer),
);
let dl = download_with_transformer(&skip_file, "https://www.rust-lang.org/", Some(transformer));
match dl.wait_until_finished().await {
Ok(_) => {
println!(" ✓ Téléchargement terminé!");
println!(" - Taille transformée: {} bytes", dl.transformed_size().await);
println!(
" - Taille transformée: {} bytes",
dl.transformed_size().await
);
}
Err(e) => {
eprintln!(" ✗ Erreur: {}", e);
@@ -254,7 +262,10 @@ async fn main() {
match dl.wait_until_finished().await {
Ok(_) => {
println!(" ✓ Téléchargement terminé!");
println!(" - Taille transformée: {} bytes", dl.transformed_size().await);
println!(
" - Taille transformée: {} bytes",
dl.transformed_size().await
);
}
Err(e) => {
eprintln!(" ✗ Erreur: {}", e);

View File

@@ -92,9 +92,7 @@ pub struct ErrorResponse {
/// Liste tous les items en cache avec leurs statistiques
///
/// 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() {
Ok(entries) => (StatusCode::OK, Json(entries)).into_response(),
Err(e) => (
@@ -192,7 +190,10 @@ pub async fn add_item<C: CacheConfig>(
.into_response();
}
match cache.add_from_url(&req.url, req.collection.as_deref()).await {
match cache
.add_from_url(&req.url, req.collection.as_deref())
.await
{
Ok(pk) => (
StatusCode::CREATED,
Json(AddItemResponse {
@@ -268,9 +269,7 @@ pub async fn delete_item<C: CacheConfig>(
/// Purge complètement le cache
///
/// Supprime tous les items et vide la base de données. Opération irréversible.
pub async fn purge_cache<C: CacheConfig>(
State(cache): State<Arc<Cache<C>>>,
) -> impl IntoResponse {
pub async fn purge_cache<C: CacheConfig>(State(cache): State<Arc<Cache<C>>>) -> impl IntoResponse {
match cache.purge().await {
Ok(_) => (
StatusCode::OK,

View File

@@ -3,9 +3,9 @@
//! Ce module fournit une interface générique pour gérer un cache de fichiers
//! avec métadonnées dans une base de données SQLite.
use crate::cache_trait::{FileCache, pk_from_url};
use crate::cache_trait::{pk_from_url, FileCache};
use crate::db::DB;
use crate::download::{Download, download_with_transformer, StreamTransformer};
use crate::download::{download_with_transformer, Download, StreamTransformer};
use anyhow::{anyhow, Result};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
@@ -26,10 +26,10 @@ pub trait CacheConfig: Send + Sync {
"file"
}
/// Cache name (ex: "covers", "audio", "cache")
fn cache_name() -> &'static str {
fn cache_name() -> &'static str {
"cache"
}
/// Default param extension ("orig")
/// Default param extension ("orig")
fn default_param() -> &'static str {
"orig"
}
@@ -155,11 +155,7 @@ impl<C: CacheConfig> Cache<C> {
// Lancer le téléchargement avec transformer
let transformer = self.transformer_factory.as_ref().map(|f| f());
let download = download_with_transformer(
&file_path,
url,
transformer,
);
let download = download_with_transformer(&file_path, url, transformer);
// Stocker dans la map des downloads en cours
{
@@ -228,7 +224,6 @@ impl<C: CacheConfig> Cache<C> {
self.add_from_url(url, collection).await
}
/// Récupère le chemin d'un fichier dans le cache
///
/// # Arguments
@@ -289,8 +284,11 @@ impl<C: CacheConfig> Cache<C> {
let file_path = self.file_path(&entry.pk);
if !file_path.exists() {
// Re-télécharger le fichier manquant
match self.add_from_url(&entry.source_url, entry.collection.as_deref()).await {
Ok(_) => {},
match self
.add_from_url(&entry.source_url, entry.collection.as_deref())
.await
{
Ok(_) => {}
Err(_) => {
// Si le téléchargement échoue, supprimer l'entrée DB
self.db.delete(&entry.pk)?;
@@ -413,7 +411,9 @@ impl<C: CacheConfig> Cache<C> {
/// * `min_size` - Taille minimale attendue en bytes
pub async fn wait_until_min_size(&self, pk: &str, min_size: u64) -> Result<()> {
if let Some(download) = self.get_download(pk).await {
download.wait_until_min_size(min_size).await
download
.wait_until_min_size(min_size)
.await
.map_err(|e| anyhow!("Download error: {}", e))
} else {
// Déjà terminé ou n'existe pas
@@ -432,7 +432,9 @@ impl<C: CacheConfig> Cache<C> {
/// * `pk` - Clé primaire du fichier
pub async fn wait_until_finished(&self, pk: &str) -> Result<()> {
if let Some(download) = self.get_download(pk).await {
download.wait_until_finished().await
download
.wait_until_finished()
.await
.map_err(|e| anyhow!("Download error: {}", e))
} else {
// Déjà terminé ou n'existe pas
@@ -460,7 +462,8 @@ impl<C: CacheConfig> Cache<C> {
///
/// Format: `{pk}.{qualifier}.{extension}`
pub fn file_path_with_qualifier(&self, pk: &str, qualifier: &str) -> PathBuf {
self.dir.join(format!("{}.{}.{}", pk, qualifier, C::file_extension()))
self.dir
.join(format!("{}.{}.{}", pk, qualifier, C::file_extension()))
}
/// Valide les données avant de les stocker
@@ -499,7 +502,9 @@ impl<C: CacheConfig> Cache<C> {
while let Ok(Some(dir_entry)) = dir_entries.next_entry().await {
if let Some(filename) = dir_entry.file_name().to_str() {
// Format: {pk}.{param}.{ext}
if filename.starts_with(&entry.pk) && filename.starts_with(&format!("{}.", entry.pk)) {
if filename.starts_with(&entry.pk)
&& filename.starts_with(&format!("{}.", entry.pk))
{
let _ = tokio::fs::remove_file(dir_entry.path()).await;
}
}
@@ -515,15 +520,18 @@ impl<C: CacheConfig> Cache<C> {
}
if removed > 0 {
tracing::info!("LRU eviction: removed {} old entries (cache size: {} -> {})",
removed, count, count - removed);
tracing::info!(
"LRU eviction: removed {} old entries (cache size: {} -> {})",
removed,
count,
count - removed
);
}
Ok(removed)
}
}
/// Implémentation du trait FileCache pour Cache
impl<C: CacheConfig> FileCache<C> for Cache<C> {
fn get_cache_dir(&self) -> &Path {

View File

@@ -3,9 +3,9 @@
//! Ce module fournit une interface générique pour gérer les métadonnées
//! des éléments en cache, avec tracking des accès et des statistiques.
use rusqlite::{Connection, params};
use serde::Serialize;
use chrono::Utc;
use rusqlite::{params, Connection};
use serde::Serialize;
use std::path::Path;
use std::sync::Mutex;
@@ -32,7 +32,10 @@ pub struct CacheEntry {
#[cfg_attr(feature = "openapi", schema(example = "2025-01-15T10:30:00Z"))]
pub last_used: Option<String>,
/// Métadonnées JSON optionnelles (ex: métadonnées audio, EXIF images, etc.)
#[cfg_attr(feature = "openapi", schema(example = r#"{"title":"Track","artist":"Artist"}"#))]
#[cfg_attr(
feature = "openapi",
schema(example = r#"{"title":"Track","artist":"Artist"}"#)
)]
pub metadata_json: Option<String>,
}
@@ -161,20 +164,16 @@ impl DB {
self.table_name
);
conn.query_row(
&sql,
[pk],
|row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
},
)
conn.query_row(&sql, [pk], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})
}
/// Met à jour le compteur d'accès et la date du dernier accès
@@ -189,10 +188,7 @@ impl DB {
self.table_name
);
conn.execute(
&sql,
params![Utc::now().to_rfc3339(), pk],
)?;
conn.execute(&sql, params![Utc::now().to_rfc3339(), pk])?;
Ok(())
}
@@ -215,17 +211,18 @@ impl DB {
let mut stmt = conn.prepare(&sql)?;
let entries = stmt.query_map([], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
let entries = stmt
.query_map([], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(entries)
}
@@ -244,17 +241,18 @@ impl DB {
let mut stmt = conn.prepare(&sql)?;
let entries = stmt.query_map([collection], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
let entries = stmt
.query_map([collection], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(entries)
}
@@ -319,17 +317,18 @@ impl DB {
let mut stmt = conn.prepare(&sql)?;
let entries = stmt.query_map([limit], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
let entries = stmt
.query_map([limit], |row| {
Ok(CacheEntry {
pk: row.get(0)?,
source_url: row.get(1)?,
collection: row.get(2)?,
hits: row.get(3)?,
last_used: row.get(4)?,
metadata_json: row.get(5)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(entries)
}

View File

@@ -1,3 +1,4 @@
use futures_util::Future;
use std::fs::File;
use std::io;
use std::path::{Path, PathBuf};
@@ -5,7 +6,6 @@ use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::RwLock;
use futures_util::Future;
/// Type pour une fonction de transformation de stream
///
@@ -248,20 +248,16 @@ async fn download_impl(
.map_err(|e| e.to_string())?;
// Lancer la requête
let response = client
.get(&url)
.send()
.await
.map_err(|e| {
let error = format!("Failed to fetch URL: {}", e);
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(async {
let mut s = state.write().await;
s.error = Some(error.clone());
});
let response = client.get(&url).send().await.map_err(|e| {
let error = format!("Failed to fetch URL: {}", e);
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(async {
let mut s = state.write().await;
s.error = Some(error.clone());
});
error
})?;
});
error
})?;
// Vérifier le statut
if !response.status().is_success() {
@@ -279,31 +275,30 @@ async fn download_impl(
}
// Créer le fichier de destination
let file = tokio::fs::File::create(&filename)
.await
.map_err(|e| {
let error = format!("Failed to create file: {}", e);
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(async {
let mut s = state.write().await;
s.error = Some(error.clone());
s.finished = true;
});
let file = tokio::fs::File::create(&filename).await.map_err(|e| {
let error = format!("Failed to create file: {}", e);
tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(async {
let mut s = state.write().await;
s.error = Some(error.clone());
s.finished = true;
});
error
})?;
});
error
})?;
// Si un transformer est fourni, l'utiliser
if let Some(transformer) = transformer {
// Créer un callback pour mettre à jour la progression
let state_clone = Arc::clone(&state);
let progress_callback: Arc<dyn Fn(u64) + Send + Sync> = Arc::new(move |transformed_bytes| {
let state = Arc::clone(&state_clone);
tokio::spawn(async move {
let mut s = state.write().await;
s.transformed_size = transformed_bytes;
let progress_callback: Arc<dyn Fn(u64) + Send + Sync> =
Arc::new(move |transformed_bytes| {
let state = Arc::clone(&state_clone);
tokio::spawn(async move {
let mut s = state.write().await;
s.transformed_size = transformed_bytes;
});
});
});
// Appeler le transformer
match transformer(response, file, progress_callback).await {
@@ -331,8 +326,8 @@ async fn default_download(
mut file: tokio::fs::File,
state: Arc<RwLock<DownloadState>>,
) -> Result<(), String> {
use tokio::io::AsyncWriteExt;
use futures_util::StreamExt;
use tokio::io::AsyncWriteExt;
let mut stream = response.bytes_stream();

View File

@@ -122,9 +122,9 @@
//! - [`pmocovers`] : Cache d'images avec conversion WebP
//! - [`pmoaudiocache`] : Cache de pistes audio
pub mod db;
pub mod cache;
pub mod cache_trait;
pub mod db;
pub mod download;
#[cfg(feature = "pmoserver")]
@@ -136,16 +136,13 @@ pub mod api;
#[cfg(feature = "openapi")]
pub mod openapi;
pub use db::{DB, CacheEntry};
pub use cache::{Cache, CacheConfig};
pub use cache_trait::{FileCache, pk_from_url};
pub use download::{Download, download, download_with_transformer, StreamTransformer};
pub use cache_trait::{pk_from_url, FileCache};
pub use db::{CacheEntry, DB};
pub use download::{download, download_with_transformer, Download, StreamTransformer};
#[cfg(feature = "pmoserver")]
pub use pmoserver_ext::{create_file_router, create_api_router, GenericCacheExt};
pub use pmoserver_ext::{create_api_router, create_file_router, GenericCacheExt};
#[cfg(all(feature = "pmoserver", feature = "openapi"))]
pub use api::{
DownloadStatus, AddItemRequest, AddItemResponse,
DeleteItemResponse, ErrorResponse,
};
pub use api::{AddItemRequest, AddItemResponse, DeleteItemResponse, DownloadStatus, ErrorResponse};

View File

@@ -49,15 +49,15 @@ use axum::{
Router,
};
#[cfg(feature = "pmoserver")]
use std::future::Future;
#[cfg(feature = "pmoserver")]
use std::pin::Pin;
#[cfg(feature = "pmoserver")]
use std::sync::Arc;
#[cfg(feature = "pmoserver")]
use tokio_util::io::ReaderStream;
#[cfg(feature = "pmoserver")]
use tracing::warn;
#[cfg(feature = "pmoserver")]
use std::pin::Pin;
#[cfg(feature = "pmoserver")]
use std::future::Future;
/// Type pour le callback de génération de param
///
@@ -75,16 +75,20 @@ use std::future::Future;
/// Les données générées ou None si le param n'est pas supporté
#[cfg(feature = "pmoserver")]
pub type ParamGenerator<C> = Arc<
dyn Fn(Arc<Cache<C>>, String, String)
-> Pin<Box<dyn Future<Output = Option<Vec<u8>>> + Send>>
+ Send + Sync
dyn Fn(Arc<Cache<C>>, String, String) -> Pin<Box<dyn Future<Output = Option<Vec<u8>>> + Send>>
+ Send
+ Sync,
>;
/// Handler générique pour GET /{cache_name}/{cache_type}/{pk}
/// Sert un fichier avec le param par défaut
#[cfg(feature = "pmoserver")]
async fn get_file<C: CacheConfig + 'static>(
State((cache, content_type, param_generator)): State<(Arc<Cache<C>>, &'static str, Option<ParamGenerator<C>>)>,
State((cache, content_type, param_generator)): State<(
Arc<Cache<C>>,
&'static str,
Option<ParamGenerator<C>>,
)>,
Path(pk): Path<String>,
) -> Response {
// Utiliser le param par défaut
@@ -96,7 +100,11 @@ async fn get_file<C: CacheConfig + 'static>(
/// Sert un fichier avec un param spécifique
#[cfg(feature = "pmoserver")]
async fn get_file_with_param<C: CacheConfig + 'static>(
State((cache, content_type, param_generator)): State<(Arc<Cache<C>>, &'static str, Option<ParamGenerator<C>>)>,
State((cache, content_type, param_generator)): State<(
Arc<Cache<C>>,
&'static str,
Option<ParamGenerator<C>>,
)>,
Path((pk, param)): Path<(String, String)>,
) -> Response {
serve_file_with_streaming(&cache, &pk, &param, content_type, param_generator).await
@@ -122,11 +130,7 @@ async fn serve_file_with_streaming<C: CacheConfig>(
if let Some(generator) = param_generator {
if let Some(data) = generator(cache.clone(), pk.to_string(), param.to_string()).await {
// Le générateur a créé les données, les servir directement
return (
StatusCode::OK,
[("content-type", content_type)],
data,
).into_response();
return (StatusCode::OK, [("content-type", content_type)], data).into_response();
}
}
}
@@ -207,12 +211,7 @@ async fn serve_complete_file(
}
match tokio::fs::read(&file_path).await {
Ok(data) => (
StatusCode::OK,
[("content-type", content_type)],
data,
)
.into_response(),
Ok(data) => (StatusCode::OK, [("content-type", content_type)], data).into_response(),
Err(e) => {
warn!("Error reading file {:?}: {}", file_path, e);
(StatusCode::INTERNAL_SERVER_ERROR, "Error reading file").into_response()
@@ -332,9 +331,7 @@ pub fn create_file_router_with_generator<C: CacheConfig + 'static>(
/// - `DELETE /{pk}` - Supprimer un item
/// - `POST /consolidate` - Consolider le cache
#[cfg(feature = "pmoserver")]
pub fn create_api_router<C: CacheConfig + 'static>(
cache: Arc<Cache<C>>,
) -> Router {
pub fn create_api_router<C: CacheConfig + 'static>(cache: Arc<Cache<C>>) -> Router {
use crate::api;
Router::new()
@@ -346,8 +343,7 @@ pub fn create_api_router<C: CacheConfig + 'static>(
)
.route(
"/{pk}",
get(api::get_item_info::<C>)
.delete(api::delete_item::<C>),
get(api::get_item_info::<C>).delete(api::delete_item::<C>),
)
.route("/{pk}/status", get(api::get_download_status::<C>))
.route("/consolidate", post(api::consolidate_cache::<C>))

View File

@@ -43,12 +43,14 @@ fn create_webp_transformer() -> StreamTransformer {
// Convertir en WebP
let img = image::load_from_memory(&bytes)
.map_err(|e| format!("Image decode error: {}", e))?;
let webp_data = crate::webp::encode_webp(&img)
.map_err(|e| format!("WebP encode error: {}", e))?;
let webp_data =
crate::webp::encode_webp(&img).map_err(|e| format!("WebP encode error: {}", e))?;
// Écrire et mettre à jour la progression
use tokio::io::AsyncWriteExt;
file.write_all(&webp_data).await.map_err(|e| e.to_string())?;
file.write_all(&webp_data)
.await
.map_err(|e| e.to_string())?;
file.flush().await.map_err(|e| e.to_string())?;
progress(webp_data.len() as u64);

View File

@@ -42,15 +42,15 @@ pub mod webp;
#[cfg(feature = "pmoserver")]
pub mod openapi;
pub use cache::{Cache, CoversConfig, new_cache};
pub use cache::{new_cache, Cache, CoversConfig};
#[cfg(feature = "pmoserver")]
pub use openapi::ApiDoc;
#[cfg(feature = "pmoserver")]
use utoipa::OpenApi;
#[cfg(feature = "pmoserver")]
use std::sync::Arc;
#[cfg(feature = "pmoserver")]
use utoipa::OpenApi;
/// Générateur de variantes d'images
///
@@ -64,7 +64,13 @@ fn create_variant_generator() -> pmocache::pmoserver_ext::ParamGenerator<CoversC
match webp::generate_variant(&cache, &pk, size).await {
Ok(data) => return Some(data),
Err(e) => {
tracing::warn!("Cannot generate variant {}x{} for {}: {}", size, size, pk, e);
tracing::warn!(
"Cannot generate variant {}x{} for {}: {}",
size,
size,
pk,
e
);
return None;
}
}
@@ -98,21 +104,26 @@ pub trait CoverCacheExt {
/// - `DELETE /api/covers/{pk}` - Supprimer une image (API REST)
/// - `GET /api/covers/{pk}/status` - Statut du téléchargement
/// - `GET /swagger-ui/covers` - Documentation interactive
async fn init_cover_cache(&mut self, cache_dir: &str, limit: usize)
-> anyhow::Result<Arc<Cache>>;
async fn init_cover_cache(
&mut self,
cache_dir: &str,
limit: usize,
) -> anyhow::Result<Arc<Cache>>;
/// Initialise le cache d'images avec la configuration par défaut
///
/// Utilise automatiquement les paramètres de `pmoconfig::Config`
async fn init_cover_cache_configured(&mut self)
-> anyhow::Result<Arc<Cache>>;
async fn init_cover_cache_configured(&mut self) -> anyhow::Result<Arc<Cache>>;
}
#[cfg(feature = "pmoserver")]
impl CoverCacheExt for pmoserver::Server {
async fn init_cover_cache(&mut self, cache_dir: &str, limit: usize)
-> anyhow::Result<Arc<Cache>> {
use pmocache::pmoserver_ext::{create_file_router_with_generator, create_api_router};
async fn init_cover_cache(
&mut self,
cache_dir: &str,
limit: usize,
) -> anyhow::Result<Arc<Cache>> {
use pmocache::pmoserver_ext::{create_api_router, create_file_router_with_generator};
let cache = Arc::new(cache::new_cache(cache_dir, limit)?);
@@ -121,7 +132,7 @@ impl CoverCacheExt for pmoserver::Server {
let file_router = create_file_router_with_generator(
cache.clone(),
"image/webp",
Some(create_variant_generator())
Some(create_variant_generator()),
);
self.add_router("/", file_router).await;
@@ -134,8 +145,7 @@ impl CoverCacheExt for pmoserver::Server {
Ok(cache)
}
async fn init_cover_cache_configured(&mut self)
-> anyhow::Result<Arc<Cache>> {
async fn init_cover_cache_configured(&mut self) -> anyhow::Result<Arc<Cache>> {
let config = pmoconfig::get_config();
let cache_dir = config.get_cover_cache_dir()?;
let limit = config.get_cover_cache_size()?;

View File

@@ -1,5 +1,5 @@
use anyhow::Result;
use image::{DynamicImage, imageops::FilterType};
use image::{imageops::FilterType, DynamicImage};
use webp::{Encoder, WebPMemory};
pub fn encode_webp(img: &DynamicImage) -> Result<Vec<u8>> {
@@ -11,34 +11,38 @@ pub fn encode_webp(img: &DynamicImage) -> Result<Vec<u8>> {
pub fn ensure_square(img: &DynamicImage, size: u32) -> DynamicImage {
let (width, height) = (img.width(), img.height());
// Calculer le ratio de mise à l'échelle
let scale = if width > height {
size as f32 / width as f32
} else {
size as f32 / height as f32
};
let new_width = (width as f32 * scale) as u32;
let new_height = (height as f32 * scale) as u32;
// Redimensionner l'image
let resized = img.resize(new_width, new_height, FilterType::Lanczos3);
// Créer une image carrée avec fond transparent
let mut square = DynamicImage::new_rgba8(size, size);
// Calculer la position pour centrer l'image redimensionnée
let x = (size - new_width) / 2;
let y = (size - new_height) / 2;
// Copier l'image redimensionnée au centre du carré
image::imageops::overlay(&mut square, &resized, x.into(), y.into());
square
}
pub async fn generate_variant(cache: &super::cache::Cache, pk: &str, size: usize) -> Result<Vec<u8>> {
pub async fn generate_variant(
cache: &super::cache::Cache,
pk: &str,
size: usize,
) -> Result<Vec<u8>> {
// Utiliser file_path_with_qualifier pour obtenir le chemin
let variant_path = cache.file_path_with_qualifier(pk, &size.to_string());
@@ -49,10 +53,7 @@ pub async fn generate_variant(cache: &super::cache::Cache, pk: &str, size: usize
let orig_path = cache.file_path_with_qualifier(pk, "orig");
// Charger l'image de manière synchrone (image::open n'est pas async)
let img = tokio::task::spawn_blocking(move || {
image::open(orig_path)
})
.await??;
let img = tokio::task::spawn_blocking(move || image::open(orig_path)).await??;
let square = ensure_square(&img, size as u32);
let webp_data = encode_webp(&square)?;

View File

@@ -2,19 +2,19 @@
//!
//! Parser et utilitaires pour le format DIDL-Lite utilisé dans UPnP/DLNA.
use bevy_reflect::Reflect;
use serde::{Deserialize, Serialize};
use std::fmt::Write;
use bevy_reflect::Reflect;
// ============= Couche d'abstraction générique =============
/// Trait pour tout parser de métadonnées média
pub trait MediaMetadataParser: Sized {
type Error: std::error::Error + Send + Sync + 'static;
/// Parse une chaîne de métadonnées
fn parse(input: &str) -> Result<Self, Self::Error>;
/// Retourne le format du parser
fn format_name() -> &'static str;
}
@@ -24,10 +24,10 @@ pub trait MediaMetadataParser: Sized {
pub struct ParsedMetadata<T> {
/// Format du document (ex: "DIDL-Lite", "RSS", etc.)
pub format: String,
/// Données parsées
pub data: T,
/// Timestamp du parsing (exclu de la réflexion car SystemTime n'implémente pas Reflect)
#[reflect(ignore)]
#[serde(skip_serializing_if = "Option::is_none")]
@@ -42,7 +42,7 @@ impl<T> ParsedMetadata<T> {
parsed_at: Some(std::time::SystemTime::now()),
}
}
/// Transforme les données avec une fonction
pub fn map<U, F>(self, f: F) -> ParsedMetadata<U>
where
@@ -66,11 +66,11 @@ pub fn parse_metadata<P: MediaMetadataParser>(input: &str) -> Result<ParsedMetad
impl MediaMetadataParser for DIDLLite {
type Error = quick_xml::de::DeError;
fn parse(input: &str) -> Result<Self, Self::Error> {
quick_xml::de::from_str(input)
}
fn format_name() -> &'static str {
"DIDL-Lite"
}
@@ -81,32 +81,31 @@ pub type DidlMetadata = ParsedMetadata<DIDLLite>;
// ============= Structures DIDL-Lite =============
/// Racine d'un document DIDL-Lite
#[derive(Debug, Clone, Serialize, Deserialize, utoipa::ToSchema, Reflect)]
#[serde(rename = "DIDL-Lite")]
pub struct DIDLLite {
#[serde(rename = "@xmlns")]
pub xmlns: String,
#[serde(rename = "@xmlns:upnp", skip_serializing_if = "Option::is_none")]
pub xmlns_upnp: Option<String>,
#[serde(rename = "@xmlns:dc", skip_serializing_if = "Option::is_none")]
pub xmlns_dc: Option<String>,
#[serde(rename = "@xmlns:dlna", skip_serializing_if = "Option::is_none")]
pub xmlns_dlna: Option<String>,
#[serde(rename = "@xmlns:sec", skip_serializing_if = "Option::is_none")]
pub xmlns_sec: Option<String>,
#[serde(rename = "@xmlns:pv", skip_serializing_if = "Option::is_none")]
pub xmlns_pv: Option<String>,
#[serde(rename = "container", default)]
pub containers: Vec<Container>,
#[serde(rename = "item", default)]
pub items: Vec<Item>,
}
@@ -116,25 +115,25 @@ pub struct DIDLLite {
pub struct Container {
#[serde(rename = "@id")]
pub id: String,
#[serde(rename = "@parentID")]
pub parent_id: String,
#[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")]
pub restricted: Option<String>,
#[serde(rename = "@childCount", skip_serializing_if = "Option::is_none")]
pub child_count: Option<String>,
#[serde(rename = "dc:title", alias = "title")]
pub title: String,
#[serde(rename = "upnp:class", alias = "class")]
pub class: String,
#[serde(rename = "container", default)]
pub containers: Vec<Container>,
#[serde(rename = "item", default)]
pub items: Vec<Item>,
}
@@ -144,46 +143,74 @@ pub struct Container {
pub struct Item {
#[serde(rename = "@id")]
pub id: String,
#[serde(rename = "@parentID")]
pub parent_id: String,
#[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")]
pub restricted: Option<String>,
#[serde(rename = "dc:title", alias = "title")]
pub title: String,
#[serde(rename = "dc:creator", alias = "creator", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "dc:creator",
alias = "creator",
skip_serializing_if = "Option::is_none"
)]
pub creator: Option<String>,
#[serde(rename = "upnp:class", alias = "class")]
pub class: String,
#[serde(rename = "upnp:artist", alias = "artist", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "upnp:artist",
alias = "artist",
skip_serializing_if = "Option::is_none"
)]
pub artist: Option<String>,
#[serde(rename = "upnp:album", alias = "album", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "upnp:album",
alias = "album",
skip_serializing_if = "Option::is_none"
)]
pub album: Option<String>,
#[serde(rename = "upnp:genre", alias = "genre", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "upnp:genre",
alias = "genre",
skip_serializing_if = "Option::is_none"
)]
pub genre: Option<String>,
#[serde(rename = "upnp:albumArtURI", alias = "albumArtURI", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "upnp:albumArtURI",
alias = "albumArtURI",
skip_serializing_if = "Option::is_none"
)]
pub album_art: Option<String>,
#[serde(skip)]
pub album_art_pk: Option<String>,
#[serde(rename = "dc:date", alias = "date", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "dc:date",
alias = "date",
skip_serializing_if = "Option::is_none"
)]
pub date: Option<String>,
#[serde(rename = "upnp:originalTrackNumber", alias = "originalTrackNumber", skip_serializing_if = "Option::is_none")]
#[serde(
rename = "upnp:originalTrackNumber",
alias = "originalTrackNumber",
skip_serializing_if = "Option::is_none"
)]
pub original_track_number: Option<String>,
#[serde(rename = "res", default)]
pub resources: Vec<Resource>,
#[serde(rename = "desc", default)]
pub descriptions: Vec<Description>,
}
@@ -193,19 +220,19 @@ pub struct Item {
pub struct Resource {
#[serde(rename = "@protocolInfo")]
pub protocol_info: String,
#[serde(rename = "@bitsPerSample", skip_serializing_if = "Option::is_none")]
pub bits_per_sample: Option<String>,
#[serde(rename = "@sampleFrequency", skip_serializing_if = "Option::is_none")]
pub sample_frequency: Option<String>,
#[serde(rename = "@nrAudioChannels", skip_serializing_if = "Option::is_none")]
pub nr_audio_channels: Option<String>,
#[serde(rename = "@duration", skip_serializing_if = "Option::is_none")]
pub duration: Option<String>,
#[serde(rename = "$text")]
pub url: String,
}
@@ -215,13 +242,13 @@ pub struct Resource {
pub struct Description {
#[serde(rename = "@id", skip_serializing_if = "Option::is_none")]
pub id: Option<String>,
#[serde(rename = "@nameSpace", skip_serializing_if = "Option::is_none")]
pub namespace: Option<String>,
#[serde(rename = "track_gain", skip_serializing_if = "Option::is_none")]
pub track_gain: Option<String>,
#[serde(rename = "track_peak", skip_serializing_if = "Option::is_none")]
pub track_peak: Option<String>,
}
@@ -233,22 +260,22 @@ impl DIDLLite {
pub fn all_containers(&self) -> impl Iterator<Item = &Container> {
AllContainersIter::new(&self.containers)
}
/// Itère sur tous les items de manière récursive
pub fn all_items(&self) -> impl Iterator<Item = &Item> {
AllItemsIter::new(&self.containers, &self.items)
}
/// Trouve un container par ID
pub fn get_container_by_id(&self, id: &str) -> Option<&Container> {
self.all_containers().find(|c| c.id == id)
}
/// Trouve un item par ID
pub fn get_item_by_id(&self, id: &str) -> Option<&Item> {
self.all_items().find(|i| i.id == id)
}
/// Filtre les containers
pub fn filter_containers<F>(&self, predicate: F) -> impl Iterator<Item = &Container>
where
@@ -256,7 +283,7 @@ impl DIDLLite {
{
self.all_containers().filter(move |c| predicate(c))
}
/// Filtre les items
pub fn filter_items<F>(&self, predicate: F) -> impl Iterator<Item = &Item>
where
@@ -264,26 +291,26 @@ impl DIDLLite {
{
self.all_items().filter(move |i| predicate(i))
}
/// Génère une représentation Markdown
pub fn to_markdown(&self) -> String {
let mut buf = String::new();
buf.push_str("### DIDL-Lite Document\n\n");
if !self.containers.is_empty() {
buf.push_str("#### Containers\n\n");
for container in &self.containers {
container.write_markdown(&mut buf, 0);
}
}
if !self.items.is_empty() {
buf.push_str("#### Items\n\n");
for item in &self.items {
item.write_markdown(&mut buf, 0);
}
}
buf
}
}
@@ -293,41 +320,41 @@ impl Container {
pub fn all_containers(&self) -> impl Iterator<Item = &Container> {
AllContainersIter::new(&self.containers)
}
/// Itère sur tous les items de ce container et ses enfants
pub fn all_items(&self) -> impl Iterator<Item = &Item> {
AllItemsIter::new(&self.containers, &self.items)
}
fn write_markdown(&self, buf: &mut String, depth: usize) {
let indent = " ".repeat(depth);
writeln!(buf, "{}- **Container**: {}", indent, self.title).unwrap();
writeln!(buf, "{} - ID: `{}`", indent, self.id).unwrap();
writeln!(buf, "{} - ParentID: `{}`", indent, self.parent_id).unwrap();
writeln!(buf, "{} - Class: `{}`", indent, self.class).unwrap();
if let Some(ref restricted) = self.restricted {
writeln!(buf, "{} - Restricted: `{}`", indent, restricted).unwrap();
}
if let Some(ref count) = self.child_count {
writeln!(buf, "{} - ChildCount: `{}`", indent, count).unwrap();
}
if !self.containers.is_empty() {
writeln!(buf, "{} - Subcontainers:", indent).unwrap();
for sub in &self.containers {
sub.write_markdown(buf, depth + 2);
}
}
if !self.items.is_empty() {
writeln!(buf, "{} - Items:", indent).unwrap();
for item in &self.items {
item.write_markdown(buf, depth + 2);
}
}
buf.push('\n');
}
}
@@ -335,21 +362,22 @@ impl Container {
impl Item {
/// Itère sur les ressources audio uniquement
pub fn audio_resources(&self) -> impl Iterator<Item = &Resource> {
self.resources.iter()
self.resources
.iter()
.filter(|r| r.protocol_info.contains("audio/"))
}
/// Retourne la ressource principale (première disponible)
pub fn primary_resource(&self) -> Option<&Resource> {
self.resources.first()
}
/// Itère sur les métadonnées sous forme de paires clé-valeur
pub fn metadata(&self) -> impl Iterator<Item = (&str, &str)> {
let mut pairs = Vec::new();
pairs.push(("title", self.title.as_str()));
if let Some(ref artist) = self.artist {
pairs.push(("artist", artist.as_str()));
}
@@ -365,7 +393,7 @@ impl Item {
if let Some(ref track) = self.original_track_number {
pairs.push(("trackNumber", track.as_str()));
}
for desc in &self.descriptions {
if let Some(ref gain) = desc.track_gain {
pairs.push(("replayGain", gain.as_str()));
@@ -374,18 +402,18 @@ impl Item {
pairs.push(("replayPeak", peak.as_str()));
}
}
pairs.into_iter()
}
fn write_markdown(&self, buf: &mut String, depth: usize) {
let indent = " ".repeat(depth);
writeln!(buf, "{}- **Item**: {}", indent, self.title).unwrap();
writeln!(buf, "{} - ID: `{}`", indent, self.id).unwrap();
writeln!(buf, "{} - ParentID: `{}`", indent, self.parent_id).unwrap();
writeln!(buf, "{} - Class: `{}`", indent, self.class).unwrap();
if let Some(ref creator) = self.creator {
writeln!(buf, "{} - Creator: {}", indent, creator).unwrap();
}
@@ -407,7 +435,7 @@ impl Item {
if let Some(ref track) = self.original_track_number {
writeln!(buf, "{} - Track: {}", indent, track).unwrap();
}
if !self.resources.is_empty() {
writeln!(buf, "{} - Resources:", indent).unwrap();
for res in &self.resources {
@@ -427,7 +455,7 @@ impl Item {
}
}
}
if !self.descriptions.is_empty() {
writeln!(buf, "{} - Descriptions:", indent).unwrap();
for desc in &self.descriptions {
@@ -442,7 +470,7 @@ impl Item {
}
}
}
buf.push('\n');
}
}
@@ -463,7 +491,7 @@ impl<'a> AllContainersIter<'a> {
impl<'a> Iterator for AllContainersIter<'a> {
type Item = &'a Container;
fn next(&mut self) -> Option<Self::Item> {
self.stack.pop().map(|container| {
// Ajouter les enfants à la pile
@@ -489,13 +517,13 @@ impl<'a> AllItemsIter<'a> {
impl<'a> Iterator for AllItemsIter<'a> {
type Item = &'a Item;
fn next(&mut self) -> Option<Self::Item> {
loop {
if let Some(item) = self.current_items.next() {
return Some(item);
}
let container = self.containers.pop()?;
self.containers.extend(container.containers.iter());
self.current_items = container.items.iter();
@@ -506,7 +534,7 @@ impl<'a> Iterator for AllItemsIter<'a> {
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_simple_didl() {
let xml = r#"
@@ -520,12 +548,12 @@ mod tests {
</item>
</DIDL-Lite>
"#;
let didl = DIDLLite::parse(xml).unwrap();
assert_eq!(didl.items.len(), 1);
assert_eq!(didl.items[0].title, "Test Song");
}
#[test]
fn test_parse_without_namespaces() {
// Teste un XML sans namespaces explicites (devices UPnP laxistes)
@@ -538,12 +566,12 @@ mod tests {
</item>
</DIDL-Lite>
"#;
let didl = DIDLLite::parse(xml).unwrap();
assert_eq!(didl.items.len(), 1);
assert_eq!(didl.items[0].title, "Test Song");
}
#[test]
fn test_generic_parser() {
let xml = r#"
@@ -552,14 +580,14 @@ mod tests {
xmlns:upnp="urn:schemas-upnp-org:metadata-1-0/upnp/">
</DIDL-Lite>
"#;
// Utiliser le parser générique
let metadata: DidlMetadata = parse_metadata(xml).unwrap();
assert_eq!(metadata.format, "DIDL-Lite");
assert!(metadata.parsed_at.is_some());
}
#[test]
fn test_metadata_map() {
let xml = r#"
@@ -568,13 +596,13 @@ mod tests {
xmlns:upnp="urn:schemas-upnp-org:metadata-1-0/upnp/">
</DIDL-Lite>
"#;
let metadata: DidlMetadata = parse_metadata(xml).unwrap();
// Transformer les données
let item_count = metadata.map(|didl| didl.items.len());
assert_eq!(item_count.format, "DIDL-Lite");
assert_eq!(item_count.data, 0);
}
}
}

View File

@@ -1,4 +1,7 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, NUMBEROFTRACKS, CURRENTTRACK, AVTRANSPORTURI, AVTRANSPORTURIMETADATA, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA, AVTRANSPORTURI,
AVTRANSPORTURIMETADATA, CURRENTTRACK, NUMBEROFTRACKS,
};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,4 +1,7 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, CURRENTTRACK, CURRENTTRACKDURATION, AVTRANSPORTURI, AVTRANSPORTURIMETADATA, RELATIVETIMEPOSITION, ABSOLUTETIMEPOSITION};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, ABSOLUTETIMEPOSITION, AVTRANSPORTURI, AVTRANSPORTURIMETADATA,
CURRENTTRACK, CURRENTTRACKDURATION, RELATIVETIMEPOSITION,
};
use pmoupnp::define_action;
define_action! {

View File

@@ -27,4 +27,3 @@ pub use seek::SEEK;
pub use setavtransportnexturi::SETNEXTAVTRANSPORTURI;
pub use setavtransporturi::SETAVTRANSPORTURI;
pub use stop::STOP;

View File

@@ -1,6 +1,6 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, TRANSPORTPLAYSPEED};
use pmoupnp::define_action;
define_action! {
pub static PLAY = "Play" {
in "InstanceID" => A_ARG_TYPE_INSTANCE_ID,

View File

@@ -1,4 +1,6 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_SEEKMODE, CURRENTTRACKDURATION};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_SEEKMODE, CURRENTTRACKDURATION,
};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,4 +1,6 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA,
};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,4 +1,6 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, AVTRANSPORTURI, AVTRANSPORTURIMETADATA};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, AVTRANSPORTURI, AVTRANSPORTURIMETADATA,
};
use pmoupnp::define_action;
define_action! {

View File

@@ -91,22 +91,21 @@
use pmoupnp::define_service;
pub mod variables;
pub mod actions;
pub mod variables;
use actions::{
GETCURRENTTRANSPORTACTIONS, GETDEVICECAPABILITIES, GETMEDIAINFO,
GETPOSITIONINFO, GETTRANSPORTINFO, GETTRANSPORTSETTINGS, NEXT, PAUSE,
PLAY, PREVIOUS, SEEK, SETNEXTAVTRANSPORTURI, SETAVTRANSPORTURI, STOP
GETCURRENTTRANSPORTACTIONS, GETDEVICECAPABILITIES, GETMEDIAINFO, GETPOSITIONINFO,
GETTRANSPORTINFO, GETTRANSPORTSETTINGS, NEXT, PAUSE, PLAY, PREVIOUS, SEEK, SETAVTRANSPORTURI,
SETNEXTAVTRANSPORTURI, STOP,
};
use variables::{
ABSOLUTETIMEPOSITION, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA,
AVTRANSPORTURI, AVTRANSPORTURIMETADATA, A_ARG_TYPE_INSTANCE_ID,
A_ARG_TYPE_PLAY_SPEED, A_ARG_TYPE_SEEKMODE, CURRENTMEDIADURATION,
CURRENTPLAYMODE, CURRENTTRACK, CURRENTTRACKDURATION, CURRENTTRACKMETADATA,
CURRENTTRACKURI, NUMBEROFTRACKS, PLAYBACKSTORAGEMEDIUM,
POSSIBLEPLAYBACKSTORAGEMEDIA, RELATIVETIMEPOSITION, SEEKMODE,
TRANSPORTPLAYSPEED, TRANSPORTSTATE, TRANSPORTSTATUS
A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_PLAY_SPEED, A_ARG_TYPE_SEEKMODE, ABSOLUTETIMEPOSITION,
AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA, AVTRANSPORTURI, AVTRANSPORTURIMETADATA,
CURRENTMEDIADURATION, CURRENTPLAYMODE, CURRENTTRACK, CURRENTTRACKDURATION,
CURRENTTRACKMETADATA, CURRENTTRACKURI, NUMBEROFTRACKS, PLAYBACKSTORAGEMEDIUM,
POSSIBLEPLAYBACKSTORAGEMEDIA, RELATIVETIMEPOSITION, SEEKMODE, TRANSPORTPLAYSPEED,
TRANSPORTSTATE, TRANSPORTSTATUS,
};
// Service AVTransport:1 conforme à la spécification UPnP AV pour MediaRenderer audio
@@ -154,4 +153,4 @@ define_service! {
STOP,
]
}
}
}

View File

@@ -1,31 +1,37 @@
use std::sync::Arc;
use pmoupnp::state_variables::{StateVariable, StateVariableError};
use pmoupnp::variable_types::StateVarType;
use bevy_reflect::Reflect;
use once_cell::sync::Lazy;
use pmodidl::{DIDLLite, MediaMetadataParser};
use pmoupnp::state_variables::{StateVariable, StateVariableError};
use pmoupnp::variable_types::StateVarType;
fn avtransporturimetadataparser(value: &str) -> Result<Box<dyn Reflect>, StateVariableError> {
// Parse DIDL-Lite
let didl = DIDLLite::parse(value)
.map_err(|e| StateVariableError::ParseError(format!("Failed to parse DIDL-Lite: {}", e)))?;
// Retourne le résultat sous forme de Box<dyn Reflect>
Ok(Box::new(didl) as Box<dyn Reflect>)
}
pub static AVTRANSPORTURIMETADATA: Lazy<Arc<StateVariable>> = Lazy::new(|| -> Arc<StateVariable> {
let mut sv = StateVariable::new(StateVarType::String, "AVTransportURIMetaData".to_string());
pub static AVTRANSPORTURIMETADATA: Lazy<Arc<StateVariable>> =
Lazy::new(|| -> Arc<StateVariable> {
let mut sv = StateVariable::new(StateVarType::String, "AVTransportURIMetaData".to_string());
sv.set_value_parser(Arc::new(avtransporturimetadataparser)).expect("Failed to set parser");
Arc::new(sv)
});
sv.set_value_parser(Arc::new(avtransporturimetadataparser))
.expect("Failed to set parser");
Arc::new(sv)
});
pub static AVTRANSPORTNEXTURIMETADATA: Lazy<Arc<StateVariable>> = Lazy::new(|| -> Arc<StateVariable> {
let mut sv = StateVariable::new(StateVarType::String, "AVTransportNextURIMetaData".to_string());
pub static AVTRANSPORTNEXTURIMETADATA: Lazy<Arc<StateVariable>> =
Lazy::new(|| -> Arc<StateVariable> {
let mut sv = StateVariable::new(
StateVarType::String,
"AVTransportNextURIMetaData".to_string(),
);
sv.set_value_parser(Arc::new(avtransporturimetadataparser)).expect("Failed to set parser");
Arc::new(sv)
});
sv.set_value_parser(Arc::new(avtransporturimetadataparser))
.expect("Failed to set parser");
Arc::new(sv)
});

View File

@@ -21,10 +21,10 @@ mod transportstatus;
pub use a_arg_type_instanceid::A_ARG_TYPE_INSTANCE_ID;
pub use a_arg_type_playspeed::A_ARG_TYPE_PLAY_SPEED;
pub use a_arg_type_seekmode::A_ARG_TYPE_SEEKMODE;
pub use avtransporturi::AVTRANSPORTURI;
pub use avtransporturi::AVTRANSPORTNEXTURI;
pub use avtransporturimetadata::AVTRANSPORTURIMETADATA;
pub use avtransporturi::AVTRANSPORTURI;
pub use avtransporturimetadata::AVTRANSPORTNEXTURIMETADATA;
pub use avtransporturimetadata::AVTRANSPORTURIMETADATA;
pub use currentmediaduration::CURRENTMEDIADURATION;
pub use currentplaymode::CURRENTPLAYMODE;
pub use currenttrackmetadata::CURRENTTRACKMETADATA;
@@ -42,6 +42,3 @@ pub use trackduration::RELATIVETIMEPOSITION;
pub use transportplayspeed::TRANSPORTPLAYSPEED;
pub use transportstate::TRANSPORTSTATE;
pub use transportstatus::TRANSPORTSTATUS;

View File

@@ -1,10 +1,14 @@
use std::sync::Arc;
use once_cell::sync::Lazy;
use pmoupnp::state_variables::StateVariable;
use pmoupnp::variable_types::StateVarType;
use once_cell::sync::Lazy;
pub static POSSIBLERECORDSTORAGEMEDIA: Lazy<Arc<StateVariable>> = Lazy::new(|| -> Arc<StateVariable> {
let sv = StateVariable::new(StateVarType::String, "PossibleRecordStorageMedia".to_string());
Arc::new(sv)
});
pub static POSSIBLERECORDSTORAGEMEDIA: Lazy<Arc<StateVariable>> =
Lazy::new(|| -> Arc<StateVariable> {
let sv = StateVariable::new(
StateVarType::String,
"PossibleRecordStorageMedia".to_string(),
);
Arc::new(sv)
});

View File

@@ -1,8 +1,8 @@
use std::sync::Arc;
use once_cell::sync::Lazy;
use pmoupnp::state_variables::StateVariable;
use pmoupnp::variable_types::StateVarType;
use once_cell::sync::Lazy;
pub static RECORDSTORAGEMEDIUM: Lazy<Arc<StateVariable>> = Lazy::new(|| -> Arc<StateVariable> {
let sv = StateVariable::new(StateVarType::String, "RecordStorageMedium".to_string());

View File

@@ -3,4 +3,3 @@ use pmoupnp::define_variable;
define_variable! {
pub static SEEKMODE: String = "SeekMode"
}

View File

@@ -11,4 +11,3 @@ define_variable! {
evented: true,
}
}

View File

@@ -17,4 +17,3 @@ define_variable! {
evented: true,
}
}

View File

@@ -6,4 +6,3 @@ define_variable! {
default: "1",
}
}

View File

@@ -7,4 +7,3 @@ define_variable! {
evented: true,
}
}

View File

@@ -1,6 +1,6 @@
use crate::connectionmanager::variables::{
A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_RCSID, A_ARG_TYPE_AVTRANSPORTID,
A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_DIRECTION, A_ARG_TYPE_CONNECTIONSTATUS
A_ARG_TYPE_AVTRANSPORTID, A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID,
};
use pmoupnp::define_action;

View File

@@ -1,4 +1,4 @@
use crate::connectionmanager::variables::{SOURCEPROTOCOLINFO, SINKPROTOCOLINFO};
use crate::connectionmanager::variables::{SINKPROTOCOLINFO, SOURCEPROTOCOLINFO};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,7 +1,7 @@
mod getprotocolinfo;
mod getcurrentconnectionids;
mod getcurrentconnectioninfo;
mod getprotocolinfo;
pub use getprotocolinfo::GETPROTOCOLINFO;
pub use getcurrentconnectionids::GETCURRENTCONNECTIONIDS;
pub use getcurrentconnectioninfo::GETCURRENTCONNECTIONINFO;
pub use getprotocolinfo::GETPROTOCOLINFO;

View File

@@ -55,14 +55,14 @@
use pmoupnp::define_service;
pub mod variables;
pub mod actions;
pub mod variables;
use actions::{GETCURRENTCONNECTIONIDS, GETCURRENTCONNECTIONINFO, GETPROTOCOLINFO};
use variables::{
A_ARG_TYPE_AVTRANSPORTID, A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID,
CURRENTCONNECTIONIDS, SINKPROTOCOLINFO, SOURCEPROTOCOLINFO
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID, CURRENTCONNECTIONIDS,
SINKPROTOCOLINFO, SOURCEPROTOCOLINFO,
};
// Service ConnectionManager:1 conforme à la spécification UPnP AV pour MediaRenderer audio

View File

@@ -1,19 +1,19 @@
mod sourceprotocolinfo;
mod sinkprotocolinfo;
mod currentconnectionids;
mod a_arg_type_avtransportid;
mod a_arg_type_connectionid;
mod a_arg_type_connectionstatus;
mod a_arg_type_direction;
mod a_arg_type_protocolinfo;
mod a_arg_type_rcsid;
mod a_arg_type_avtransportid;
mod currentconnectionids;
mod sinkprotocolinfo;
mod sourceprotocolinfo;
pub use sourceprotocolinfo::SOURCEPROTOCOLINFO;
pub use sinkprotocolinfo::SINKPROTOCOLINFO;
pub use currentconnectionids::CURRENTCONNECTIONIDS;
pub use a_arg_type_avtransportid::A_ARG_TYPE_AVTRANSPORTID;
pub use a_arg_type_connectionid::A_ARG_TYPE_CONNECTIONID;
pub use a_arg_type_connectionstatus::A_ARG_TYPE_CONNECTIONSTATUS;
pub use a_arg_type_direction::A_ARG_TYPE_DIRECTION;
pub use a_arg_type_protocolinfo::A_ARG_TYPE_PROTOCOLINFO;
pub use a_arg_type_rcsid::A_ARG_TYPE_RCSID;
pub use a_arg_type_avtransportid::A_ARG_TYPE_AVTRANSPORTID;
pub use currentconnectionids::CURRENTCONNECTIONIDS;
pub use sinkprotocolinfo::SINKPROTOCOLINFO;
pub use sourceprotocolinfo::SOURCEPROTOCOLINFO;

View File

@@ -3,12 +3,11 @@
use once_cell::sync::Lazy;
use std::sync::Arc;
use pmoupnp::devices::Device;
use crate::{
avtransport::AVTTRANSPORT,
avtransport::AVTTRANSPORT, connectionmanager::CONNECTIONMANAGER,
renderingcontrol::RENDERINGCONTROL,
connectionmanager::CONNECTIONMANAGER,
};
use pmoupnp::devices::Device;
/// Device MediaRenderer UPnP.
///
@@ -54,13 +53,16 @@ pub static MEDIA_RENDERER: Lazy<Arc<Device>> = Lazy::new(|| {
device.set_udn_prefix("pmomusic".to_string());
// Ajouter les trois services obligatoires
device.add_service(Arc::clone(&AVTTRANSPORT))
device
.add_service(Arc::clone(&AVTTRANSPORT))
.expect("Failed to add AVTransport service");
device.add_service(Arc::clone(&RENDERINGCONTROL))
device
.add_service(Arc::clone(&RENDERINGCONTROL))
.expect("Failed to add RenderingControl service");
device.add_service(Arc::clone(&CONNECTIONMANAGER))
device
.add_service(Arc::clone(&CONNECTIONMANAGER))
.expect("Failed to add ConnectionManager service");
Arc::new(device)

View File

@@ -29,7 +29,7 @@
pub mod avtransport;
pub mod connectionmanager;
pub mod renderingcontrol;
pub mod device;
pub mod renderingcontrol;
pub use device::MEDIA_RENDERER;

View File

@@ -1,4 +1,4 @@
use crate::renderingcontrol::variables::{A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_CHANNEL, MUTE};
use crate::renderingcontrol::variables::{A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID, MUTE};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,4 +1,4 @@
use crate::renderingcontrol::variables::{A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_CHANNEL, VOLUME};
use crate::renderingcontrol::variables::{A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID, VOLUME};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,9 +1,9 @@
mod getvolume;
mod setvolume;
mod getmute;
mod getvolume;
mod setmute;
mod setvolume;
pub use getvolume::GETVOLUME;
pub use setvolume::SETVOLUME;
pub use getmute::GETMUTE;
pub use getvolume::GETVOLUME;
pub use setmute::SETMUTE;
pub use setvolume::SETVOLUME;

View File

@@ -1,4 +1,4 @@
use crate::renderingcontrol::variables::{A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_CHANNEL, MUTE};
use crate::renderingcontrol::variables::{A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID, MUTE};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,4 +1,4 @@
use crate::renderingcontrol::variables::{A_ARG_TYPE_INSTANCE_ID, A_ARG_TYPE_CHANNEL, VOLUME};
use crate::renderingcontrol::variables::{A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID, VOLUME};
use pmoupnp::define_action;
define_action! {

View File

@@ -51,8 +51,8 @@
use pmoupnp::define_service;
pub mod variables;
pub mod actions;
pub mod variables;
use actions::{GETMUTE, GETVOLUME, SETMUTE, SETVOLUME};
use variables::{A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID, MUTE, VOLUME};

View File

@@ -1,9 +1,9 @@
mod a_arg_type_instanceid;
mod a_arg_type_channel;
mod volume;
mod a_arg_type_instanceid;
mod mute;
mod volume;
pub use a_arg_type_instanceid::A_ARG_TYPE_INSTANCE_ID;
pub use a_arg_type_channel::A_ARG_TYPE_CHANNEL;
pub use volume::VOLUME;
pub use a_arg_type_instanceid::A_ARG_TYPE_INSTANCE_ID;
pub use mute::MUTE;
pub use volume::VOLUME;

View File

@@ -1,6 +1,6 @@
use crate::connectionmanager::variables::{
A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_RCSID, A_ARG_TYPE_AVTRANSPORTID,
A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_DIRECTION, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_AVTRANSPORTID, A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID,
};
use pmoupnp::define_action;

View File

@@ -1,4 +1,4 @@
use crate::connectionmanager::variables::{SOURCEPROTOCOLINFO, SINKPROTOCOLINFO};
use crate::connectionmanager::variables::{SINKPROTOCOLINFO, SOURCEPROTOCOLINFO};
use pmoupnp::define_action;
define_action! {

View File

@@ -1,7 +1,7 @@
mod getprotocolinfo;
mod getcurrentconnectionids;
mod getcurrentconnectioninfo;
mod getprotocolinfo;
pub use getprotocolinfo::GETPROTOCOLINFO;
pub use getcurrentconnectionids::GETCURRENTCONNECTIONIDS;
pub use getcurrentconnectioninfo::GETCURRENTCONNECTIONINFO;
pub use getprotocolinfo::GETPROTOCOLINFO;

View File

@@ -65,14 +65,14 @@
use pmoupnp::define_service;
pub mod variables;
pub mod actions;
pub mod variables;
use actions::{GETCURRENTCONNECTIONIDS, GETCURRENTCONNECTIONINFO, GETPROTOCOLINFO};
use variables::{
A_ARG_TYPE_AVTRANSPORTID, A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID,
CURRENTCONNECTIONIDS, SINKPROTOCOLINFO, SOURCEPROTOCOLINFO
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID, CURRENTCONNECTIONIDS,
SINKPROTOCOLINFO, SOURCEPROTOCOLINFO,
};
// Service ConnectionManager:1 conforme à la spécification UPnP AV pour MediaServer

View File

@@ -1,19 +1,19 @@
mod a_arg_type_avtransportid;
mod a_arg_type_connectionid;
mod a_arg_type_connectionstatus;
mod a_arg_type_direction;
mod a_arg_type_protocolinfo;
mod a_arg_type_rcsid;
mod a_arg_type_avtransportid;
mod currentconnectionids;
mod sourceprotocolinfo;
mod sinkprotocolinfo;
mod sourceprotocolinfo;
pub use a_arg_type_avtransportid::A_ARG_TYPE_AVTRANSPORTID;
pub use a_arg_type_connectionid::A_ARG_TYPE_CONNECTIONID;
pub use a_arg_type_connectionstatus::A_ARG_TYPE_CONNECTIONSTATUS;
pub use a_arg_type_direction::A_ARG_TYPE_DIRECTION;
pub use a_arg_type_protocolinfo::A_ARG_TYPE_PROTOCOLINFO;
pub use a_arg_type_rcsid::A_ARG_TYPE_RCSID;
pub use a_arg_type_avtransportid::A_ARG_TYPE_AVTRANSPORTID;
pub use currentconnectionids::CURRENTCONNECTIONIDS;
pub use sourceprotocolinfo::SOURCEPROTOCOLINFO;
pub use sinkprotocolinfo::SINKPROTOCOLINFO;
pub use sourceprotocolinfo::SOURCEPROTOCOLINFO;

View File

@@ -10,8 +10,8 @@
//! - **Search** : Recherche dans les sources qui le supportent
//! - **Update ID** : Suivi des changements pour les notifications UPnP
use pmosource::api::{list_all_sources, get_source as get_source_from_registry};
use pmodidl::{Container, DIDLLite};
use pmosource::api::{get_source as get_source_from_registry, list_all_sources};
use pmosource::{BrowseResult, MusicSource};
use std::sync::Arc;
@@ -28,8 +28,7 @@ fn to_didl_lite(containers: &[Container], items: &[pmodidl::Item]) -> Result<Str
items: items.to_vec(),
};
quick_xml::se::to_string(&didl)
.map_err(|e| format!("Failed to serialize DIDL-Lite: {}", e))
quick_xml::se::to_string(&didl).map_err(|e| format!("Failed to serialize DIDL-Lite: {}", e))
}
/// Handler pour le service ContentDirectory
@@ -208,11 +207,7 @@ impl ContentHandler {
requested_count as usize
};
let paginated: Vec<Container> = containers
.into_iter()
.skip(start)
.take(count)
.collect();
let paginated: Vec<Container> = containers.into_iter().skip(start).take(count).collect();
let returned = paginated.len();
let didl = to_didl_lite(&paginated, &[])?;
@@ -278,11 +273,7 @@ impl ContentHandler {
// On commence dans les items
containers.clear();
let item_start = start - total_containers;
items = items
.into_iter()
.skip(item_start)
.take(count)
.collect();
items = items.into_iter().skip(item_start).take(count).collect();
}
let returned = (containers.len() + items.len()) as u32;

View File

@@ -1,10 +1,9 @@
use crate::contentdirectory::handlers;
use crate::contentdirectory::variables::{
A_ARG_TYPE_OBJECTID, A_ARG_TYPE_BROWSEFLAG, A_ARG_TYPE_FILTER,
A_ARG_TYPE_SORTCRITERIA, A_ARG_TYPE_INDEX, A_ARG_TYPE_COUNT,
A_ARG_TYPE_RESULT, A_ARG_TYPE_UPDATEID,
A_ARG_TYPE_BROWSEFLAG, A_ARG_TYPE_COUNT, A_ARG_TYPE_FILTER, A_ARG_TYPE_INDEX,
A_ARG_TYPE_OBJECTID, A_ARG_TYPE_RESULT, A_ARG_TYPE_SORTCRITERIA, A_ARG_TYPE_UPDATEID,
};
use pmoupnp::define_action;
use crate::contentdirectory::handlers;
define_action! {
pub static BROWSE = "Browse" stateless {

View File

@@ -1,6 +1,6 @@
use crate::contentdirectory::handlers;
use crate::contentdirectory::variables::SEARCHCAPABILITIES;
use pmoupnp::define_action;
use crate::contentdirectory::handlers;
define_action! {
pub static GETSEARCHCAPABILITIES = "GetSearchCapabilities" stateless {

View File

@@ -1,6 +1,6 @@
use crate::contentdirectory::handlers;
use crate::contentdirectory::variables::SORTCAPABILITIES;
use pmoupnp::define_action;
use crate::contentdirectory::handlers;
define_action! {
pub static GETSORTCAPABILITIES = "GetSortCapabilities" stateless {

View File

@@ -1,6 +1,6 @@
use crate::contentdirectory::handlers;
use crate::contentdirectory::variables::SYSTEMUPDATEID;
use pmoupnp::define_action;
use crate::contentdirectory::handlers;
define_action! {
pub static GETSYSTEMUPDATEID = "GetSystemUpdateID" stateless {

View File

@@ -1,11 +1,11 @@
mod browse;
mod search;
mod getsearchcapabilities;
mod getsortcapabilities;
mod getsystemupdateid;
mod search;
pub use browse::BROWSE;
pub use search::SEARCH;
pub use getsearchcapabilities::GETSEARCHCAPABILITIES;
pub use getsortcapabilities::GETSORTCAPABILITIES;
pub use getsystemupdateid::GETSYSTEMUPDATEID;
pub use search::SEARCH;

View File

@@ -1,10 +1,9 @@
use crate::contentdirectory::handlers;
use crate::contentdirectory::variables::{
A_ARG_TYPE_OBJECTID, A_ARG_TYPE_SEARCHCRITERIA, A_ARG_TYPE_FILTER,
A_ARG_TYPE_SORTCRITERIA, A_ARG_TYPE_INDEX, A_ARG_TYPE_COUNT,
A_ARG_TYPE_RESULT, A_ARG_TYPE_UPDATEID,
A_ARG_TYPE_COUNT, A_ARG_TYPE_FILTER, A_ARG_TYPE_INDEX, A_ARG_TYPE_OBJECTID, A_ARG_TYPE_RESULT,
A_ARG_TYPE_SEARCHCRITERIA, A_ARG_TYPE_SORTCRITERIA, A_ARG_TYPE_UPDATEID,
};
use pmoupnp::define_action;
use crate::contentdirectory::handlers;
define_action! {
pub static SEARCH = "Search" stateless {

View File

@@ -23,9 +23,9 @@
//! - [`get_sort_capabilities_handler`] : Capacités de tri supportées
//! - [`get_system_update_id_handler`] : ID de mise à jour du système
use pmoupnp::{action_handler, get, set};
use pmoupnp::actions::{ActionError, ActionHandler};
use crate::content_handler::ContentHandler;
use pmoupnp::actions::{ActionError, ActionHandler};
use pmoupnp::{action_handler, get, set};
use tracing::{debug, error};
/// Handler pour l'action Browse.
@@ -76,7 +76,10 @@ pub fn browse_handler() -> ActionHandler {
set!(&mut data, "TotalMatches", total);
set!(&mut data, "UpdateID", update_id);
debug!("✅ Browse completed: returned={}, total={}", returned, total);
debug!(
"✅ Browse completed: returned={}, total={}",
returned, total
);
Ok(data)
})
}
@@ -128,7 +131,10 @@ pub fn search_handler() -> ActionHandler {
set!(&mut data, "TotalMatches", total);
set!(&mut data, "UpdateID", update_id);
debug!("✅ Search completed: returned={}, total={}", returned, total);
debug!(
"✅ Search completed: returned={}, total={}",
returned, total
);
Ok(data)
})
}

View File

@@ -74,18 +74,15 @@
use pmoupnp::define_service;
pub mod variables;
pub mod actions;
pub mod handlers;
pub mod variables;
use actions::{
BROWSE, SEARCH, GETSEARCHCAPABILITIES, GETSORTCAPABILITIES, GETSYSTEMUPDATEID
};
use actions::{BROWSE, GETSEARCHCAPABILITIES, GETSORTCAPABILITIES, GETSYSTEMUPDATEID, SEARCH};
use variables::{
A_ARG_TYPE_OBJECTID, A_ARG_TYPE_BROWSEFLAG, A_ARG_TYPE_FILTER,
A_ARG_TYPE_SORTCRITERIA, A_ARG_TYPE_INDEX, A_ARG_TYPE_COUNT,
A_ARG_TYPE_UPDATEID, A_ARG_TYPE_RESULT, A_ARG_TYPE_SEARCHCRITERIA,
SEARCHCAPABILITIES, SORTCAPABILITIES, SYSTEMUPDATEID
A_ARG_TYPE_BROWSEFLAG, A_ARG_TYPE_COUNT, A_ARG_TYPE_FILTER, A_ARG_TYPE_INDEX,
A_ARG_TYPE_OBJECTID, A_ARG_TYPE_RESULT, A_ARG_TYPE_SEARCHCRITERIA, A_ARG_TYPE_SORTCRITERIA,
A_ARG_TYPE_UPDATEID, SEARCHCAPABILITIES, SORTCAPABILITIES, SYSTEMUPDATEID,
};
// Service ContentDirectory:1 conforme à la spécification UPnP AV pour MediaServer

View File

@@ -1,25 +1,25 @@
mod a_arg_type_objectid;
mod a_arg_type_browseflag;
mod a_arg_type_filter;
mod a_arg_type_sortcriteria;
mod a_arg_type_index;
mod a_arg_type_count;
mod a_arg_type_updateid;
mod a_arg_type_filter;
mod a_arg_type_index;
mod a_arg_type_objectid;
mod a_arg_type_result;
mod a_arg_type_searchcriteria;
mod a_arg_type_sortcriteria;
mod a_arg_type_updateid;
mod searchcapabilities;
mod sortcapabilities;
mod systemupdateid;
pub use a_arg_type_objectid::A_ARG_TYPE_OBJECTID;
pub use a_arg_type_browseflag::A_ARG_TYPE_BROWSEFLAG;
pub use a_arg_type_filter::A_ARG_TYPE_FILTER;
pub use a_arg_type_sortcriteria::A_ARG_TYPE_SORTCRITERIA;
pub use a_arg_type_index::A_ARG_TYPE_INDEX;
pub use a_arg_type_count::A_ARG_TYPE_COUNT;
pub use a_arg_type_updateid::A_ARG_TYPE_UPDATEID;
pub use a_arg_type_filter::A_ARG_TYPE_FILTER;
pub use a_arg_type_index::A_ARG_TYPE_INDEX;
pub use a_arg_type_objectid::A_ARG_TYPE_OBJECTID;
pub use a_arg_type_result::A_ARG_TYPE_RESULT;
pub use a_arg_type_searchcriteria::A_ARG_TYPE_SEARCHCRITERIA;
pub use a_arg_type_sortcriteria::A_ARG_TYPE_SORTCRITERIA;
pub use a_arg_type_updateid::A_ARG_TYPE_UPDATEID;
pub use searchcapabilities::SEARCHCAPABILITIES;
pub use sortcapabilities::SORTCAPABILITIES;
pub use systemupdateid::SYSTEMUPDATEID;

View File

@@ -3,11 +3,8 @@
use once_cell::sync::Lazy;
use std::sync::Arc;
use crate::{connectionmanager::CONNECTIONMANAGER, contentdirectory::CONTENTDIRECTORY};
use pmoupnp::devices::Device;
use crate::{
contentdirectory::CONTENTDIRECTORY,
connectionmanager::CONNECTIONMANAGER,
};
/// Device MediaServer UPnP.
///
@@ -52,10 +49,12 @@ pub static MEDIA_SERVER: Lazy<Arc<Device>> = Lazy::new(|| {
device.set_udn_prefix("pmomusic".to_string());
// Ajouter les deux services obligatoires
device.add_service(Arc::clone(&CONTENTDIRECTORY))
device
.add_service(Arc::clone(&CONTENTDIRECTORY))
.expect("Failed to add ContentDirectory service");
device.add_service(Arc::clone(&CONNECTIONMANAGER))
device
.add_service(Arc::clone(&CONNECTIONMANAGER))
.expect("Failed to add ConnectionManager service");
Arc::new(device)

View File

@@ -63,23 +63,23 @@
//! server.register_qobuz_from_config().await?;
//! ```
pub mod contentdirectory;
pub mod connectionmanager;
pub mod device;
pub mod source_registry;
pub mod server_ext;
pub mod content_handler;
pub mod contentdirectory;
pub mod device;
pub mod server_ext;
pub mod source_registry;
pub mod sources;
// API REST pour l'enregistrement des sources (requires features qobuz/paradise)
#[cfg(any(feature = "qobuz", feature = "paradise"))]
pub mod sources_api;
pub use device::MEDIA_SERVER;
pub use source_registry::SourceRegistry;
pub use server_ext::{MediaServerExt, get_source_registry, MusicSourceExt};
pub use content_handler::ContentHandler;
pub use sources::{SourcesExt, SourceInitError};
pub use device::MEDIA_SERVER;
pub use server_ext::{MediaServerExt, MusicSourceExt, get_source_registry};
pub use source_registry::SourceRegistry;
pub use sources::{SourceInitError, SourcesExt};
// Re-export sources when features are enabled
#[cfg(feature = "qobuz")]

View File

@@ -7,8 +7,8 @@
//! méthodes spécifiques au MediaServer UPnP.
use async_trait::async_trait;
use pmosource::MusicSource;
use pmoserver::Server;
use pmosource::MusicSource;
use std::sync::Arc;
// Réexporter le trait de base de pmosource
@@ -23,7 +23,10 @@ pub use pmosource::MusicSourceExt;
///
/// let sources = pmosource::api::list_all_sources().await;
/// ```
#[deprecated(since = "0.2.0", note = "Use pmosource::api::list_all_sources() directly")]
#[deprecated(
since = "0.2.0",
note = "Use pmosource::api::list_all_sources() directly"
)]
pub async fn get_source_registry() -> Vec<Arc<dyn MusicSource>> {
pmosource::api::list_all_sources().await
}

View File

@@ -210,8 +210,8 @@ impl Default for SourceRegistry {
#[cfg(test)]
mod tests {
use super::*;
use pmosource::{MusicSource, Result, BrowseResult};
use pmodidl::{Container, Item};
use pmosource::{BrowseResult, MusicSource, Result};
use std::time::SystemTime;
#[derive(Debug)]
@@ -308,9 +308,15 @@ mod tests {
async fn test_list_all() {
let registry = SourceRegistry::new();
registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await;
registry.register(Arc::new(TestSource::new("test-2", "Test 2"))).await;
registry.register(Arc::new(TestSource::new("test-3", "Test 3"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "Test 1")))
.await;
registry
.register(Arc::new(TestSource::new("test-2", "Test 2")))
.await;
registry
.register(Arc::new(TestSource::new("test-3", "Test 3")))
.await;
let sources = registry.list_all().await;
assert_eq!(sources.len(), 3);
@@ -321,17 +327,23 @@ mod tests {
let registry = SourceRegistry::new();
assert_eq!(registry.count().await, 0);
registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "Test 1")))
.await;
assert_eq!(registry.count().await, 1);
registry.register(Arc::new(TestSource::new("test-2", "Test 2"))).await;
registry
.register(Arc::new(TestSource::new("test-2", "Test 2")))
.await;
assert_eq!(registry.count().await, 2);
}
#[tokio::test]
async fn test_remove() {
let registry = SourceRegistry::new();
registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "Test 1")))
.await;
assert!(registry.contains("test-1").await);
assert!(registry.remove("test-1").await);
@@ -345,7 +357,9 @@ mod tests {
assert!(!registry.contains("test-1").await);
registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "Test 1")))
.await;
assert!(registry.contains("test-1").await);
assert!(!registry.contains("test-2").await);
@@ -355,11 +369,15 @@ mod tests {
async fn test_replace_source() {
let registry = SourceRegistry::new();
registry.register(Arc::new(TestSource::new("test-1", "Old Name"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "Old Name")))
.await;
let old = registry.get("test-1").await.unwrap();
assert_eq!(old.name(), "Old Name");
registry.register(Arc::new(TestSource::new("test-1", "New Name"))).await;
registry
.register(Arc::new(TestSource::new("test-1", "New Name")))
.await;
let new = registry.get("test-1").await.unwrap();
assert_eq!(new.name(), "New Name");

View File

@@ -3,8 +3,8 @@
//! Ce module fournit des helpers pour créer et enregistrer facilement des sources
//! musicales préconfigurées à partir de la configuration système.
use pmosource::MusicSourceExt;
use pmoserver::Server;
use pmosource::MusicSourceExt;
use std::sync::Arc;
/// Erreur lors de l'initialisation d'une source
@@ -93,7 +93,11 @@ pub trait SourcesExt {
/// server.register_qobuz_with_credentials("user@example.com", "password").await?;
/// ```
#[cfg(feature = "qobuz")]
async fn register_qobuz_with_credentials(&mut self, username: &str, password: &str) -> Result<()>;
async fn register_qobuz_with_credentials(
&mut self,
username: &str,
password: &str,
) -> Result<()>;
/// Enregistre la source Radio Paradise
///
@@ -141,7 +145,11 @@ impl SourcesExt for Server {
}
#[cfg(feature = "qobuz")]
async fn register_qobuz_with_credentials(&mut self, username: &str, password: &str) -> Result<()> {
async fn register_qobuz_with_credentials(
&mut self,
username: &str,
password: &str,
) -> Result<()> {
use pmoqobuz::{QobuzClient, QobuzSource};
tracing::info!("Initializing Qobuz source with explicit credentials...");
@@ -165,18 +173,19 @@ impl SourcesExt for Server {
#[cfg(feature = "paradise")]
async fn register_paradise(&mut self) -> Result<()> {
use pmoparadise::{RadioParadiseClient, RadioParadiseSource, RadioParadiseExt};
use pmoparadise::{RadioParadiseClient, RadioParadiseExt, RadioParadiseSource};
tracing::info!("Initializing Radio Paradise source...");
// Créer le client (Radio Paradise ne nécessite pas d'authentification)
let client = RadioParadiseClient::new()
.await
.map_err(|e| SourceInitError::ParadiseError(format!("Failed to create client: {}", e)))?;
let client = RadioParadiseClient::new().await.map_err(|e| {
SourceInitError::ParadiseError(format!("Failed to create client: {}", e))
})?;
// Créer la source depuis le registry avec capacité FIFO par défaut
let source = RadioParadiseSource::from_registry_default(client)
.map_err(|e| SourceInitError::ParadiseError(format!("Failed to create source: {}", e)))?;
let source = RadioParadiseSource::from_registry_default(client).map_err(|e| {
SourceInitError::ParadiseError(format!("Failed to create source: {}", e))
})?;
// Enregistrer la source
// Note: La FIFO sera peuplée automatiquement lors du premier browse

View File

@@ -13,13 +13,7 @@
//! Ces endpoints sont définis ici plutôt que dans `pmosource` pour éviter les
//! dépendances circulaires (pmoqobuz et pmoparadise dépendent de pmosource).
use axum::{
extract::Json,
http::StatusCode,
response::IntoResponse,
routing::post,
Router,
};
use axum::{Router, extract::Json, http::StatusCode, response::IntoResponse, routing::post};
use pmosource::MusicSource;
use serde::{Deserialize, Serialize};
use std::sync::Arc;

View File

@@ -43,11 +43,13 @@ async fn main() -> Result<()> {
// Display all tracks
println!("Available Tracks:");
for (index, song) in block.songs_ordered() {
println!(" {}. {} - {} ({:.1}s)",
index,
song.artist,
song.title,
song.duration as f64 / 1000.0);
println!(
" {}. {} - {} ({:.1}s)",
index,
song.artist,
song.title,
song.duration as f64 / 1000.0
);
}
println!();
@@ -67,7 +69,10 @@ async fn main() -> Result<()> {
println!("Track Metadata:");
println!(" Sample Rate: {} Hz", track_stream.metadata.sample_rate);
println!(" Channels: {}", track_stream.metadata.channels);
println!(" Bits Per Sample: {}", track_stream.metadata.bits_per_sample);
println!(
" Bits Per Sample: {}",
track_stream.metadata.bits_per_sample
);
println!(" Total Samples: {}", track_stream.metadata.total_samples);
println!();
@@ -86,10 +91,15 @@ async fn main() -> Result<()> {
let (start, duration) = client.track_position_seconds(&block, index)?;
println!("Track {}: {} - {}", index, song.artist, song.title);
println!(" mpv command:");
println!(" mpv --start={:.3} --length={:.3} '{}'", start, duration, block.url);
println!(
" mpv --start={:.3} --length={:.3} '{}'",
start, duration, block.url
);
println!(" ffmpeg command (extract to file):");
println!(" ffmpeg -ss {:.3} -t {:.3} -i '{}' -c copy track_{}.flac",
start, duration, block.url, index);
println!(
" ffmpeg -ss {:.3} -t {:.3} -i '{}' -c copy track_{}.flac",
start, duration, block.url, index
);
println!();
}

View File

@@ -48,9 +48,11 @@ async fn main() -> Result<()> {
if let Some(rating) = song.rating {
println!(" Rating: {:.1}/10", rating);
}
println!(" Duration: {}:{:02}",
song.duration / 60000,
(song.duration % 60000) / 1000);
println!(
" Duration: {}:{:02}",
song.duration / 60000,
(song.duration % 60000) / 1000
);
// Display cover URL
if let Some(cover) = &song.cover {

View File

@@ -5,7 +5,7 @@
//! - Accessing the embedded WebP image
//! - Optionally saving it to a file
use pmoparadise::{RadioParadiseSource, RadioParadiseClient};
use pmoparadise::{RadioParadiseClient, RadioParadiseSource};
use pmosource::MusicSource;
use std::fs;
use std::io::Write;

View File

@@ -38,7 +38,10 @@ async fn main() -> Result<()> {
eprintln!("Current Block:");
eprintln!(" Event: {}", current_block.event);
eprintln!(" Songs: {}", current_block.song_count());
eprintln!(" Duration: {:.1} minutes", current_block.length as f64 / 60000.0);
eprintln!(
" Duration: {:.1} minutes",
current_block.length as f64 / 60000.0
);
eprintln!(" URL: {}\n", current_block.url);
// Display tracklist
@@ -51,7 +54,10 @@ async fn main() -> Result<()> {
// Prefetch next block in advance
eprintln!("Prefetching next block...");
client.prefetch_next(&current_block).await?;
eprintln!("Next block prefetched: {}\n", client.next_block_url().unwrap());
eprintln!(
"Next block prefetched: {}\n",
client.next_block_url().unwrap()
);
// Stream the block
eprintln!("Streaming block... (writing to stdout)");
@@ -71,12 +77,18 @@ async fn main() -> Result<()> {
// Progress indicator (to stderr so it doesn't interfere with piped audio)
if total_bytes % (1024 * 1024) == 0 {
eprintln!(" Downloaded: {:.1} MB", total_bytes as f64 / 1024.0 / 1024.0);
eprintln!(
" Downloaded: {:.1} MB",
total_bytes as f64 / 1024.0 / 1024.0
);
}
}
eprintln!("\nBlock streaming complete!");
eprintln!("Total downloaded: {:.2} MB", total_bytes as f64 / 1024.0 / 1024.0);
eprintln!(
"Total downloaded: {:.2} MB",
total_bytes as f64 / 1024.0 / 1024.0
);
// In a real application, you would now:
// 1. Get the next block using prefetched metadata

View File

@@ -48,7 +48,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("Look for 'Radio Paradise FLAC' in your DLNA/UPnP clients.");
println!();
println!("ContentDirectory service available at:");
println!(" http://localhost:8080/upnp/device/{}/service/ContentDirectory", server.udn());
println!(
" http://localhost:8080/upnp/device/{}/service/ContentDirectory",
server.udn()
);
println!();
println!("Press Ctrl+C to stop the server.");
println!();

View File

@@ -8,9 +8,9 @@
//! cargo run --example with_cache --features cache
//! ```
use pmoparadise::{RadioParadiseClient, RadioParadiseSource};
use pmocovers::Cache as CoverCache;
use pmoaudiocache::AudioCache;
use pmocovers::Cache as CoverCache;
use pmoparadise::{RadioParadiseClient, RadioParadiseSource};
use pmosource::MusicSource;
use std::sync::Arc;
use tokio::time::{sleep, Duration};
@@ -64,7 +64,13 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Add current song to the source
println!(" Adding current track to FIFO with caching...");
if let Some(song) = &now_playing.current_song {
source.add_song(block.clone(), song, now_playing.current_song_index.unwrap_or(0)).await?;
source
.add_song(
block.clone(),
song,
now_playing.current_song_index.unwrap_or(0),
)
.await?;
println!("✅ Track added and caching started!");
println!(" - Cover image will be cached to: ./cache/covers/");
println!(" - Audio will be cached to: ./cache/audio/\n");
@@ -78,7 +84,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("\n📋 Items in FIFO:");
let items = source.get_items(0, 10).await?;
for (i, item) in items.iter().enumerate() {
println!(" {}. {} - {}",
println!(
" {}. {} - {}",
i + 1,
item.artist.as_deref().unwrap_or("Unknown"),
item.title

View File

@@ -133,11 +133,7 @@ impl RadioParadiseClient {
#[cfg(feature = "logging")]
tracing::debug!("Fetching block: {}", url);
let response = self.client
.get(url)
.timeout(self.timeout)
.send()
.await?;
let response = self.client.get(url).timeout(self.timeout).send().await?;
if !response.status().is_success() {
return Err(Error::other(format!(
@@ -381,6 +377,9 @@ mod tests {
fn test_cover_url() {
let client = RadioParadiseClient::with_client(Client::new());
let url = client.cover_url("test.jpg").unwrap();
assert_eq!(url.as_str(), "https://img.radioparadise.com/covers/l/test.jpg");
assert_eq!(
url.as_str(),
"https://img.radioparadise.com/covers/l/test.jpg"
);
}
}

View File

@@ -269,10 +269,12 @@ pub use stream::BlockStream;
pub use track::{TrackMetadata, TrackStream};
#[cfg(feature = "mediaserver")]
pub use mediaserver::{RadioParadiseMediaServer, MediaServerBuilder};
pub use mediaserver::{MediaServerBuilder, RadioParadiseMediaServer};
#[cfg(feature = "pmoserver")]
pub use pmoserver_ext::{RadioParadiseExt, RadioParadiseState, RadioParadiseApiDoc, create_api_router};
pub use pmoserver_ext::{
create_api_router, RadioParadiseApiDoc, RadioParadiseExt, RadioParadiseState,
};
// Version information
pub const VERSION: &str = env!("CARGO_PKG_VERSION");

View File

@@ -1,7 +1,7 @@
//! ConnectionManager service implementation
use pmoupnp::services::Service;
use pmoupnp::actions::Action;
use pmoupnp::services::Service;
use pmoupnp::state_variables::StateVariable;
use std::sync::Arc;
@@ -15,23 +15,20 @@ pub fn create_connection_manager_service() -> Service {
service.set_service_id("urn:upnp-org:serviceId:ConnectionManager".to_string());
// State variables
let source_protocol_info = StateVariable::new(
"SourceProtocolInfo".to_string(),
"string".to_string(),
).with_send_events(true)
.with_default_value(get_protocol_info());
let source_protocol_info =
StateVariable::new("SourceProtocolInfo".to_string(), "string".to_string())
.with_send_events(true)
.with_default_value(get_protocol_info());
let sink_protocol_info = StateVariable::new(
"SinkProtocolInfo".to_string(),
"string".to_string(),
).with_send_events(true)
.with_default_value("".to_string());
let sink_protocol_info =
StateVariable::new("SinkProtocolInfo".to_string(), "string".to_string())
.with_send_events(true)
.with_default_value("".to_string());
let current_connection_ids = StateVariable::new(
"CurrentConnectionIDs".to_string(),
"string".to_string(),
).with_send_events(true)
.with_default_value("0".to_string());
let current_connection_ids =
StateVariable::new("CurrentConnectionIDs".to_string(), "string".to_string())
.with_send_events(true)
.with_default_value("0".to_string());
service.add_state_variable(Arc::new(source_protocol_info));
service.add_state_variable(Arc::new(sink_protocol_info));
@@ -39,14 +36,8 @@ pub fn create_connection_manager_service() -> Service {
// GetProtocolInfo action
let mut get_protocol_info = Action::new("GetProtocolInfo".to_string());
get_protocol_info.add_output_argument(
"Source".to_string(),
"SourceProtocolInfo".to_string(),
);
get_protocol_info.add_output_argument(
"Sink".to_string(),
"SinkProtocolInfo".to_string(),
);
get_protocol_info.add_output_argument("Source".to_string(), "SourceProtocolInfo".to_string());
get_protocol_info.add_output_argument("Sink".to_string(), "SinkProtocolInfo".to_string());
service.add_action(Arc::new(get_protocol_info));
// GetCurrentConnectionIDs action
@@ -63,10 +54,7 @@ pub fn create_connection_manager_service() -> Service {
"ConnectionID".to_string(),
"A_ARG_TYPE_ConnectionID".to_string(),
);
get_connection_info.add_output_argument(
"RcsID".to_string(),
"A_ARG_TYPE_RcsID".to_string(),
);
get_connection_info.add_output_argument("RcsID".to_string(), "A_ARG_TYPE_RcsID".to_string());
get_connection_info.add_output_argument(
"AVTransportID".to_string(),
"A_ARG_TYPE_AVTransportID".to_string(),
@@ -83,10 +71,8 @@ pub fn create_connection_manager_service() -> Service {
"PeerConnectionID".to_string(),
"A_ARG_TYPE_ConnectionID".to_string(),
);
get_connection_info.add_output_argument(
"Direction".to_string(),
"A_ARG_TYPE_Direction".to_string(),
);
get_connection_info
.add_output_argument("Direction".to_string(), "A_ARG_TYPE_Direction".to_string());
get_connection_info.add_output_argument(
"Status".to_string(),
"A_ARG_TYPE_ConnectionStatus".to_string(),
@@ -94,34 +80,42 @@ pub fn create_connection_manager_service() -> Service {
service.add_action(Arc::new(get_connection_info));
// Additional state variables for arguments
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_ConnectionID".to_string(), "i4".to_string())
));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_RcsID".to_string(), "i4".to_string())
));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_AVTransportID".to_string(), "i4".to_string())
));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_ProtocolInfo".to_string(), "string".to_string())
));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_ConnectionManager".to_string(), "string".to_string())
));
service.add_state_variable(Arc::new(StateVariable::new(
"A_ARG_TYPE_ConnectionID".to_string(),
"i4".to_string(),
)));
service.add_state_variable(Arc::new(StateVariable::new(
"A_ARG_TYPE_RcsID".to_string(),
"i4".to_string(),
)));
service.add_state_variable(Arc::new(StateVariable::new(
"A_ARG_TYPE_AVTransportID".to_string(),
"i4".to_string(),
)));
service.add_state_variable(Arc::new(StateVariable::new(
"A_ARG_TYPE_ProtocolInfo".to_string(),
"string".to_string(),
)));
service.add_state_variable(Arc::new(StateVariable::new(
"A_ARG_TYPE_ConnectionManager".to_string(),
"string".to_string(),
)));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_Direction".to_string(), "string".to_string())
.with_allowed_values(vec!["Input".to_string(), "Output".to_string()])
.with_allowed_values(vec!["Input".to_string(), "Output".to_string()]),
));
service.add_state_variable(Arc::new(
StateVariable::new("A_ARG_TYPE_ConnectionStatus".to_string(), "string".to_string())
.with_allowed_values(vec![
"OK".to_string(),
"ContentFormatMismatch".to_string(),
"InsufficientBandwidth".to_string(),
"UnreliableChannel".to_string(),
"Unknown".to_string(),
])
StateVariable::new(
"A_ARG_TYPE_ConnectionStatus".to_string(),
"string".to_string(),
)
.with_allowed_values(vec![
"OK".to_string(),
"ContentFormatMismatch".to_string(),
"InsufficientBandwidth".to_string(),
"UnreliableChannel".to_string(),
"Unknown".to_string(),
]),
));
service
@@ -143,7 +137,8 @@ fn get_protocol_info() -> String {
"http-get:*:audio/mpeg:*",
"http-get:*:audio/mp3:*",
"http-get:*:audio/x-mp3:*",
].join(",")
]
.join(",")
}
#[cfg(test)]
@@ -153,8 +148,14 @@ mod tests {
#[test]
fn test_create_connection_manager() {
let service = create_connection_manager_service();
assert_eq!(service.service_type(), "urn:schemas-upnp-org:service:ConnectionManager:1");
assert_eq!(service.service_id(), "urn:upnp-org:serviceId:ConnectionManager");
assert_eq!(
service.service_type(),
"urn:schemas-upnp-org:service:ConnectionManager:1"
);
assert_eq!(
service.service_id(),
"urn:upnp-org:serviceId:ConnectionManager"
);
}
#[test]

Some files were not shown because too many files have changed in this diff Show More