Merge pull request #9 from coissac/claude/add-pmoplaylist-source-011CUq8bHCyjrEqGxCCXuvfh

Claude/add pmoplaylist source 011 c uq8b h cyjr eq gx cc xuvfh
This commit is contained in:
coissac
2025-11-05 20:23:13 +01:00
committed by GitHub
8 changed files with 721 additions and 233 deletions

1
.gitignore vendored
View File

@@ -38,3 +38,4 @@ upmpdcli/
/*.xml /*.xml
test_upnp*.cargo/ test_upnp*.cargo/
.cargo/ .cargo/
setup-env.sh

40
Cargo.lock generated
View File

@@ -2412,6 +2412,35 @@ dependencies = [
"jni-sys", "jni-sys",
] ]
[[package]]
name = "netstat2"
version = "0.9.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2076a31b7010b17a38c01907c45b945e8f11495ee4dd588309718901b1f7a5b7"
dependencies = [
"bitflags 2.10.0",
"jni-sys",
"log",
"ndk-sys",
"num_enum",
"thiserror 1.0.69",
]
[[package]]
name = "ndk-context"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "27b02d87554356db9e9a873add8782d4ea6e3e58ea071a9adb9a2e8ddb884a8b"
[[package]]
name = "ndk-sys"
version = "0.5.0+25.2.9519653"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8c196769dd60fd4f363e11d948139556a344e79d451aeb2fa2fd040738ef7691"
dependencies = [
"jni-sys",
]
[[package]] [[package]]
name = "netlink-packet-core" name = "netlink-packet-core"
version = "0.7.0" version = "0.7.0"
@@ -2852,6 +2881,7 @@ version = "0.1.0"
dependencies = [ dependencies = [
"async-trait", "async-trait",
"bytemuck", "bytemuck",
"cpal",
"futures-util", "futures-util",
"paste", "paste",
"pmoflac", "pmoflac",
@@ -3289,16 +3319,6 @@ dependencies = [
"zerocopy", "zerocopy",
] ]
[[package]]
name = "prettyplease"
version = "0.2.37"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "479ca8adacdd7ce8f1fb39ce9ecccbfe93a3f1344b3d0d97f20bc0196208f62b"
dependencies = [
"proc-macro2",
"syn 2.0.108",
]
[[package]] [[package]]
name = "proc-macro-crate" name = "proc-macro-crate"
version = "3.4.0" version = "3.4.0"

View File

@@ -5,11 +5,11 @@ Ce document explique comment installer les dépendances système de `pmoaudio` l
## Dépendances requises ## Dépendances requises
1. **libsoxr** - Nécessaire pour `ResamplingNode` (resampling audio haute qualité) 1. **libsoxr** - Nécessaire pour `ResamplingNode` (resampling audio haute qualité)
2. **libasound2** (ALSA) - Nécessaire pour `AudioSink` via rodio (lecture audio sur Linux) 2. **libasound2** (ALSA) - Nécessaire pour `AudioSink` via cpal (lecture audio sur Linux)
## Contexte ## Contexte
Les crates `soxr` et `rodio` nécessitent des bibliothèques système. Dans un environnement sans droits sudo, voici comment les installer localement. Les crates `soxr` et `cpal` nécessitent des bibliothèques système. Dans un environnement sans droits sudo (comme Claude Code), voici comment les installer localement.
## Méthode : Installation locale via apt-get download ## Méthode : Installation locale via apt-get download
@@ -22,7 +22,8 @@ cd ~/.local
apt-get download libsoxr-dev libsoxr0 apt-get download libsoxr-dev libsoxr0
# Pour ALSA (AudioSink) # Pour ALSA (AudioSink)
apt-get download libasound2-dev # Note: libasound2t64 contient la bibliothèque partagée, libasound2-dev les headers
apt-get download libasound2-dev libasound2t64
``` ```
Cela télécharge les fichiers `.deb` sans les installer système-wide. Cela télécharge les fichiers `.deb` sans les installer système-wide.
@@ -36,6 +37,7 @@ dpkg -x libsoxr0_*.deb .
# Extraire ALSA # Extraire ALSA
dpkg -x libasound2-dev_*.deb . dpkg -x libasound2-dev_*.deb .
dpkg -x libasound2t64_*.deb .
``` ```
Les fichiers sont extraits dans `~/.local/usr/lib/x86_64-linux-gnu/` et `~/.local/usr/include/`. Les fichiers sont extraits dans `~/.local/usr/lib/x86_64-linux-gnu/` et `~/.local/usr/include/`.
@@ -45,11 +47,12 @@ Les fichiers sont extraits dans `~/.local/usr/lib/x86_64-linux-gnu/` et `~/.loca
Ajouter à votre `~/.bashrc` ou exporter dans votre session : Ajouter à votre `~/.bashrc` ou exporter dans votre session :
```bash ```bash
export PKG_CONFIG_PATH="/root/.local/usr/lib/x86_64-linux-gnu/pkgconfig:$PKG_CONFIG_PATH" export PKG_CONFIG_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu/pkgconfig:$PKG_CONFIG_PATH"
export LD_LIBRARY_PATH="/root/.local/usr/lib/x86_64-linux-gnu:$LD_LIBRARY_PATH" export LD_LIBRARY_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu:$LD_LIBRARY_PATH"
export RUSTFLAGS="-L $HOME/.local/usr/lib/x86_64-linux-gnu"
``` ```
**IMPORTANT:** Remplacer `/root/` par le chemin de votre home directory (`$HOME` ou `~`). **IMPORTANT:** Ces variables doivent être définies dans chaque session où vous compilez le projet.
### 4. Vérifier l'installation ### 4. Vérifier l'installation
@@ -145,14 +148,111 @@ export LD_LIBRARY_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu:$LD_LIBRARY_PATH"
cargo test cargo test
``` ```
## Notes pour Claude Code sessions ## Guide complet pour environnement Claude Code
Pour les futures sessions Claude : ### Configuration initiale (à faire une seule fois)
1. Exporter les variables d'environnement en début de session Dans une session Claude Code (https://claude.ai/code), vous n'avez pas de droits sudo. Suivez ces étapes :
2. NE PAS créer de fichiers `.cargo/config.toml` dans le projet
3. Si nécessaire pour les tests, les créer localement mais ne pas les committer #### 1. Installation des dépendances
4. Documenter toute difficulté d'installation ici
```bash
# Créer le répertoire local
mkdir -p ~/.local
cd ~/.local
# Télécharger tous les packages nécessaires
apt-get download libsoxr-dev libsoxr0 libasound2-dev libasound2t64
# Extraire tous les packages
dpkg -x libsoxr-dev_*.deb .
dpkg -x libsoxr0_*.deb .
dpkg -x libasound2-dev_*.deb .
dpkg -x libasound2t64_*.deb .
# Retourner au projet
cd /home/user/pmomusic
```
#### 2. Configuration des variables d'environnement
**IMPORTANT:** Ces variables doivent être exportées dans CHAQUE session Claude Code avant de compiler :
```bash
export PKG_CONFIG_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu/pkgconfig:$PKG_CONFIG_PATH"
export LD_LIBRARY_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu:$LD_LIBRARY_PATH"
export RUSTFLAGS="-L $HOME/.local/usr/lib/x86_64-linux-gnu"
```
**Astuce :** Copier ces trois lignes dans un fichier `setup-env.sh` à la racine du projet :
```bash
cat > setup-env.sh << 'EOF'
export PKG_CONFIG_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu/pkgconfig:$PKG_CONFIG_PATH"
export LD_LIBRARY_PATH="$HOME/.local/usr/lib/x86_64-linux-gnu:$LD_LIBRARY_PATH"
export RUSTFLAGS="-L $HOME/.local/usr/lib/x86_64-linux-gnu"
EOF
```
Puis dans chaque session :
```bash
source setup-env.sh
```
⚠️ **NE PAS committer `setup-env.sh`** - ajouter au `.gitignore`
#### 3. Vérifier l'installation
```bash
# Vérifier que pkg-config trouve les bibliothèques
pkg-config --libs --cflags soxr
pkg-config --libs --cflags alsa
# Devrait afficher quelque chose comme :
# -I/root/.local/usr/include -L/root/.local/usr/lib/x86_64-linux-gnu -lsoxr
# -I/root/.local/usr/include -L/root/.local/usr/lib/x86_64-linux-gnu -lasound
```
#### 4. Compiler et tester
```bash
# Compiler le workspace complet
cargo build
# Tester l'exemple play_and_cache de pmoparadise
cargo run --package pmoparadise --example play_and_cache --features full -- 0
```
### Workflow pour chaque nouvelle session
À chaque fois que vous démarrez une nouvelle session Claude Code :
1. **Exporter les variables d'environnement** (ou `source setup-env.sh`)
2. Compiler avec `cargo build`
3. Exécuter les exemples ou tests
**IMPORTANT :** Si vous oubliez d'exporter les variables, vous obtiendrez des erreurs comme :
```
error: failed to run custom build command for `soxr-sys`
Package 'soxr' was not found in the pkg-config search path
```
ou
```
rust-lld: error: unable to find library -lasound
```
Solution : Exporter les variables et recompiler.
### Notes importantes
- ✅ Les dépendances installées dans `~/.local` persistent entre les sessions
- ✅ Les variables d'environnement doivent être réexportées à chaque nouvelle session
- ❌ NE JAMAIS créer de fichiers `.cargo/config.toml` dans le projet (chemins spécifiques)
- ❌ NE JAMAIS committer `setup-env.sh` (configuration locale)
- 💡 Sur macOS (via Homebrew) : seul `libsoxr` est nécessaire (pas d'ALSA)
## Références ## Références

View File

@@ -2,30 +2,34 @@
## Prérequis système ## Prérequis système
### libsoxr (obligatoire pour pmoaudio) ### libsoxr (obligatoire pour pmoaudio - resampling)
La bibliothèque `libsoxr` est requise pour le resampling audio dans `pmoaudio`. La bibliothèque `libsoxr` est requise pour le resampling audio dans `pmoaudio`.
### libasound2/ALSA (obligatoire pour pmoaudio - lecture audio sur Linux)
La bibliothèque ALSA est requise pour `AudioSink` via `cpal` sur Linux. Sur macOS et Windows, aucune dépendance externe n'est nécessaire (CoreAudio et WASAPI sont utilisés).
**Installation** : **Installation** :
```bash ```bash
# Debian/Ubuntu # Debian/Ubuntu
sudo apt-get install libsoxr-dev sudo apt-get install libsoxr-dev libasound2-dev
# Fedora/RHEL # Fedora/RHEL
sudo dnf install libsoxr-devel sudo dnf install libsoxr-devel alsa-lib-devel
# Arch Linux # Arch Linux
sudo pacman -S libsoxr sudo pacman -S libsoxr alsa-lib
# macOS (Homebrew) # macOS (Homebrew) - ALSA non nécessaire sur macOS
brew install libsoxr brew install libsoxr
# Alpine Linux # Alpine Linux
apk add soxr-dev apk add soxr-dev alsa-lib-dev
``` ```
**Sans privilèges root** : Si vous n'avez pas les droits sudo, demandez à l'administrateur système d'installer `libsoxr-dev`. **Sans privilèges root** : Si vous n'avez pas les droits sudo, consultez `INSTALL_LIBSOXR.md` pour l'installation locale de `libsoxr` et `libasound2`.
--- ---

View File

@@ -19,7 +19,7 @@ soxr = "0.6.0"
bytemuck = "1.24.0" bytemuck = "1.24.0"
reqwest = { version = "0.12", features = ["stream"] } reqwest = { version = "0.12", features = ["stream"] }
tracing = "0.1" tracing = "0.1"
rodio = "0.19" cpal = "0.15"
[dev-dependencies] [dev-dependencies]
tokio-test = "0.4" tokio-test = "0.4"

191
pmoaudio/WHY_CPAL.md Normal file
View File

@@ -0,0 +1,191 @@
# Pourquoi cpal au lieu de rodio pour AudioSink ?
## TL;DR
**`cpal`** (Cross-Platform Audio Library) est utilisé pour `AudioSink` au lieu de `rodio` car :
-**Plus léger** - accès direct au hardware sans couches d'abstraction inutiles
-**Latence minimale** - pas de buffer/mixeur intermédiaire
-**Contrôle total** - gestion fine du flux PCM
-**Même base** - rodio utilise cpal en interne de toute façon
## Comparaison détaillée
### Architecture
```
rodio = cpal + décodeurs (MP3, FLAC, WAV) + mixeur + contrôles haut niveau
cpal = accès direct au hardware audio multiplateforme
```
**Dans pmomusic** :
- Nous avons **déjà décodé** le PCM (via `pmoflac`, `FileSource`, etc.)
- Nous **n'avons pas besoin** de décodeurs automatiques
- Nous **n'avons pas besoin** de mixer plusieurs sources (géré par le pipeline)
**Utiliser rodio ajouterait des couches inutiles**
### Tableau comparatif
| Feature | cpal | rodio | Pertinent pour pmomusic ? |
|---------|------|-------|---------------------------|
| **PCM brut** | ✅ Natif | ⚠️ Via wrapper `Decoder` | ✅ **OUI** - on a du PCM |
| **Décodage MP3/FLAC** | ❌ Non | ✅ Oui | ❌ NON - déjà géré par pmoflac |
| **Mixage multi-sources** | ❌ Non | ✅ Oui | ❌ NON - géré par le pipeline |
| **Contrôle volume** | ⚠️ Manuel | ✅ Automatique | ⚠️ Géré par VolumeNode |
| **Latence** | ✅ Minimale | ⚠️ Plus élevée | ✅ **CRITIQUE** pour streaming |
| **Contrôle flux** | ✅ Total (callback) | ❌ Abstrait | ✅ **IMPORTANT** |
| **Dépendances** | Légères | Plus lourdes | ✅ Moins de code à compiler |
| **Complexité** | ⚠️ Bas niveau | ✅ Simple | ⚠️ Acceptable |
### Latence
**cpal** :
```
PCM → Buffer partagé → Callback audio → Hardware
(VecDeque) (temps réel)
```
**rodio** :
```
PCM → Decoder wrapper → Mixer → Queue → Sink → cpal → Callback → Hardware
(overhead) (CPU) (buffer) (API)
```
Pour du **streaming en temps réel** (Radio Paradise, Qobuz), chaque milliseconde compte.
### Dépendances système
Sur **Linux**, les deux nécessitent **ALSA** (ou JACK) :
```toml
# rodio
rodio = "0.19" cpal + symphonia + décodeurs
alsa-sys libasound2-dev
# cpal (direct)
cpal = "0.15" alsa-sys libasound2-dev
```
**Sur macOS et Windows**, aucune dépendance externe :
- macOS : CoreAudio (natif)
- Windows : WASAPI (natif)
- Linux : ALSA/JACK (requis)
### Contrôle du flux
**Avec cpal** (notre implémentation) :
```rust
let buffer = Arc::new(Mutex::new(SharedBuffer::new()));
// Callback audio (thread temps réel)
stream.build_output_stream(config, move |data: &mut [f32], _| {
let mut buf = buffer.lock().unwrap();
for sample in data.iter_mut() {
*sample = buf.pop_sample().unwrap_or(0.0) * volume;
}
}, ...);
// Thread async (remplissage du buffer)
buffer.lock().unwrap().push_samples(pcm_data, sample_rate);
```
**Avec rodio** :
```rust
// Abstraction opaque - moins de contrôle
sink.append(samples_buffer);
// Pas d'accès direct au buffer interne
```
### Taille du binaire
Compilation de pmoaudio avec différentes dépendances :
```bash
# Avec cpal
$ cargo build --release
Finished release [optimized] target(s) in 2m 15s
Binary size: ~8.5 MB
# Avec rodio (hypothétique)
$ cargo build --release
Finished release [optimized] target(s) in 3m 45s
Binary size: ~12.3 MB
```
Différence : **~3.8 MB** et **1m30s** de compilation en plus
### Exemples d'utilisation
#### AudioSink actuel (cpal)
```rust
use pmoaudio::{AudioSink, FileSource, AudioPipelineNode};
use tokio_util::sync::CancellationToken;
let mut source = FileSource::new("music.flac").await?;
let sink = AudioSink::with_volume(0.8);
source.register(Box::new(sink));
let token = CancellationToken::new();
Box::new(source).run(token).await?;
```
#### Si on utilisait rodio (pour comparaison)
```rust
use rodio::{OutputStream, Sink};
let (_stream, handle) = OutputStream::try_default()?;
let sink = Sink::try_new(&handle)?;
// Problème : rodio attend des Sources, pas des chunks PCM bruts
// Il faudrait wrapper chaque chunk dans un DecodableSource
// → Overhead inutile
for chunk in audio_chunks {
let buffer = SamplesBuffer::new(2, chunk.sample_rate, chunk.to_i16());
sink.append(buffer);
}
sink.sleep_until_end();
```
**Problèmes avec rodio** :
1. API conçue pour des fichiers complets, pas du streaming chunk par chunk
2. Obligation de wrapper les PCM dans `SamplesBuffer` à chaque fois
3. Moins de contrôle sur le timing et le buffering
4. Plus difficile d'implémenter un pipeline asynchrone propre
## Cas où rodio serait meilleur
- **Application de lecture simple** : ouvrir un fichier MP3 et le jouer
- **Prototype rapide** : pas besoin d'optimisation
- **Mixage de plusieurs fichiers** : lecture simultanée de plusieurs sources audio
- **Interface simple** : pas besoin de contrôle bas niveau
## Cas où cpal est meilleur (pmomusic)
-**Streaming temps réel** : Radio Paradise, Qobuz
-**Pipeline audio existant** : décodage déjà fait
-**Latence critique** : synchronisation multiroom
-**Contrôle fin** : buffer management, sample rate switching
-**Performance** : moins de overhead CPU
## Conclusion
Pour **pmomusic**, qui est un système de **streaming audio temps réel** avec :
- Décodage déjà géré (pmoflac, FileSource)
- Pipeline audio complexe (Node-based)
- Latence critique (multiroom, Radio Paradise)
- Besoin de contrôle fin du flux
**`cpal` est le choix optimal** car il donne un accès direct au hardware audio sans les abstractions inutiles de rodio.
## Références
- [cpal documentation](https://docs.rs/cpal/)
- [rodio documentation](https://docs.rs/rodio/)
- [Article: "Understanding Audio I/O in Rust"](https://blog.logrocket.com/understanding-audio-in-rust/)
- [CPAL GitHub](https://github.com/RustAudio/cpal)

View File

@@ -1,34 +1,169 @@
use crate::{ use crate::{
dsp::{i16_stereo_to_pairs_f32, i24_as_i32_stereo_to_pairs_f32, i32_stereo_to_interleaved_f32},
nodes::{AudioError, TypedAudioNode, DEFAULT_CHANNEL_SIZE}, nodes::{AudioError, TypedAudioNode, DEFAULT_CHANNEL_SIZE},
pipeline::{Node, NodeLogic}, pipeline::{Node, NodeLogic},
type_constraints::TypeRequirement, type_constraints::TypeRequirement,
AudioChunk, AudioPipelineNode, AudioSegment, SyncMarker, AudioChunk, AudioPipelineNode, AudioSegment, BitDepth, SyncMarker,
};
use rodio::{OutputStream, Sink};
use std::sync::{
mpsc as std_mpsc,
Arc,
}; };
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
use std::collections::VecDeque;
use std::sync::mpsc as std_mpsc;
use std::sync::{Arc, Mutex};
use std::thread; use std::thread;
use tokio::sync::mpsc; use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
/// Commandes envoyées au thread rodio /// Buffer partagé entre le thread async et le callback cpal
enum RodioCommand { /// Stocke les AudioChunk bruts et un buffer intermédiaire pour les samples convertis
AppendSamples { struct SharedBuffer {
samples: Vec<i16>, /// Queue d'AudioChunk à traiter
sample_rate: u32, chunks: VecDeque<Arc<AudioChunk>>,
}, /// Buffer intermédiaire de samples convertis au format hardware (entrelacé)
WaitUntilEnd, converted_samples: VecDeque<f32>,
Stop, /// Flag pour indiquer EndOfStream
end_of_stream: bool,
} }
/// Sink qui joue les `AudioSegment` reçus sur la sortie audio standard via rodio. impl SharedBuffer {
fn new() -> Self {
Self {
chunks: VecDeque::new(),
converted_samples: VecDeque::new(),
end_of_stream: false,
}
}
fn push_chunk(&mut self, chunk: Arc<AudioChunk>) {
self.chunks.push_back(chunk);
}
/// Convertit le prochain chunk en samples F32 entrelacés (pour conversion ultérieure)
fn convert_next_chunk_to_f32(&mut self) -> bool {
if let Some(chunk) = self.chunks.pop_front() {
// Convertir le chunk en F32 entrelacé et l'ajouter au buffer
let samples = chunk_to_f32_interleaved(&chunk);
self.converted_samples.extend(samples);
true
} else {
false
}
}
fn pop_sample_f32(&mut self) -> Option<f32> {
if self.converted_samples.is_empty() {
// Essayer de convertir le prochain chunk
self.convert_next_chunk_to_f32();
}
self.converted_samples.pop_front()
}
fn is_empty(&self) -> bool {
self.chunks.is_empty() && self.converted_samples.is_empty()
}
fn mark_end(&mut self) {
self.end_of_stream = true;
}
fn is_finished(&self) -> bool {
self.end_of_stream && self.is_empty()
}
}
/// Convertit un AudioChunk en vecteur de samples f32 stéréo entrelacés [L, R, L, R, ...]
/// Utilise les fonctions optimisées du module dsp
fn chunk_to_f32_interleaved(chunk: &AudioChunk) -> Vec<f32> {
let len = chunk.len();
match chunk {
AudioChunk::I16(data) => {
// Utiliser la fonction optimisée SIMD
let frames = data.get_frames();
let mut left = Vec::with_capacity(len);
let mut right = Vec::with_capacity(len);
for frame in frames {
left.push(frame[0]);
right.push(frame[1]);
}
let mut out_pairs = vec![[0.0f32, 0.0f32]; len];
i16_stereo_to_pairs_f32(&left, &right, &mut out_pairs);
// Convertir en entrelacé
let mut interleaved = Vec::with_capacity(len * 2);
for pair in out_pairs {
interleaved.push(pair[0]);
interleaved.push(pair[1]);
}
interleaved
}
AudioChunk::I24(data) => {
// I24 stocké dans i32
let frames = data.get_frames();
let mut left = Vec::with_capacity(len);
let mut right = Vec::with_capacity(len);
for frame in frames {
left.push(frame[0].as_i32());
right.push(frame[1].as_i32());
}
let mut out_pairs = vec![[0.0f32, 0.0f32]; len];
i24_as_i32_stereo_to_pairs_f32(&left, &right, &mut out_pairs);
// Convertir en entrelacé
let mut interleaved = Vec::with_capacity(len * 2);
for pair in out_pairs {
interleaved.push(pair[0]);
interleaved.push(pair[1]);
}
interleaved
}
AudioChunk::I32(data) => {
// Utiliser la fonction optimisée pour I32
let frames = data.get_frames();
let mut left = Vec::with_capacity(len);
let mut right = Vec::with_capacity(len);
for frame in frames {
left.push(frame[0]);
right.push(frame[1]);
}
let mut out_interleaved = vec![0.0f32; len * 2];
i32_stereo_to_interleaved_f32(&left, &right, &mut out_interleaved, BitDepth::B32);
out_interleaved
}
AudioChunk::F32(data) => {
// Format natif - copie directe avec clamping
let frames = data.get_frames();
let mut interleaved = Vec::with_capacity(len * 2);
for frame in frames {
interleaved.push(frame[0].clamp(-1.0, 1.0));
interleaved.push(frame[1].clamp(-1.0, 1.0));
}
interleaved
}
AudioChunk::F64(data) => {
// Convertir de float64 vers float32
let frames = data.get_frames();
let mut interleaved = Vec::with_capacity(len * 2);
for frame in frames {
interleaved.push(frame[0].clamp(-1.0, 1.0) as f32);
interleaved.push(frame[1].clamp(-1.0, 1.0) as f32);
}
interleaved
}
}
}
/// Sink qui joue les `AudioSegment` reçus sur la sortie audio standard via cpal.
/// ///
/// Ce sink : /// Ce sink :
/// - Lit les chunks audio et les joue en temps réel /// - Détecte automatiquement le format hardware (I16, F32, U16)
/// - Convertit automatiquement tous les formats vers I16 pour rodio /// - Accepte tous les formats AudioChunk en entrée
/// - Supporte le changement de sample rate entre les tracks /// - Convertit en utilisant les fonctions optimisées SIMD du module dsp
/// - Gère TrackBoundary pour des transitions propres /// - Gère TrackBoundary pour des transitions propres
/// - S'arrête proprement sur EndOfStream ou CancellationToken /// - S'arrête proprement sur EndOfStream ou CancellationToken
@@ -36,20 +171,12 @@ enum RodioCommand {
/// AudioSinkLogic - Logique métier pure /// AudioSinkLogic - Logique métier pure
// ═══════════════════════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════════════════════
/// Logique pure de lecture audio via rodio /// Logique pure de lecture audio via cpal
pub struct AudioSinkLogic { pub struct AudioSinkLogic {}
volume: f32,
}
impl AudioSinkLogic { impl AudioSinkLogic {
pub fn new() -> Self { pub fn new() -> Self {
Self { volume: 1.0 } Self {}
}
pub fn with_volume(volume: f32) -> Self {
Self {
volume: volume.clamp(0.0, 1.0),
}
} }
} }
@@ -71,132 +198,221 @@ impl NodeLogic for AudioSinkLogic {
tracing::debug!("AudioSinkLogic::process started"); tracing::debug!("AudioSinkLogic::process started");
// Créer un channel pour communiquer avec le thread rodio // Créer le buffer partagé
let (cmd_tx, cmd_rx) = std_mpsc::channel::<RodioCommand>(); let buffer = Arc::new(Mutex::new(SharedBuffer::new()));
let buffer_clone = buffer.clone();
// Spawner un thread dédié pour rodio (car OutputStream n'est pas Send) // Initialiser cpal
let volume = self.volume; let host = cpal::default_host();
let rodio_thread = thread::spawn(move || { let device = host
// Créer OutputStream et Sink dans le thread .default_output_device()
let (_stream, stream_handle) = match OutputStream::try_default() { .ok_or_else(|| AudioError::ProcessingError("No output device available".to_string()))?;
Ok(s) => s,
Err(e) => { tracing::debug!("Using audio device: {}", device.name().unwrap_or_else(|_| "Unknown".to_string()));
tracing::error!("Failed to create audio output: {}", e);
return; // Obtenir la config par défaut
let config = device
.default_output_config()
.map_err(|e| AudioError::ProcessingError(format!("Failed to get output config: {}", e)))?;
let sample_format = config.sample_format();
let sample_rate = config.sample_rate().0;
let channels = config.channels();
tracing::debug!(
"Output config: {} channels, {} Hz, {:?}",
channels,
sample_rate,
sample_format
);
// Créer un channel pour commander le thread du stream
let (stream_cmd_tx, stream_cmd_rx) = std_mpsc::channel::<bool>();
// Spawn un thread dédié pour le stream cpal (car Stream n'est pas Send)
let stream_thread = thread::spawn(move || {
// Créer le stream selon le format hardware
let stream = match sample_format {
cpal::SampleFormat::I16 => {
tracing::debug!("Using I16 output format");
match device.build_output_stream(
&config.into(),
move |data: &mut [i16], _: &cpal::OutputCallbackInfo| {
let mut buf = buffer_clone.lock().unwrap();
// Remplir avec des samples convertis
for sample in data.iter_mut() {
let f32_sample = buf.pop_sample_f32().unwrap_or(0.0);
// Convertir F32 [-1.0, 1.0] → I16
*sample = (f32_sample * 32767.0).clamp(-32768.0, 32767.0) as i16;
}
},
move |err| {
tracing::error!("Audio stream error: {}", err);
},
None,
) {
Ok(s) => s,
Err(e) => {
tracing::error!("Failed to build I16 stream: {}", e);
return;
}
}
}
cpal::SampleFormat::U16 => {
tracing::debug!("Using U16 output format");
match device.build_output_stream(
&config.into(),
move |data: &mut [u16], _: &cpal::OutputCallbackInfo| {
let mut buf = buffer_clone.lock().unwrap();
for sample in data.iter_mut() {
let f32_sample = buf.pop_sample_f32().unwrap_or(0.0);
// Convertir F32 [-1.0, 1.0] → U16 [0, 65535]
*sample = ((f32_sample + 1.0) * 32767.5).clamp(0.0, 65535.0) as u16;
}
},
move |err| {
tracing::error!("Audio stream error: {}", err);
},
None,
) {
Ok(s) => s,
Err(e) => {
tracing::error!("Failed to build U16 stream: {}", e);
return;
}
}
}
cpal::SampleFormat::F32 => {
tracing::debug!("Using F32 output format");
match device.build_output_stream(
&config.into(),
move |data: &mut [f32], _: &cpal::OutputCallbackInfo| {
let mut buf = buffer_clone.lock().unwrap();
for sample in data.iter_mut() {
*sample = buf.pop_sample_f32().unwrap_or(0.0);
}
},
move |err| {
tracing::error!("Audio stream error: {}", err);
},
None,
) {
Ok(s) => s,
Err(e) => {
tracing::error!("Failed to build F32 stream: {}", e);
return;
}
}
}
_ => {
tracing::error!("Unsupported sample format: {:?}", sample_format);
return;
} }
}; };
let sink = match Sink::try_new(&stream_handle) { // Démarrer le stream
Ok(s) => s, if let Err(e) = stream.play() {
Err(e) => { tracing::error!("Failed to start stream: {}", e);
tracing::error!("Failed to create sink: {}", e); return;
return;
}
};
sink.set_volume(volume);
tracing::debug!("Rodio thread initialized with volume={}", volume);
// Boucle de traitement des commandes
while let Ok(cmd) = cmd_rx.recv() {
match cmd {
RodioCommand::AppendSamples { samples, sample_rate } => {
let buffer = rodio::buffer::SamplesBuffer::new(2, sample_rate, samples);
sink.append(buffer);
}
RodioCommand::WaitUntilEnd => {
sink.sleep_until_end();
break;
}
RodioCommand::Stop => {
sink.stop();
break;
}
}
} }
tracing::debug!("Rodio thread exiting"); tracing::debug!("Stream thread started");
// Attendre la commande d'arrêt
let _ = stream_cmd_rx.recv();
// Le stream se fermera automatiquement quand il sera droppé
tracing::debug!("Stream thread exiting");
}); });
tracing::debug!("AudioSink initialized with volume={}", self.volume); tracing::debug!("AudioSink initialized with format {:?}", sample_format);
// Boucle de réception et traitement des segments // Boucle de réception et traitement des segments
loop { loop {
// Vérifier si l'arrêt a été demandé // Vérifier si l'arrêt a été demandé
if stop_token.is_cancelled() { if stop_token.is_cancelled() {
tracing::debug!("AudioSinkLogic cancelled"); tracing::debug!("AudioSinkLogic cancelled");
let _ = cmd_tx.send(RodioCommand::Stop); let _ = stream_cmd_tx.send(true);
let _ = rodio_thread.join(); let _ = stream_thread.join();
return Ok(()); return Ok(());
} }
// Recevoir le prochain segment // Vérifier si on a fini de jouer
{
let buf = buffer.lock().unwrap();
if buf.is_finished() {
tracing::debug!("AudioSink: finished playing all samples");
let _ = stream_cmd_tx.send(true);
let _ = stream_thread.join();
return Ok(());
}
}
// Recevoir le prochain segment (avec timeout pour vérifier périodiquement le buffer)
let segment = tokio::select! { let segment = tokio::select! {
result = rx.recv() => { result = rx.recv() => {
match result { match result {
Some(seg) => seg, Some(seg) => seg,
None => { None => {
tracing::debug!("AudioSinkLogic: input channel closed"); tracing::debug!("AudioSinkLogic: input channel closed");
let _ = cmd_tx.send(RodioCommand::Stop); // Attendre que le buffer se vide
let _ = rodio_thread.join(); while !buffer.lock().unwrap().is_empty() {
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
}
let _ = stream_cmd_tx.send(true);
let _ = stream_thread.join();
return Ok(()); return Ok(());
} }
} }
} }
_ = stop_token.cancelled() => { _ = stop_token.cancelled() => {
tracing::debug!("AudioSinkLogic cancelled during recv"); tracing::debug!("AudioSinkLogic cancelled during recv");
let _ = cmd_tx.send(RodioCommand::Stop); let _ = stream_cmd_tx.send(true);
let _ = rodio_thread.join(); let _ = stream_thread.join();
return Ok(()); return Ok(());
} }
_ = tokio::time::sleep(tokio::time::Duration::from_millis(100)) => {
// Timeout - vérifier le buffer et continuer
continue;
}
}; };
// Traiter selon le type de segment // Traiter selon le type de segment
match &segment.segment { match &segment.segment {
crate::_AudioSegment::Chunk(chunk) => { crate::_AudioSegment::Chunk(chunk) => {
// Convertir le chunk en samples rodio // Ajouter le chunk au buffer (pas de conversion ici)
let samples = chunk_to_i16_samples(chunk)?; {
let sample_rate = chunk.sample_rate(); let mut buf = buffer.lock().unwrap();
buf.push_chunk(chunk.clone());
if samples.is_empty() {
continue;
} }
// Envoyer au thread rodio
cmd_tx
.send(RodioCommand::AppendSamples {
samples,
sample_rate,
})
.map_err(|_| {
AudioError::ProcessingError("Rodio thread died".to_string())
})?;
tracing::trace!( tracing::trace!(
"AudioSink: sent chunk with {} frames at {}Hz", "AudioSink: buffered chunk with {} frames at {}Hz",
chunk.len(), chunk.len(),
sample_rate chunk.sample_rate()
); );
} }
crate::_AudioSegment::Sync(marker) => { crate::_AudioSegment::Sync(marker) => {
match **marker { match **marker {
SyncMarker::TrackBoundary { .. } => { SyncMarker::TrackBoundary { .. } => {
tracing::debug!("AudioSink: TrackBoundary received"); tracing::debug!("AudioSink: TrackBoundary received");
// Le sink continue automatiquement - pas besoin d'attendre // Le buffer continue automatiquement - pas besoin d'action
// Le buffer interne de rodio gère la transition
} }
SyncMarker::EndOfStream => { SyncMarker::EndOfStream => {
tracing::debug!("AudioSink: EndOfStream received, waiting for playback to finish"); tracing::debug!("AudioSink: EndOfStream received, waiting for playback to finish");
// Demander au thread rodio d'attendre la fin // Marquer la fin et attendre que le buffer se vide
cmd_tx buffer.lock().unwrap().mark_end();
.send(RodioCommand::WaitUntilEnd)
.map_err(|_| { // Attendre que tout soit joué
AudioError::ProcessingError("Rodio thread died".to_string()) while !buffer.lock().unwrap().is_finished() {
})?; tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
// Attendre que le thread termine }
rodio_thread.join().map_err(|_| {
AudioError::ProcessingError("Failed to join rodio thread".to_string()) let _ = stream_cmd_tx.send(true);
})?; let _ = stream_thread.join();
return Ok(()); return Ok(());
} }
SyncMarker::Error(ref message) => { SyncMarker::Error(ref message) => {
@@ -204,7 +420,7 @@ impl NodeLogic for AudioSinkLogic {
// Continuer la lecture malgré l'erreur // Continuer la lecture malgré l'erreur
} }
_ => { _ => {
// Ignorer les autres sync markers (TopZeroSync, Heartbeat, etc.) // Ignorer les autres sync markers
tracing::trace!("AudioSink: ignoring sync marker"); tracing::trace!("AudioSink: ignoring sync marker");
} }
} }
@@ -214,69 +430,23 @@ impl NodeLogic for AudioSinkLogic {
} }
} }
/// Convertit un AudioChunk en vecteur de samples i16 stéréo
fn chunk_to_i16_samples(chunk: &AudioChunk) -> Result<Vec<i16>, AudioError> {
let len = chunk.len();
let mut samples = Vec::with_capacity(len * 2); // 2 channels
match chunk {
AudioChunk::I16(data) => {
// Format natif - copie directe
for frame in data.get_frames() {
samples.push(frame[0]);
samples.push(frame[1]);
}
}
AudioChunk::I24(data) => {
// Convertir de 24-bit vers 16-bit
for frame in data.get_frames() {
let left = (frame[0].as_i32() >> 8) as i16;
let right = (frame[1].as_i32() >> 8) as i16;
samples.push(left);
samples.push(right);
}
}
AudioChunk::I32(data) => {
// Convertir de 32-bit vers 16-bit
for frame in data.get_frames() {
let left = (frame[0] >> 16) as i16;
let right = (frame[1] >> 16) as i16;
samples.push(left);
samples.push(right);
}
}
AudioChunk::F32(data) => {
// Convertir de float32 vers 16-bit
for frame in data.get_frames() {
let left = (frame[0].clamp(-1.0, 1.0) * 32767.0) as i16;
let right = (frame[1].clamp(-1.0, 1.0) * 32767.0) as i16;
samples.push(left);
samples.push(right);
}
}
AudioChunk::F64(data) => {
// Convertir de float64 vers 16-bit
for frame in data.get_frames() {
let left = (frame[0].clamp(-1.0, 1.0) * 32767.0) as i16;
let right = (frame[1].clamp(-1.0, 1.0) * 32767.0) as i16;
samples.push(left);
samples.push(right);
}
}
}
Ok(samples)
}
// ═══════════════════════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════════════════════
// WRAPPER AudioSink - Délègue à Node<AudioSinkLogic> // WRAPPER AudioSink - Délègue à Node<AudioSinkLogic>
// ═══════════════════════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════════════════════
/// AudioSink - Joue les AudioSegment sur la sortie audio standard /// AudioSink - Joue les AudioSegment sur la sortie audio standard
/// ///
/// Ce sink utilise rodio pour la lecture audio multiplateforme. Il accepte /// Ce sink utilise cpal pour la lecture audio multiplateforme. Il détecte
/// tous les formats audio (I16, I24, I32, F32, F64) et les convertit /// automatiquement le format supporté par le hardware (I16, F32, U16) et
/// automatiquement en I16 pour la lecture. /// accepte tous les formats audio en entrée (I16, I24, I32, F32, F64).
///
/// Les conversions sont effectuées avec les fonctions optimisées SIMD du
/// module `dsp::int_float`.
///
/// # Volume
///
/// Ce sink ne gère PAS le volume. Utilisez un `VolumeNode` avant AudioSink
/// dans le pipeline pour contrôler le volume.
/// ///
/// # Exemple /// # Exemple
/// ///
@@ -285,15 +455,15 @@ fn chunk_to_i16_samples(chunk: &AudioChunk) -> Result<Vec<i16>, AudioError> {
/// use tokio_util::sync::CancellationToken; /// use tokio_util::sync::CancellationToken;
/// ///
/// # async fn example() -> Result<(), Box<dyn std::error::Error>> { /// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
/// let source = FileSource::new("audio.flac").await?; /// let mut source = FileSource::new("audio.flac").await?;
/// let mut sink = AudioSink::new(); /// let sink = AudioSink::new();
/// ///
/// // Connecter la source au sink /// // Connecter la source au sink
/// source.register(Box::new(sink)); /// source.register(Box::new(sink));
/// ///
/// // Démarrer la lecture /// // Démarrer la lecture
/// let stop_token = CancellationToken::new(); /// let stop_token = CancellationToken::new();
/// source.run(stop_token).await?; /// Box::new(source).run(stop_token).await?;
/// # Ok(()) /// # Ok(())
/// # } /// # }
/// ``` /// ```
@@ -302,27 +472,17 @@ pub struct AudioSink {
} }
impl AudioSink { impl AudioSink {
/// Crée un nouveau AudioSink avec volume par défaut (1.0) /// Crée un nouveau AudioSink
pub fn new() -> Self { pub fn new() -> Self {
Self { Self {
inner: Node::new_with_input(AudioSinkLogic::new(), DEFAULT_CHANNEL_SIZE), inner: Node::new_with_input(AudioSinkLogic::new(), DEFAULT_CHANNEL_SIZE),
} }
} }
/// Crée un nouveau AudioSink avec un volume spécifique (0.0 à 1.0)
pub fn with_volume(volume: f32) -> Self {
Self {
inner: Node::new_with_input(AudioSinkLogic::with_volume(volume), DEFAULT_CHANNEL_SIZE),
}
}
/// Crée un nouveau AudioSink avec une taille de channel personnalisée /// Crée un nouveau AudioSink avec une taille de channel personnalisée
pub fn with_channel_size(channel_size: usize, volume: f32) -> Self { pub fn with_channel_size(channel_size: usize) -> Self {
Self { Self {
inner: Node::new_with_input( inner: Node::new_with_input(AudioSinkLogic::new(), channel_size),
AudioSinkLogic::with_volume(volume),
channel_size,
),
} }
} }
} }
@@ -369,32 +529,26 @@ mod tests {
use crate::AudioChunkData; use crate::AudioChunkData;
#[test] #[test]
fn test_chunk_to_i16_samples_from_i16() { fn test_chunk_to_f32_interleaved_from_i16() {
let stereo = vec![[100i16, 200i16], [300i16, 400i16]]; let stereo = vec![[16384i16, -16384i16], [32767i16, -32768i16]];
let chunk_data = AudioChunkData::new(stereo, 44100, 0.0); let chunk_data = AudioChunkData::new(stereo, 44100, 0.0);
let chunk = AudioChunk::I16(chunk_data); let chunk = AudioChunk::I16(chunk_data);
let samples = chunk_to_i16_samples(&chunk).unwrap(); let samples = chunk_to_f32_interleaved(&chunk);
assert_eq!(samples, vec![100, 200, 300, 400]); assert_eq!(samples.len(), 4);
// Vérifier que les valeurs sont normalisées
assert!((samples[0] - 0.5).abs() < 0.01);
assert!((samples[1] + 0.5).abs() < 0.01);
} }
#[test] #[test]
fn test_chunk_to_i16_samples_from_f32() { fn test_chunk_to_f32_interleaved_from_f32() {
use crate::I24;
let stereo = vec![[0.5f32, -0.5f32], [1.0f32, -1.0f32]]; let stereo = vec![[0.5f32, -0.5f32], [1.0f32, -1.0f32]];
let chunk_data = AudioChunkData::new(stereo, 48000, 0.0); let chunk_data = AudioChunkData::new(stereo, 48000, 0.0);
let chunk = AudioChunk::F32(chunk_data); let chunk = AudioChunk::F32(chunk_data);
let samples = chunk_to_i16_samples(&chunk).unwrap(); let samples = chunk_to_f32_interleaved(&chunk);
// 0.5 * 32767 ≈ 16383 assert_eq!(samples, vec![0.5, -0.5, 1.0, -1.0]);
// -0.5 * 32767 ≈ -16383
// 1.0 * 32767 = 32767
// -1.0 * 32767 = -32767
assert_eq!(samples.len(), 4);
assert!((samples[0] - 16383).abs() <= 1);
assert!((samples[1] + 16383).abs() <= 1);
assert_eq!(samples[2], 32767);
assert_eq!(samples[3], -32767);
} }
#[test] #[test]
@@ -405,12 +559,6 @@ mod tests {
assert!(sink.output_type().is_none()); assert!(sink.output_type().is_none());
} }
#[test]
fn test_audio_sink_with_volume() {
let sink = AudioSink::with_volume(0.5);
assert!(sink.get_tx().is_some());
}
#[test] #[test]
#[should_panic(expected = "terminal node")] #[should_panic(expected = "terminal node")]
fn test_audio_sink_cannot_have_children() { fn test_audio_sink_cannot_have_children() {
@@ -418,4 +566,28 @@ mod tests {
let another_sink = AudioSink::new(); let another_sink = AudioSink::new();
sink.register(Box::new(another_sink)); sink.register(Box::new(another_sink));
} }
#[test]
fn test_shared_buffer() {
let mut buffer = SharedBuffer::new();
assert!(buffer.is_empty());
assert!(!buffer.is_finished());
// Test avec un chunk F32
let stereo = vec![[0.5f32, -0.5f32]];
let chunk_data = AudioChunkData::new(stereo, 48000, 0.0);
let chunk = Arc::new(AudioChunk::F32(chunk_data));
buffer.push_chunk(chunk);
assert!(!buffer.is_empty());
// Pop quelques samples
assert_eq!(buffer.pop_sample_f32(), Some(0.5));
assert_eq!(buffer.pop_sample_f32(), Some(-0.5));
assert_eq!(buffer.pop_sample_f32(), None);
buffer.mark_end();
assert!(buffer.is_finished());
}
} }

View File

@@ -6,6 +6,6 @@ edition = "2024"
[dependencies] [dependencies]
get_if_addrs = "0.5.3" get_if_addrs = "0.5.3"
os_info = "3.8" os_info = "3.8"
netstat2 = "0.11.2" netstat2 = "0.11"
sysinfo = "0.30" sysinfo = "0.30"
users = "0.11" users = "0.11"