From d3dc49ef52b3f5e0500e11b27292e38781d947ff Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Sat, 11 Oct 2025 00:33:13 +0200 Subject: [PATCH] =?UTF-8?q?Cr=C3=A9ation=20du=20module=20pmoaudio?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pmoaudio/Cargo.toml | 11 ++ pmoaudio/README.md | 165 ++++++++++++++++ pmoaudio/examples/multiroom_demo.rs | 77 ++++++++ pmoaudio/examples/pipeline_demo.rs | 126 ++++++++++++ pmoaudio/examples/simple_pipeline.rs | 53 +++++ pmoaudio/examples/streaming_demo.rs | 52 +++++ pmoaudio/src/audio_chunk.rs | 173 +++++++++++++++++ pmoaudio/src/lib.rs | 90 +++++++++ pmoaudio/src/nodes/buffer_node.rs | 241 +++++++++++++++++++++++ pmoaudio/src/nodes/decoder_node.rs | 147 ++++++++++++++ pmoaudio/src/nodes/dsp_node.rs | 236 ++++++++++++++++++++++ pmoaudio/src/nodes/mod.rs | 146 ++++++++++++++ pmoaudio/src/nodes/sink_node.rs | 199 +++++++++++++++++++ pmoaudio/src/nodes/source_node.rs | 169 ++++++++++++++++ pmoaudio/src/nodes/timer_node.rs | 281 +++++++++++++++++++++++++++ pmoaudio/tests/integration_test.rs | 152 +++++++++++++++ 16 files changed, 2318 insertions(+) create mode 100644 pmoaudio/Cargo.toml create mode 100644 pmoaudio/README.md create mode 100644 pmoaudio/examples/multiroom_demo.rs create mode 100644 pmoaudio/examples/pipeline_demo.rs create mode 100644 pmoaudio/examples/simple_pipeline.rs create mode 100644 pmoaudio/examples/streaming_demo.rs create mode 100644 pmoaudio/src/audio_chunk.rs create mode 100644 pmoaudio/src/lib.rs create mode 100644 pmoaudio/src/nodes/buffer_node.rs create mode 100644 pmoaudio/src/nodes/decoder_node.rs create mode 100644 pmoaudio/src/nodes/dsp_node.rs create mode 100644 pmoaudio/src/nodes/mod.rs create mode 100644 pmoaudio/src/nodes/sink_node.rs create mode 100644 pmoaudio/src/nodes/source_node.rs create mode 100644 pmoaudio/src/nodes/timer_node.rs create mode 100644 pmoaudio/tests/integration_test.rs diff --git a/pmoaudio/Cargo.toml b/pmoaudio/Cargo.toml new file mode 100644 index 00000000..3f00b898 --- /dev/null +++ b/pmoaudio/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "pmoaudio" +version = "0.1.0" +edition = "2021" + +[dependencies] +tokio = { version = "1.42", features = ["full"] } +async-trait = "0.1" + +[dev-dependencies] +tokio-test = "0.4" diff --git a/pmoaudio/README.md b/pmoaudio/README.md new file mode 100644 index 00000000..aceb6a76 --- /dev/null +++ b/pmoaudio/README.md @@ -0,0 +1,165 @@ +# PMOAudio + +Pipeline audio stéréo async optimisé pour Rust, utilisant Tokio. + +## Caractéristiques + +- **Pipeline push-based async** : Tous les nodes utilisent Tokio pour un traitement non-bloquant +- **Zero-copy optimisé** : Les données audio sont partagées via `Arc>` pour éviter les clonages inutiles +- **Support multiroom** : BufferNode avec buffer circulaire et offsets indépendants par abonné +- **TimerNode** : Calcul de position temporelle en temps réel +- **Backpressure** : Channels bounded avec `try_send` pour éviter les blocages + +## Architecture + +### AudioChunk + +Structure de données pour un chunk audio stéréo : + +```rust +pub struct AudioChunk { + pub order: u64, // Numéro d'ordre + pub left: Arc>, // Canal gauche (partagé) + pub right: Arc>, // Canal droit (partagé) + pub sample_rate: u32, // Taux d'échantillonnage +} +``` + +Les données sont wrappées dans `Arc` pour permettre le partage sans copie entre plusieurs abonnés. + +### Nodes + +#### SingleSubscriberNode +- Un seul abonné +- Pas de clone inutile du Arc + +#### MultiSubscriberNode +- Plusieurs abonnés +- Partage le même `Arc` avec tous + +#### SourceNode +- Génère ou lit des chunks audio +- Version mock avec génération de sinusoïdes pour tests + +#### DecoderNode +- Décode les chunks audio +- Supporte le passthrough et le resampling (mock) + +#### DspNode +- Applique des transformations DSP +- Clone les données uniquement si modification nécessaire +- Exemple : gain, filtrage + +#### BufferNode +- Buffer circulaire (`VecDeque>`) +- Support multiroom avec offsets indépendants +- `try_send` non-bloquant pour éviter de bloquer la source + +#### TimerNode +- Node passthrough qui ne modifie pas les données +- Incrémente un compteur de samples +- Calcule la position : `position_sec = elapsed_samples / sample_rate` +- Fournit un `TimerHandle` pour monitoring + +#### SinkNode +- Node terminal qui consomme les chunks +- Versions : silent, logging, stats, mock file writer + +## Pipeline type + +``` +SourceNode → DecoderNode → DSPNode → BufferNode → TimerNode → SinkNode(s) + ↓ + Multiroom Sinks + (avec offsets) +``` + +## Exemples + +### Pipeline simple + +```rust +use pmoaudio::{SinkNode, SourceNode, TimerNode}; + +#[tokio::main] +async fn main() { + let (mut timer, timer_tx) = TimerNode::new(10); + let (sink, sink_tx) = SinkNode::new("Output".to_string(), 10); + + timer.add_subscriber(sink_tx); + let timer_handle = timer.get_position_handle(); + + tokio::spawn(async move { timer.run().await.unwrap() }); + + let sink_handle = tokio::spawn(async move { + sink.run_with_stats().await.unwrap() + }); + + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(timer_tx); + source.generate_chunks(30, 4800, 48000, 440.0).await.unwrap(); + }); + + sink_handle.await.unwrap(); +} +``` + +### Multiroom + +```rust +let (buffer, buffer_tx) = BufferNode::new(50, 10); + +let (sink1, sink1_tx) = SinkNode::new("Room 1".to_string(), 10); +let (sink2, sink2_tx) = SinkNode::new("Room 2".to_string(), 10); + +buffer.add_subscriber_with_offset(sink1_tx, 0).await; // Pas de délai +buffer.add_subscriber_with_offset(sink2_tx, 5).await; // 5 chunks de retard +``` + +## Lancer les exemples + +```bash +# Pipeline simple +cargo run --example simple_pipeline + +# Pipeline complet avec tous les nodes +cargo run --example pipeline_demo + +# Configuration multiroom +cargo run --example multiroom_demo + +# Streaming avec timing réel +cargo run --example streaming_demo +``` + +## Tests + +```bash +cargo test +``` + +20 tests unitaires couvrant : +- Propagation des chunks +- Calcul de position par TimerNode +- BufferNode multi-abonné avec offsets +- Arc sharing et zero-copy +- DSP avec gain et filtrage +- Resampling + +## Optimisations + +1. **Arc sharing** : Les `AudioChunk` sont clonés via `Arc::clone()` qui ne clone que le pointeur +2. **Copy-on-Write** : Les DSP nodes clonent les données uniquement si modification nécessaire +3. **Bounded channels** : Backpressure automatique +4. **try_send** : Non-bloquant pour BufferNode, permet de sauter des chunks si un abonné est saturé +5. **RwLock** : Pour partage concurrent du compteur TimerNode + +## Dépendances + +- `tokio` : Runtime async et channels +- `async-trait` : Traits async + +## License + +MIT diff --git a/pmoaudio/examples/multiroom_demo.rs b/pmoaudio/examples/multiroom_demo.rs new file mode 100644 index 00000000..8464d8a7 --- /dev/null +++ b/pmoaudio/examples/multiroom_demo.rs @@ -0,0 +1,77 @@ +//! Exemple de configuration multiroom avec BufferNode +//! +//! Démontre l'utilisation du buffer circulaire pour synchroniser +//! plusieurs sorties avec des délais différents + +use pmoaudio::{BufferNode, SinkNode, SourceNode}; + +#[tokio::main] +async fn main() { + println!("=== Multiroom Demo ===\n"); + + // Buffer avec capacité pour gérer les délais + let (buffer, buffer_tx) = BufferNode::new(50, 10); + + // Créer 3 sorties avec délais différents + let (sink1, sink1_tx) = SinkNode::new("Room 1 (no delay)".to_string(), 10); + let (sink2, sink2_tx) = SinkNode::new("Room 2 (5 chunks delay)".to_string(), 10); + let (sink3, sink3_tx) = SinkNode::new("Room 3 (10 chunks delay)".to_string(), 10); + + buffer.add_subscriber_with_offset(sink1_tx, 0).await; + buffer.add_subscriber_with_offset(sink2_tx, 5).await; + buffer.add_subscriber_with_offset(sink3_tx, 10).await; + + // Spawn buffer et sinks + tokio::spawn(async move { + buffer.run().await.unwrap(); + }); + + let sink1_handle = tokio::spawn(async move { + let stats = sink1.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + let sink2_handle = tokio::spawn(async move { + let stats = sink2.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + let sink3_handle = tokio::spawn(async move { + let stats = sink3.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + // Générer de l'audio dans une tâche séparée + println!("Generating audio for multiroom playback...\n"); + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(buffer_tx); + source.generate_chunks(30, 4800, 48000, 440.0).await.unwrap(); + }); + + println!("Waiting for all rooms to finish...\n"); + + // Attendre toutes les sorties + let stats1 = sink1_handle.await.unwrap(); + let stats2 = sink2_handle.await.unwrap(); + let stats3 = sink3_handle.await.unwrap(); + + println!("\n=== Multiroom Summary ==="); + println!( + "{}: {} chunks received", + stats1.name, stats1.chunks_received + ); + println!( + "{}: {} chunks received", + stats2.name, stats2.chunks_received + ); + println!( + "{}: {} chunks received", + stats3.name, stats3.chunks_received + ); + + println!("\nNote: Delayed rooms receive fewer chunks due to the offset"); +} diff --git a/pmoaudio/examples/pipeline_demo.rs b/pmoaudio/examples/pipeline_demo.rs new file mode 100644 index 00000000..73962184 --- /dev/null +++ b/pmoaudio/examples/pipeline_demo.rs @@ -0,0 +1,126 @@ +//! Exemple de pipeline audio stéréo complet avec tous les nodes +//! +//! Pipeline: SourceNode → DecoderNode → DspNode → BufferNode → TimerNode → SinkNode(s) + +use pmoaudio::{BufferNode, DecoderNode, DspNode, SinkNode, SourceNode, TimerNode}; + +#[tokio::main] +async fn main() { + println!("=== PMOAudio Pipeline Demo ===\n"); + + // Créer le pipeline de nodes + + // 2. DecoderNode - passthrough dans cet exemple + let (mut decoder, decoder_tx) = DecoderNode::new(10); + + // 3. DspNode - applique un gain de 0.5 + let (mut dsp, dsp_tx) = DspNode::new(10, 0.5); + + // 4. BufferNode - buffer circulaire pour multiroom + let (mut buffer, buffer_tx) = BufferNode::new(100, 10); + + // 5. TimerNode - calcule la position temporelle + let (mut timer, timer_tx) = TimerNode::new(10); + + // 6. SinkNodes - deux destinations finales + let (sink1, sink1_tx) = SinkNode::new("Main Output".to_string(), 10); + let (sink2, sink2_tx) = SinkNode::new("Secondary Output".to_string(), 10); + + // Ajouter un abonné au BufferNode avec offset (multiroom simulation) + let (sink3, sink3_tx) = SinkNode::new("Delayed Output".to_string(), 10); + buffer.add_subscriber_with_offset(sink3_tx, 5).await; // 5 chunks de retard + + // Connecter le pipeline + decoder.add_subscriber(dsp_tx); + dsp.add_subscriber(buffer_tx); + buffer.add_next_subscriber(timer_tx); // BufferNode -> TimerNode + timer.add_subscriber(sink1_tx); + timer.add_subscriber(sink2_tx); + + // Obtenir un handle pour lire la position du TimerNode + let timer_handle = timer.get_position_handle(); + + // Spawn tous les nodes + let decoder_handle = tokio::spawn(async move { + decoder.run_passthrough().await.unwrap(); + }); + + let dsp_handle = tokio::spawn(async move { + dsp.run().await.unwrap(); + }); + + let buffer_handle = tokio::spawn(async move { + buffer.run().await.unwrap(); + }); + + let timer_handle_task = tokio::spawn(async move { + timer.run().await.unwrap(); + }); + + let sink1_handle = tokio::spawn(async move { + let stats = sink1.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + let sink2_handle = tokio::spawn(async move { + sink2.run_silent().await.unwrap(); + }); + + let sink3_handle = tokio::spawn(async move { + let stats = sink3.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + // Spawn une tâche pour afficher la position périodiquement + let position_monitor = tokio::spawn(async move { + for _ in 0..10 { + tokio::time::sleep(tokio::time::Duration::from_millis(200)).await; + let position = timer_handle.position_sec().await; + let samples = timer_handle.elapsed_samples().await; + println!("Position: {:.3} sec ({} samples)", position, samples); + } + }); + + // Générer des chunks audio + println!("Generating audio chunks...\n"); + let chunk_size = 4800; // 100ms à 48kHz + let sample_rate = 48000; + let frequency = 440.0; // La 440Hz + + // Source node dans une tâche séparée + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(decoder_tx); + + // Générer 50 chunks (environ 5 secondes) + source + .generate_chunks(50, chunk_size, sample_rate, frequency) + .await + .unwrap(); + + println!("\nChunks sent. Processing...\n"); + }); + + // Attendre que tous les nodes terminent + decoder_handle.await.unwrap(); + dsp_handle.await.unwrap(); + buffer_handle.await.unwrap(); + timer_handle_task.await.unwrap(); + + let stats1 = sink1_handle.await.unwrap(); + sink2_handle.await.unwrap(); + let stats3 = sink3_handle.await.unwrap(); + position_monitor.await.unwrap(); + + println!("\n=== Pipeline Demo Complete ==="); + println!( + "Main output processed: {} chunks, {:.3} sec", + stats1.chunks_received, stats1.total_duration_sec + ); + println!( + "Delayed output processed: {} chunks, {:.3} sec", + stats3.chunks_received, stats3.total_duration_sec + ); +} diff --git a/pmoaudio/examples/simple_pipeline.rs b/pmoaudio/examples/simple_pipeline.rs new file mode 100644 index 00000000..f7a527be --- /dev/null +++ b/pmoaudio/examples/simple_pipeline.rs @@ -0,0 +1,53 @@ +//! Exemple simple de pipeline audio : Source → Timer → Sink +//! +//! Démontre l'utilisation basique du pipeline avec calcul de position + +use pmoaudio::{SinkNode, SourceNode, TimerNode}; + +#[tokio::main] +async fn main() { + println!("=== Simple Pipeline Example ===\n"); + + // Créer les nodes + let (mut timer, timer_tx) = TimerNode::new(10); + let (sink, sink_tx) = SinkNode::new("Output".to_string(), 10); + + // Connecter + timer.add_subscriber(sink_tx); + + // Handle pour monitorer la position + let timer_handle = timer.get_position_handle(); + + // Spawn timer et sink + tokio::spawn(async move { + timer.run().await.unwrap(); + }); + + let sink_handle = tokio::spawn(async move { + let stats = sink.run_with_stats().await.unwrap(); + stats.display(); + stats + }); + + // Générer quelques secondes d'audio dans une tâche séparée + println!("Generating 440Hz sine wave...\n"); + + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(timer_tx); + + source + .generate_chunks(30, 4800, 48000, 440.0) // ~3 secondes + .await + .unwrap(); + + // La source est drop ici, fermant le channel + }); + + // Attendre la fin + let stats = sink_handle.await.unwrap(); + + let final_position = timer_handle.position_sec().await; + println!("\nFinal position: {:.3} seconds", final_position); + println!("Total duration: {:.3} seconds", stats.total_duration_sec); +} diff --git a/pmoaudio/examples/streaming_demo.rs b/pmoaudio/examples/streaming_demo.rs new file mode 100644 index 00000000..802b4b19 --- /dev/null +++ b/pmoaudio/examples/streaming_demo.rs @@ -0,0 +1,52 @@ +//! Exemple de streaming audio en temps réel +//! +//! Démontre l'utilisation du pipeline avec génération de chunks +//! en temps réel avec timing approprié + +use pmoaudio::{SinkNode, SourceNode, TimerNode}; + +#[tokio::main] +async fn main() { + println!("=== Streaming Demo ===\n"); + println!("Streaming audio in real-time for 3 seconds...\n"); + + let mut source = SourceNode::new(); + let (mut timer, timer_tx) = TimerNode::new(20); + let (sink, sink_tx) = SinkNode::new("Streaming Output".to_string(), 20); + + source.add_subscriber(timer_tx); + timer.add_subscriber(sink_tx); + + let timer_handle = timer.get_position_handle(); + + // Spawn le pipeline + tokio::spawn(async move { + timer.run().await.unwrap(); + }); + + let sink_handle = tokio::spawn(async move { + sink.run_with_logging().await.unwrap(); + }); + + // Monitor la position + let monitor_handle = tokio::spawn(async move { + for _ in 0..15 { + tokio::time::sleep(tokio::time::Duration::from_millis(200)).await; + let position = timer_handle.position_sec().await; + println!("Playback position: {:.3} sec", position); + } + }); + + // Stream des chunks avec timing réel + // 100ms par chunk à 48kHz = 4800 samples + source + .stream_chunks(4800, 48000, 440.0, 3000) // 3 secondes + .await + .unwrap(); + + println!("\nStreaming complete."); + + // Attendre la fin + sink_handle.await.unwrap(); + monitor_handle.await.unwrap(); +} diff --git a/pmoaudio/src/audio_chunk.rs b/pmoaudio/src/audio_chunk.rs new file mode 100644 index 00000000..ad8f547e --- /dev/null +++ b/pmoaudio/src/audio_chunk.rs @@ -0,0 +1,173 @@ +use std::sync::Arc; + +/// Représente un chunk audio stéréo avec données partagées via Arc +/// +/// Cette structure encapsule des données audio stéréo (canaux gauche et droit) +/// en utilisant `Arc>` pour permettre le partage efficace entre plusieurs +/// consumers sans copier les données audio. +/// +/// # Optimisation zero-copy +/// +/// Les données audio sont wrappées dans `Arc`, ce qui signifie que: +/// - Le clonage d'un `AudioChunk` ne clone que les pointeurs Arc (très rapide) +/// - Les données audio réelles ne sont copiées que si nécessaire (Copy-on-Write) +/// - Plusieurs nodes peuvent partager le même chunk simultanément +/// +/// # Exemples +/// +/// ``` +/// use pmoaudio::AudioChunk; +/// +/// // Créer un chunk avec des données générées +/// let left = vec![0.0, 0.1, 0.2, 0.3]; +/// let right = vec![0.0, 0.1, 0.2, 0.3]; +/// let chunk = AudioChunk::new(0, left, right, 48000); +/// +/// assert_eq!(chunk.len(), 4); +/// assert_eq!(chunk.sample_rate, 48000); +/// ``` +#[derive(Debug, Clone)] +pub struct AudioChunk { + /// Numéro d'ordre du chunk dans le flux + /// + /// Permet de suivre l'ordre des chunks et détecter les pertes éventuelles + pub order: u64, + + /// Canal gauche (partagé via Arc pour éviter les clonages) + /// + /// Les samples sont en format float 32-bit, normalement entre -1.0 et 1.0 + pub left: Arc>, + + /// Canal droit (partagé via Arc pour éviter les clonages) + /// + /// Les samples sont en format float 32-bit, normalement entre -1.0 et 1.0 + pub right: Arc>, + + /// Taux d'échantillonnage en Hz + /// + /// Valeurs typiques: 44100, 48000, 96000, 192000 + pub sample_rate: u32, +} + +impl AudioChunk { + /// Crée un nouveau chunk audio + /// + /// Les vecteurs sont automatiquement wrappés dans `Arc`. + /// + /// # Arguments + /// + /// * `order` - Numéro d'ordre du chunk dans le flux + /// * `left` - Samples du canal gauche + /// * `right` - Samples du canal droit + /// * `sample_rate` - Taux d'échantillonnage en Hz + /// + /// # Exemples + /// + /// ``` + /// use pmoaudio::AudioChunk; + /// + /// let chunk = AudioChunk::new( + /// 0, + /// vec![0.0, 0.5, 1.0], + /// vec![0.0, 0.5, 1.0], + /// 48000 + /// ); + /// ``` + pub fn new(order: u64, left: Vec, right: Vec, sample_rate: u32) -> Self { + Self { + order, + left: Arc::new(left), + right: Arc::new(right), + sample_rate, + } + } + + /// Crée un chunk à partir de données déjà wrappées dans Arc + /// + /// Utile pour éviter un double wrapping si les données sont déjà dans Arc. + pub fn from_arc( + order: u64, + left: Arc>, + right: Arc>, + sample_rate: u32, + ) -> Self { + Self { + order, + left, + right, + sample_rate, + } + } + + /// Retourne le nombre d'échantillons par canal + /// + /// # Exemples + /// + /// ``` + /// use pmoaudio::AudioChunk; + /// + /// let chunk = AudioChunk::new(0, vec![0.0; 1000], vec![0.0; 1000], 48000); + /// assert_eq!(chunk.len(), 1000); + /// ``` + pub fn len(&self) -> usize { + self.left.len() + } + + /// Vérifie si le chunk est vide + pub fn is_empty(&self) -> bool { + self.left.is_empty() + } + + /// Clone les données pour permettre une modification (Copy-on-Write) + /// + /// Cette méthode doit être appelée uniquement si vous avez besoin de modifier + /// les données audio. Pour une simple lecture, utilisez directement les champs + /// `left` et `right`. + /// + /// # Exemples + /// + /// ``` + /// use pmoaudio::AudioChunk; + /// + /// let chunk = AudioChunk::new(0, vec![1.0, 2.0], vec![3.0, 4.0], 48000); + /// let (mut left, mut right) = chunk.clone_data(); + /// + /// // Modifier les données + /// for sample in &mut left { + /// *sample *= 0.5; + /// } + /// ``` + pub fn clone_data(&self) -> (Vec, Vec) { + ((*self.left).clone(), (*self.right).clone()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_audio_chunk_creation() { + let left = vec![0.0, 0.1, 0.2]; + let right = vec![0.0, 0.1, 0.2]; + let chunk = AudioChunk::new(0, left, right, 48000); + + assert_eq!(chunk.order, 0); + assert_eq!(chunk.len(), 3); + assert_eq!(chunk.sample_rate, 48000); + assert!(!chunk.is_empty()); + } + + #[test] + fn test_audio_chunk_arc_sharing() { + let left = Arc::new(vec![0.0, 0.1, 0.2]); + let right = Arc::new(vec![0.0, 0.1, 0.2]); + + let chunk1 = AudioChunk::from_arc(0, left.clone(), right.clone(), 48000); + let chunk2 = chunk1.clone(); + + // Vérifier que les Arc pointent vers les mêmes données + assert!(Arc::ptr_eq(&chunk1.left, &chunk2.left)); + assert!(Arc::ptr_eq(&chunk1.right, &chunk2.right)); + } +} diff --git a/pmoaudio/src/lib.rs b/pmoaudio/src/lib.rs new file mode 100644 index 00000000..9e428cf2 --- /dev/null +++ b/pmoaudio/src/lib.rs @@ -0,0 +1,90 @@ +//! PMOAudio - Pipeline audio stéréo async optimisé +//! +//! Cette crate fournit un pipeline audio push-based async utilisant Tokio, +//! optimisé pour minimiser les clonages de données via `Arc>`. +//! +//! # Architecture +//! +//! Le pipeline est composé de nodes asynchrones qui communiquent via des channels Tokio. +//! Les données audio sont encapsulées dans des [`AudioChunk`] et partagées via `Arc` pour +//! éviter les copies inutiles. +//! +//! ## Pipeline type +//! +//! ```text +//! SourceNode → DecoderNode → DSPNode → BufferNode → TimerNode → SinkNode(s) +//! ↓ +//! Multiroom Sinks +//! (avec offsets) +//! ``` +//! +//! # Exemples +//! +//! ## Pipeline simple +//! +//! ```no_run +//! use pmoaudio::{SinkNode, SourceNode, TimerNode}; +//! +//! #[tokio::main] +//! async fn main() { +//! let (mut timer, timer_tx) = TimerNode::new(10); +//! let (sink, sink_tx) = SinkNode::new("Output".to_string(), 10); +//! +//! timer.add_subscriber(sink_tx); +//! +//! tokio::spawn(async move { timer.run().await.unwrap() }); +//! let sink_handle = tokio::spawn(async move { +//! sink.run_with_stats().await.unwrap() +//! }); +//! +//! tokio::spawn(async move { +//! let mut source = SourceNode::new(); +//! source.add_subscriber(timer_tx); +//! source.generate_chunks(30, 4800, 48000, 440.0).await.unwrap(); +//! }); +//! +//! sink_handle.await.unwrap(); +//! } +//! ``` +//! +//! ## Configuration multiroom +//! +//! ```no_run +//! use pmoaudio::{BufferNode, SinkNode}; +//! +//! #[tokio::main] +//! async fn main() { +//! let (buffer, buffer_tx) = BufferNode::new(50, 10); +//! +//! let (sink1, sink1_tx) = SinkNode::new("Room 1".to_string(), 10); +//! let (sink2, sink2_tx) = SinkNode::new("Room 2".to_string(), 10); +//! +//! // Room 1 sans délai, Room 2 avec 5 chunks de retard +//! buffer.add_subscriber_with_offset(sink1_tx, 0).await; +//! buffer.add_subscriber_with_offset(sink2_tx, 5).await; +//! +//! tokio::spawn(async move { buffer.run().await.unwrap() }); +//! // ... spawn sinks et source +//! } +//! ``` +//! +//! # Optimisations +//! +//! - **Zero-copy** : Les [`AudioChunk`] sont partagés via `Arc`, seul le pointeur est cloné +//! - **Copy-on-Write** : Les nodes DSP clonent les données uniquement si modification nécessaire +//! - **Backpressure** : Channels bounded avec `try_send` pour éviter les blocages +//! - **RwLock** : Pour partage concurrent du compteur [`TimerNode`] + +mod audio_chunk; +mod nodes; + +pub use audio_chunk::AudioChunk; +pub use nodes::{ + buffer_node::BufferNode, + decoder_node::DecoderNode, + dsp_node::DspNode, + sink_node::{SinkNode, SinkStats}, + source_node::SourceNode, + timer_node::{TimerHandle, TimerNode}, + AudioError, AudioNode, MultiSubscriberNode, SingleSubscriberNode, +}; diff --git a/pmoaudio/src/nodes/buffer_node.rs b/pmoaudio/src/nodes/buffer_node.rs new file mode 100644 index 00000000..b93703bd --- /dev/null +++ b/pmoaudio/src/nodes/buffer_node.rs @@ -0,0 +1,241 @@ +use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}}; +use std::collections::VecDeque; +use std::sync::Arc; +use tokio::sync::{mpsc, RwLock}; + +/// Subscriber avec son propre offset dans le buffer +struct BufferSubscriber { + tx: mpsc::Sender>, + offset: usize, // Position dans le buffer circulaire +} + +/// BufferNode avec buffer circulaire pour support multiroom +/// +/// Ce node maintient un buffer circulaire de chunks et permet à plusieurs +/// abonnés de lire avec des offsets différents, ce qui est idéal pour des +/// configurations multiroom où différentes pièces peuvent avoir un léger +/// délai de synchronisation. +/// +/// # Fonctionnement +/// +/// - Le buffer est implémenté avec un `VecDeque` de taille fixe +/// - Chaque abonné peut avoir un offset indépendant (en nombre de chunks) +/// - Utilise `try_send` pour éviter de bloquer si un abonné est saturé +/// +/// # Exemples +/// +/// ```no_run +/// use pmoaudio::{BufferNode, SinkNode}; +/// +/// #[tokio::main] +/// async fn main() { +/// let (buffer, buffer_tx) = BufferNode::new(50, 10); +/// +/// let (sink1, sink1_tx) = SinkNode::new("Room 1".to_string(), 10); +/// let (sink2, sink2_tx) = SinkNode::new("Room 2".to_string(), 10); +/// +/// // Room 1 sans délai +/// buffer.add_subscriber_with_offset(sink1_tx, 0).await; +/// +/// // Room 2 avec 5 chunks de retard +/// buffer.add_subscriber_with_offset(sink2_tx, 5).await; +/// +/// tokio::spawn(async move { buffer.run().await.unwrap() }); +/// // ... spawn sinks et source +/// } +/// ``` +pub struct BufferNode { + buffer: Arc>>>, + subscribers: Arc>>, + buffer_size: usize, + rx: mpsc::Receiver>, + next_subscribers: MultiSubscriberNode, // Pour passer au node suivant +} + +impl BufferNode { + /// Crée un nouveau BufferNode + /// + /// # Arguments + /// * `buffer_size` - Taille maximale du buffer circulaire + /// * `channel_size` - Taille du channel bounded pour backpressure + pub fn new(buffer_size: usize, channel_size: usize) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + let node = Self { + buffer: Arc::new(RwLock::new(VecDeque::with_capacity(buffer_size))), + subscribers: Arc::new(RwLock::new(Vec::new())), + buffer_size, + rx, + next_subscribers: MultiSubscriberNode::new(), + }; + + (node, tx) + } + + /// Ajoute un abonné avec un offset spécifique (pour multiroom) + pub async fn add_subscriber_with_offset( + &self, + tx: mpsc::Sender>, + offset: usize, + ) { + let mut subs = self.subscribers.write().await; + subs.push(BufferSubscriber { tx, offset }); + } + + /// Ajoute un abonné sans offset (commence au chunk courant) + pub async fn add_subscriber(&self, tx: mpsc::Sender>) { + self.add_subscriber_with_offset(tx, 0).await; + } + + /// Ajoute un abonné pour le node suivant (sans buffer) + pub fn add_next_subscriber(&mut self, tx: mpsc::Sender>) { + self.next_subscribers.add_subscriber(tx); + } + + /// Démarre la boucle de traitement du BufferNode + pub async fn run(mut self) -> Result<(), AudioError> { + let mut chunk_index = 0usize; + + while let Some(chunk) = self.rx.recv().await { + // Ajouter au buffer circulaire + { + let mut buffer = self.buffer.write().await; + if buffer.len() >= self.buffer_size { + buffer.pop_front(); + } + buffer.push_back(chunk.clone()); + } + + // Envoyer aux abonnés avec offset + { + let buffer = self.buffer.read().await; + let mut subs = self.subscribers.write().await; + + for sub in subs.iter_mut() { + // Calculer l'index dans le buffer en fonction de l'offset + let target_index = if chunk_index >= sub.offset { + chunk_index - sub.offset + } else { + continue; // Pas encore assez de données + }; + + // Vérifier si le chunk est disponible dans le buffer + let buffer_age = chunk_index - target_index; + if buffer_age < buffer.len() { + let chunk_to_send = &buffer[buffer.len() - buffer_age - 1]; + // try_send non-bloquant pour éviter de bloquer la source + let _ = sub.tx.try_send(chunk_to_send.clone()); + } + } + } + + // Push vers les nodes suivants sans buffer + self.next_subscribers.try_push(chunk).await?; + + chunk_index += 1; + } + + Ok(()) + } + + /// Version avec push synchrone au lieu de try_push + pub async fn run_blocking(mut self) -> Result<(), AudioError> { + let mut chunk_index = 0usize; + + while let Some(chunk) = self.rx.recv().await { + // Ajouter au buffer circulaire + { + let mut buffer = self.buffer.write().await; + if buffer.len() >= self.buffer_size { + buffer.pop_front(); + } + buffer.push_back(chunk.clone()); + } + + // Envoyer aux abonnés avec offset + { + let buffer = self.buffer.read().await; + let subs = self.subscribers.read().await; + + for sub in subs.iter() { + let target_index = if chunk_index >= sub.offset { + chunk_index - sub.offset + } else { + continue; + }; + + let buffer_age = chunk_index - target_index; + if buffer_age < buffer.len() { + let chunk_to_send = &buffer[buffer.len() - buffer_age - 1]; + let _ = sub.tx.send(chunk_to_send.clone()).await; + } + } + } + + // Push vers les nodes suivants + for _ in 0..self.next_subscribers.subscribers.len() { + self.next_subscribers.push(chunk.clone()).await?; + } + + chunk_index += 1; + } + + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_buffer_node_basic() { + let (mut node, tx) = BufferNode::new(10, 5); + let (out_tx, mut out_rx) = mpsc::channel(5); + + node.add_next_subscriber(out_tx); + + // Spawn le node + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer des chunks + for i in 0..3 { + let chunk = AudioChunk::new(i, vec![0.0; 100], vec![0.0; 100], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + // Recevoir les chunks + for i in 0..3 { + let chunk = out_rx.recv().await.unwrap(); + assert_eq!(chunk.order, i); + } + } + + #[tokio::test] + async fn test_buffer_node_with_offset() { + let (node, tx) = BufferNode::new(10, 10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + // Ajouter un abonné avec offset de 2 chunks + node.add_subscriber_with_offset(out_tx, 2).await; + + // Spawn le node + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer 5 chunks + for i in 0..5 { + let chunk = AudioChunk::new(i, vec![0.0; 100], vec![0.0; 100], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; + + // L'abonné devrait recevoir les chunks 0, 1, 2 (avec 2 chunks de retard) + let chunk = out_rx.try_recv().unwrap(); + assert_eq!(chunk.order, 0); + } +} diff --git a/pmoaudio/src/nodes/decoder_node.rs b/pmoaudio/src/nodes/decoder_node.rs new file mode 100644 index 00000000..9cbe1e10 --- /dev/null +++ b/pmoaudio/src/nodes/decoder_node.rs @@ -0,0 +1,147 @@ +use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}}; +use std::sync::Arc; +use tokio::sync::mpsc; + +/// DecoderNode - Décode des chunks audio +/// +/// Version mock qui passe simplement les chunks (ou simule un décodage simple) +pub struct DecoderNode { + rx: mpsc::Receiver>, + subscribers: MultiSubscriberNode, +} + +impl DecoderNode { + pub fn new(channel_size: usize) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + let node = Self { + rx, + subscribers: MultiSubscriberNode::new(), + }; + + (node, tx) + } + + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.add_subscriber(tx); + } + + /// Mode passthrough - passe les chunks sans modification + pub async fn run_passthrough(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + self.subscribers.push(chunk).await?; + } + Ok(()) + } + + /// Mode mock décodage - simule un changement de sample rate + pub async fn run_with_resampling(mut self, target_sample_rate: u32) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + if chunk.sample_rate == target_sample_rate { + // Pas besoin de resampling + self.subscribers.push(chunk).await?; + } else { + // Simuler un resampling (mock simple) + let ratio = target_sample_rate as f64 / chunk.sample_rate as f64; + let new_len = (chunk.len() as f64 * ratio) as usize; + + let (left_data, right_data) = chunk.clone_data(); + let mut new_left = Vec::with_capacity(new_len); + let mut new_right = Vec::with_capacity(new_len); + + // Resampling linéaire simple (mock) + for i in 0..new_len { + let src_pos = i as f64 / ratio; + let src_idx = src_pos as usize; + + 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; + + new_left.push(left_sample); + new_right.push(right_sample); + } else if src_idx < left_data.len() { + new_left.push(left_data[src_idx]); + new_right.push(right_data[src_idx]); + } + } + + let new_chunk = AudioChunk::new(chunk.order, new_left, new_right, target_sample_rate); + self.subscribers.push(Arc::new(new_chunk)).await?; + } + } + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_decoder_passthrough() { + let (mut node, tx) = DecoderNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run_passthrough().await.unwrap(); + }); + + // Envoyer un chunk + let chunk = AudioChunk::new(0, vec![1.0, 2.0, 3.0], vec![4.0, 5.0, 6.0], 48000); + let chunk_arc = Arc::new(chunk); + tx.send(chunk_arc.clone()).await.unwrap(); + + // Recevoir le chunk + let received = out_rx.recv().await.unwrap(); + assert!(Arc::ptr_eq(&chunk_arc, &received)); + } + + #[tokio::test] + async fn test_decoder_resampling() { + let (mut node, tx) = DecoderNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run_with_resampling(96000).await.unwrap(); + }); + + // Envoyer un chunk à 48000 Hz + let chunk = AudioChunk::new(0, vec![1.0; 100], vec![1.0; 100], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + + // Recevoir le chunk resampleé + let received = out_rx.recv().await.unwrap(); + assert_eq!(received.sample_rate, 96000); + // Le chunk devrait être environ 2x plus grand + assert!(received.len() > 150 && received.len() < 250); + } + + #[tokio::test] + async fn test_decoder_no_resampling_needed() { + let (mut node, tx) = DecoderNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run_with_resampling(48000).await.unwrap(); + }); + + // Envoyer un chunk déjà au bon sample rate + let chunk = AudioChunk::new(0, vec![1.0; 100], vec![1.0; 100], 48000); + let chunk_arc = Arc::new(chunk); + tx.send(chunk_arc.clone()).await.unwrap(); + + // Le chunk devrait être passé sans modification + let received = out_rx.recv().await.unwrap(); + assert!(Arc::ptr_eq(&chunk_arc, &received)); + } +} diff --git a/pmoaudio/src/nodes/dsp_node.rs b/pmoaudio/src/nodes/dsp_node.rs new file mode 100644 index 00000000..9d88086c --- /dev/null +++ b/pmoaudio/src/nodes/dsp_node.rs @@ -0,0 +1,236 @@ +use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}}; +use std::sync::Arc; +use tokio::sync::mpsc; + +/// DspNode - Applique des transformations DSP aux chunks audio +/// +/// Clone les données uniquement si elles doivent être modifiées +pub struct DspNode { + rx: mpsc::Receiver>, + subscribers: MultiSubscriberNode, + gain: f32, +} + +impl DspNode { + pub fn new(channel_size: usize, gain: f32) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + let node = Self { + rx, + subscribers: MultiSubscriberNode::new(), + gain, + }; + + (node, tx) + } + + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.add_subscriber(tx); + } + + /// Applique le gain aux chunks + pub async fn run(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + if (self.gain - 1.0).abs() < f32::EPSILON { + // Gain = 1.0, pas de transformation nécessaire + self.subscribers.push(chunk).await?; + } else { + // Clone les données pour les modifier + let (mut left_data, mut right_data) = chunk.clone_data(); + + // Appliquer le gain + for sample in &mut left_data { + *sample *= self.gain; + } + for sample in &mut right_data { + *sample *= self.gain; + } + + let new_chunk = AudioChunk::new( + chunk.order, + left_data, + right_data, + chunk.sample_rate, + ); + + self.subscribers.push(Arc::new(new_chunk)).await?; + } + } + Ok(()) + } + + /// Met à jour le gain dynamiquement (nécessite un `Arc>` dans une version réelle) + pub fn set_gain(&mut self, gain: f32) { + self.gain = gain; + } +} + +/// DspNode avec filtre passe-bas simple (mock) +#[allow(dead_code)] +pub struct LowPassDspNode { + rx: mpsc::Receiver>, + subscribers: MultiSubscriberNode, + alpha: f32, // Coefficient du filtre + prev_left: f32, + prev_right: f32, +} + +impl LowPassDspNode { + #[allow(dead_code)] + pub fn new(channel_size: usize, cutoff_ratio: f32) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + // Filtre RC simple: alpha = dt / (RC + dt) + // cutoff_ratio entre 0 (tout couper) et 1 (tout passer) + let alpha = cutoff_ratio.clamp(0.0, 1.0); + + let node = Self { + rx, + subscribers: MultiSubscriberNode::new(), + alpha, + prev_left: 0.0, + prev_right: 0.0, + }; + + (node, tx) + } + + #[allow(dead_code)] + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.add_subscriber(tx); + } + + #[allow(dead_code)] + pub async fn run(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + let (left_data, right_data) = chunk.clone_data(); + let mut new_left = Vec::with_capacity(left_data.len()); + let mut new_right = Vec::with_capacity(right_data.len()); + + // Appliquer le filtre + for &sample in &left_data { + self.prev_left = self.prev_left + self.alpha * (sample - self.prev_left); + new_left.push(self.prev_left); + } + + for &sample in &right_data { + self.prev_right = self.prev_right + self.alpha * (sample - self.prev_right); + new_right.push(self.prev_right); + } + + let new_chunk = AudioChunk::new( + chunk.order, + new_left, + new_right, + chunk.sample_rate, + ); + + self.subscribers.push(Arc::new(new_chunk)).await?; + } + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_dsp_node_unity_gain() { + let (mut node, tx) = DspNode::new(10, 1.0); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer un chunk + let chunk = AudioChunk::new(0, vec![1.0, 2.0, 3.0], vec![4.0, 5.0, 6.0], 48000); + let chunk_arc = Arc::new(chunk); + tx.send(chunk_arc.clone()).await.unwrap(); + + // Avec gain = 1.0, le chunk ne devrait pas être cloné + let received = out_rx.recv().await.unwrap(); + assert!(Arc::ptr_eq(&chunk_arc, &received)); + } + + #[tokio::test] + async fn test_dsp_node_gain() { + let (mut node, tx) = DspNode::new(10, 2.0); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer un chunk + let chunk = AudioChunk::new(0, vec![1.0, 2.0, 3.0], vec![4.0, 5.0, 6.0], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + + // Vérifier que le gain a été appliqué + let received = out_rx.recv().await.unwrap(); + assert_eq!(received.left[0], 2.0); + assert_eq!(received.left[1], 4.0); + assert_eq!(received.left[2], 6.0); + assert_eq!(received.right[0], 8.0); + assert_eq!(received.right[1], 10.0); + assert_eq!(received.right[2], 12.0); + } + + #[tokio::test] + async fn test_lowpass_dsp_node() { + let (mut node, tx) = LowPassDspNode::new(10, 0.5); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer un chunk avec un signal carré + let chunk = AudioChunk::new( + 0, + vec![1.0, 1.0, 1.0, -1.0, -1.0, -1.0], + vec![1.0, 1.0, 1.0, -1.0, -1.0, -1.0], + 48000, + ); + tx.send(Arc::new(chunk)).await.unwrap(); + + // Le filtre devrait lisser le signal + let received = out_rx.recv().await.unwrap(); + + // Vérifier que le signal est lissé (valeurs intermédiaires) + assert!(received.left[0].abs() < 1.0); // Premier échantillon lissé + assert!(received.left[2].abs() < 1.0); // Signal ne devrait pas atteindre 1.0 immédiatement + } + + #[tokio::test] + async fn test_dsp_node_multiple_subscribers() { + let (mut node, tx) = DspNode::new(10, 0.5); + let (out_tx1, mut out_rx1) = mpsc::channel(10); + let (out_tx2, mut out_rx2) = mpsc::channel(10); + + node.add_subscriber(out_tx1); + node.add_subscriber(out_tx2); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + let chunk = AudioChunk::new(0, vec![2.0, 4.0], vec![2.0, 4.0], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + + // Les deux abonnés devraient recevoir le même Arc + let received1 = out_rx1.recv().await.unwrap(); + let received2 = out_rx2.recv().await.unwrap(); + + assert!(Arc::ptr_eq(&received1, &received2)); + assert_eq!(received1.left[0], 1.0); // 2.0 * 0.5 + assert_eq!(received1.left[1], 2.0); // 4.0 * 0.5 + } +} diff --git a/pmoaudio/src/nodes/mod.rs b/pmoaudio/src/nodes/mod.rs new file mode 100644 index 00000000..3a5ade38 --- /dev/null +++ b/pmoaudio/src/nodes/mod.rs @@ -0,0 +1,146 @@ +//! Nodes du pipeline audio +//! +//! Ce module contient tous les types de nodes disponibles pour construire +//! un pipeline audio, ainsi que les traits et structures de support. + +use crate::AudioChunk; +use std::sync::Arc; +use tokio::sync::mpsc; + +pub mod buffer_node; +pub mod decoder_node; +pub mod dsp_node; +pub mod sink_node; +pub mod source_node; +pub mod timer_node; + +/// Trait de base pour tous les nodes audio +/// +/// Tous les nodes du pipeline implémentent ce trait pour permettre +/// une interface uniforme de traitement des chunks audio. +#[async_trait::async_trait] +pub trait AudioNode: Send + Sync { + /// Push un chunk vers ce node + /// + /// # Erreurs + /// + /// Retourne `AudioError::SendError` si l'envoi échoue + async fn push(&mut self, chunk: Arc) -> Result<(), AudioError>; + + /// Ferme le node proprement + async fn close(&mut self); +} + +/// Node avec un seul abonné (pas de clone inutile) +/// +/// Optimisé pour les cas où un node n'a qu'un seul destinataire. +/// Le Arc du chunk est simplement transféré sans clonage supplémentaire. +/// +/// # Exemples +/// +/// ``` +/// use pmoaudio::SingleSubscriberNode; +/// use tokio::sync::mpsc; +/// +/// let (tx, rx) = mpsc::channel(10); +/// let node = SingleSubscriberNode::new(tx); +/// ``` +pub struct SingleSubscriberNode { + tx: mpsc::Sender>, +} + +impl SingleSubscriberNode { + pub fn new(tx: mpsc::Sender>) -> Self { + Self { tx } + } + + pub async fn push(&self, chunk: Arc) -> Result<(), AudioError> { + self.tx + .send(chunk) + .await + .map_err(|_| AudioError::SendError) + } +} + +/// Node avec plusieurs abonnés (partage le même Arc) +/// +/// Permet de broadcaster un chunk à plusieurs destinations. +/// Tous les abonnés reçoivent le même `Arc`, donc pas de copie +/// des données audio - seul le compteur de référence Arc est incrémenté. +/// +/// # Exemples +/// +/// ``` +/// use pmoaudio::MultiSubscriberNode; +/// use tokio::sync::mpsc; +/// +/// let mut node = MultiSubscriberNode::new(); +/// let (tx1, rx1) = mpsc::channel(10); +/// let (tx2, rx2) = mpsc::channel(10); +/// +/// node.add_subscriber(tx1); +/// node.add_subscriber(tx2); +/// // Les deux abonnés recevront les mêmes chunks +/// ``` +pub struct MultiSubscriberNode { + subscribers: Vec>>, +} + +impl MultiSubscriberNode { + pub fn new() -> Self { + Self { + subscribers: Vec::new(), + } + } + + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.push(tx); + } + + pub async fn push(&self, chunk: Arc) -> Result<(), AudioError> { + for tx in &self.subscribers { + // On partage le même Arc avec tous les abonnés + tx.send(chunk.clone()) + .await + .map_err(|_| AudioError::SendError)?; + } + Ok(()) + } + + pub async fn try_push(&self, chunk: Arc) -> Result<(), AudioError> { + for tx in &self.subscribers { + // try_send non-bloquant, ignore si saturé + let _ = tx.try_send(chunk.clone()); + } + Ok(()) + } +} + +impl Default for MultiSubscriberNode { + fn default() -> Self { + Self::new() + } +} + +/// Erreurs possibles dans le pipeline audio +#[derive(Debug, Clone)] +pub enum AudioError { + /// Échec d'envoi d'un chunk à travers un channel + SendError, + /// Échec de réception d'un chunk depuis un channel + ReceiveError, + /// Erreur de traitement avec message descriptif + ProcessingError(String), +} + +impl std::fmt::Display for AudioError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + AudioError::SendError => write!(f, "Failed to send audio chunk"), + AudioError::ReceiveError => write!(f, "Failed to receive audio chunk"), + AudioError::ProcessingError(msg) => write!(f, "Processing error: {}", msg), + } + } +} + +impl std::error::Error for AudioError {} diff --git a/pmoaudio/src/nodes/sink_node.rs b/pmoaudio/src/nodes/sink_node.rs new file mode 100644 index 00000000..fff47537 --- /dev/null +++ b/pmoaudio/src/nodes/sink_node.rs @@ -0,0 +1,199 @@ +use crate::{AudioChunk, nodes::AudioError}; +use std::sync::Arc; +use tokio::sync::mpsc; + +/// SinkNode - Node terminal qui consomme les chunks audio +/// +/// Version mock pour tests et logging +pub struct SinkNode { + rx: mpsc::Receiver>, + name: String, +} + +impl SinkNode { + pub fn new(name: String, channel_size: usize) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + let node = Self { rx, name }; + + (node, tx) + } + + /// Version silencieuse - consomme les chunks sans action + pub async fn run_silent(mut self) -> Result<(), AudioError> { + while let Some(_chunk) = self.rx.recv().await { + // Ne rien faire, juste consommer + } + Ok(()) + } + + /// Version avec logging + pub async fn run_with_logging(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + println!( + "[{}] Received chunk #{} - {} samples @ {} Hz", + self.name, + chunk.order, + chunk.len(), + chunk.sample_rate + ); + } + Ok(()) + } + + /// Version avec statistiques + pub async fn run_with_stats(mut self) -> Result { + let mut stats = SinkStats::new(self.name.clone()); + + while let Some(chunk) = self.rx.recv().await { + stats.process_chunk(&chunk); + } + + Ok(stats) + } + + /// Version mock pour écriture dans un fichier (simule l'écriture) + pub async fn run_mock_file_writer(mut self) -> Result { + let mut total_samples = 0; + + while let Some(chunk) = self.rx.recv().await { + total_samples += chunk.len(); + // Simuler l'écriture avec un petit délai + tokio::time::sleep(tokio::time::Duration::from_micros(10)).await; + } + + Ok(total_samples) + } +} + +/// Statistiques collectées par un SinkNode +#[derive(Debug, Clone)] +pub struct SinkStats { + pub name: String, + pub chunks_received: u64, + pub total_samples: u64, + pub total_duration_sec: f64, + pub peak_left: f32, + pub peak_right: f32, + pub rms_left: f64, + pub rms_right: f64, +} + +impl SinkStats { + pub fn new(name: String) -> Self { + Self { + name, + chunks_received: 0, + total_samples: 0, + total_duration_sec: 0.0, + peak_left: 0.0, + peak_right: 0.0, + rms_left: 0.0, + rms_right: 0.0, + } + } + + pub fn process_chunk(&mut self, chunk: &AudioChunk) { + self.chunks_received += 1; + self.total_samples += chunk.len() as u64; + self.total_duration_sec += chunk.len() as f64 / chunk.sample_rate as f64; + + // Calculer les peaks + for &sample in chunk.left.iter() { + if sample.abs() > self.peak_left { + self.peak_left = sample.abs(); + } + } + + for &sample in chunk.right.iter() { + if sample.abs() > self.peak_right { + self.peak_right = sample.abs(); + } + } + + // Calculer RMS (moyenne des carrés) + 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 + + 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 + + sum_squares_right) + / self.total_samples as f64) + .sqrt(); + } + + pub fn display(&self) { + println!("\n=== Sink Statistics: {} ===", self.name); + println!("Chunks received: {}", self.chunks_received); + println!("Total samples: {}", self.total_samples); + println!("Total duration: {:.3} sec", self.total_duration_sec); + println!("Peak L/R: {:.3} / {:.3}", self.peak_left, self.peak_right); + println!("RMS L/R: {:.3} / {:.3}", self.rms_left, self.rms_right); + println!("========================\n"); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_sink_node_silent() { + let (node, tx) = SinkNode::new("test".to_string(), 10); + + let handle = tokio::spawn(async move { node.run_silent().await }); + + // Envoyer quelques chunks + for i in 0..3 { + let chunk = AudioChunk::new(i, vec![0.0; 100], vec![0.0; 100], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + drop(tx); + handle.await.unwrap().unwrap(); + } + + #[tokio::test] + async fn test_sink_node_stats() { + let (node, tx) = SinkNode::new("test".to_string(), 10); + + let handle = tokio::spawn(async move { node.run_with_stats().await }); + + // Envoyer des chunks avec signal connu + for i in 0..3 { + let chunk = AudioChunk::new(i, vec![1.0; 1000], vec![0.5; 1000], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + drop(tx); + let stats = handle.await.unwrap().unwrap(); + + assert_eq!(stats.chunks_received, 3); + assert_eq!(stats.total_samples, 3000); + assert_eq!(stats.peak_left, 1.0); + assert_eq!(stats.peak_right, 0.5); + assert!((stats.rms_left - 1.0).abs() < 0.001); + assert!((stats.rms_right - 0.5).abs() < 0.001); + } + + #[tokio::test] + async fn test_sink_node_file_writer() { + let (node, tx) = SinkNode::new("writer".to_string(), 10); + + let handle = tokio::spawn(async move { node.run_mock_file_writer().await }); + + // Envoyer des chunks + for i in 0..5 { + let chunk = AudioChunk::new(i, vec![0.0; 100], vec![0.0; 100], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + drop(tx); + let total_samples = handle.await.unwrap().unwrap(); + + assert_eq!(total_samples, 500); + } +} diff --git a/pmoaudio/src/nodes/source_node.rs b/pmoaudio/src/nodes/source_node.rs new file mode 100644 index 00000000..a93a9fc9 --- /dev/null +++ b/pmoaudio/src/nodes/source_node.rs @@ -0,0 +1,169 @@ +use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}}; +use std::sync::Arc; +use tokio::sync::mpsc; + +/// SourceNode - Génère ou lit des chunks audio depuis une source +/// +/// Ce node est la source du pipeline. Version mock pour tests. +pub struct SourceNode { + subscribers: MultiSubscriberNode, +} + +impl SourceNode { + pub fn new() -> Self { + Self { + subscribers: MultiSubscriberNode::new(), + } + } + + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.add_subscriber(tx); + } + + /// Génère un chunk de test avec une forme d'onde sinusoïdale + pub fn generate_test_chunk( + order: u64, + size: usize, + sample_rate: u32, + frequency: f32, + ) -> AudioChunk { + let mut left = Vec::with_capacity(size); + let mut right = Vec::with_capacity(size); + + for i in 0..size { + let t = (order * size as u64 + i as u64) as f32 / sample_rate as f32; + let sample = (2.0 * std::f32::consts::PI * frequency * t).sin(); + left.push(sample); + right.push(sample * 0.8); // Légèrement différent pour la stéréo + } + + AudioChunk::new(order, left, right, sample_rate) + } + + /// Génère et envoie des chunks de test + pub async fn generate_chunks( + &self, + count: u64, + chunk_size: usize, + sample_rate: u32, + frequency: f32, + ) -> Result<(), AudioError> { + for i in 0..count { + let chunk = Self::generate_test_chunk(i, chunk_size, sample_rate, frequency); + self.subscribers.push(Arc::new(chunk)).await?; + } + Ok(()) + } + + /// Génère des chunks silencieux + pub async fn generate_silence( + &self, + count: u64, + chunk_size: usize, + 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, + ); + self.subscribers.push(Arc::new(chunk)).await?; + } + Ok(()) + } + + /// Version streaming : génère des chunks continuellement avec délai + pub async fn stream_chunks( + &self, + chunk_size: usize, + sample_rate: u32, + frequency: f32, + duration_ms: u64, + ) -> Result<(), AudioError> { + let chunk_duration_ms = (chunk_size as f64 / sample_rate as f64 * 1000.0) as u64; + let mut order = 0u64; + + let start = tokio::time::Instant::now(); + let duration = tokio::time::Duration::from_millis(duration_ms); + + while start.elapsed() < duration { + let chunk = Self::generate_test_chunk(order, chunk_size, sample_rate, frequency); + self.subscribers.push(Arc::new(chunk)).await?; + + order += 1; + + // Attendre pour simuler le timing réel + tokio::time::sleep(tokio::time::Duration::from_millis(chunk_duration_ms)).await; + } + + Ok(()) + } +} + +impl Default for SourceNode { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_source_node_generation() { + let mut source = SourceNode::new(); + let (tx, mut rx) = mpsc::channel(10); + + source.add_subscriber(tx); + + // Générer 3 chunks + source.generate_chunks(3, 100, 48000, 440.0).await.unwrap(); + + // Vérifier la réception + for i in 0..3 { + let chunk = rx.recv().await.unwrap(); + assert_eq!(chunk.order, i); + assert_eq!(chunk.len(), 100); + assert_eq!(chunk.sample_rate, 48000); + } + } + + #[test] + fn test_sine_wave_generation() { + let chunk = SourceNode::generate_test_chunk(0, 48000, 48000, 440.0); + + // Vérifier qu'on a bien une sinusoïde + // À 440 Hz avec 48000 samples/s, on devrait avoir 440 cycles + let left = &*chunk.left; + + // Trouver les passages par zéro + let mut zero_crossings = 0; + for i in 1..left.len() { + if (left[i - 1] < 0.0 && left[i] >= 0.0) || (left[i - 1] >= 0.0 && left[i] < 0.0) { + zero_crossings += 1; + } + } + + // 440 cycles = 880 passages par zéro (approximativement) + assert!(zero_crossings > 850 && zero_crossings < 910); + } + + #[tokio::test] + async fn test_source_node_silence() { + let mut source = SourceNode::new(); + let (tx, mut rx) = mpsc::channel(10); + + source.add_subscriber(tx); + + source.generate_silence(2, 100, 48000).await.unwrap(); + + for _ in 0..2 { + let chunk = rx.recv().await.unwrap(); + assert!(chunk.left.iter().all(|&x| x == 0.0)); + assert!(chunk.right.iter().all(|&x| x == 0.0)); + } + } +} diff --git a/pmoaudio/src/nodes/timer_node.rs b/pmoaudio/src/nodes/timer_node.rs new file mode 100644 index 00000000..6123323f --- /dev/null +++ b/pmoaudio/src/nodes/timer_node.rs @@ -0,0 +1,281 @@ +use crate::{AudioChunk, nodes::{AudioError, MultiSubscriberNode}}; +use std::sync::Arc; +use tokio::sync::{mpsc, RwLock}; + +/// TimerNode - Node passthrough qui calcule la position temporelle +/// +/// Ce node ne modifie pas les données audio, il les passe directement +/// aux abonnés tout en maintenant un compteur de samples pour calculer +/// la position en secondes. +/// +/// # Fonctionnement +/// +/// Pour chaque chunk reçu: +/// 1. Incrémente `elapsed_samples += chunk.len()` +/// 2. Calcule `position_sec = elapsed_samples / sample_rate` +/// 3. Push le chunk (sans modification) vers les abonnés +/// +/// # Utilisation +/// +/// Le TimerNode fournit un [`TimerHandle`] qui permet de lire la position +/// depuis d'autres threads/tasks sans bloquer le pipeline. +/// +/// # Exemples +/// +/// ```no_run +/// use pmoaudio::TimerNode; +/// +/// #[tokio::main] +/// async fn main() { +/// let (mut timer, timer_tx) = TimerNode::new(10); +/// let handle = timer.get_position_handle(); +/// +/// tokio::spawn(async move { +/// timer.run().await.unwrap(); +/// }); +/// +/// // Lire la position depuis un autre thread +/// let position = handle.position_sec().await; +/// println!("Position: {:.2} sec", position); +/// } +/// ``` +pub struct TimerNode { + rx: mpsc::Receiver>, + subscribers: MultiSubscriberNode, + elapsed_samples: Arc>, + current_sample_rate: Arc>, +} + +impl TimerNode { + /// Crée un nouveau TimerNode + pub fn new(channel_size: usize) -> (Self, mpsc::Sender>) { + let (tx, rx) = mpsc::channel(channel_size); + + let node = Self { + rx, + subscribers: MultiSubscriberNode::new(), + elapsed_samples: Arc::new(RwLock::new(0)), + current_sample_rate: Arc::new(RwLock::new(48000)), // Default + }; + + (node, tx) + } + + /// Ajoute un abonné + pub fn add_subscriber(&mut self, tx: mpsc::Sender>) { + self.subscribers.add_subscriber(tx); + } + + /// Retourne la position actuelle en secondes + pub async fn position_sec(&self) -> f64 { + let elapsed = *self.elapsed_samples.read().await; + let sample_rate = *self.current_sample_rate.read().await; + elapsed as f64 / sample_rate as f64 + } + + /// Retourne le nombre total d'échantillons écoulés + pub async fn elapsed_samples(&self) -> u64 { + *self.elapsed_samples.read().await + } + + /// Reset le compteur + pub async fn reset(&self) { + let mut elapsed = self.elapsed_samples.write().await; + *elapsed = 0; + } + + /// Démarre la boucle de traitement du TimerNode + pub async fn run(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + // Mettre à jour le sample rate si nécessaire + { + let mut sr = self.current_sample_rate.write().await; + if *sr != chunk.sample_rate { + *sr = chunk.sample_rate; + } + } + + // Incrémenter le compteur d'échantillons + { + let mut elapsed = self.elapsed_samples.write().await; + *elapsed += chunk.len() as u64; + } + + // Push immédiatement le même chunk vers les abonnés (passthrough) + self.subscribers.push(chunk).await?; + } + + Ok(()) + } + + /// Version non-bloquante avec try_push + pub async fn run_nonblocking(mut self) -> Result<(), AudioError> { + while let Some(chunk) = self.rx.recv().await { + { + let mut sr = self.current_sample_rate.write().await; + if *sr != chunk.sample_rate { + *sr = chunk.sample_rate; + } + } + + { + let mut elapsed = self.elapsed_samples.write().await; + *elapsed += chunk.len() as u64; + } + + self.subscribers.try_push(chunk).await?; + } + + Ok(()) + } + + /// Retourne un handle pour lire la position depuis d'autres threads + pub fn get_position_handle(&self) -> TimerHandle { + TimerHandle { + elapsed_samples: self.elapsed_samples.clone(), + current_sample_rate: self.current_sample_rate.clone(), + } + } +} + +/// Handle pour lire la position du TimerNode depuis d'autres threads +/// +/// Ce handle peut être cloné et utilisé depuis plusieurs threads/tasks +/// pour monitorer la position de lecture sans bloquer le pipeline. +/// +/// # Exemples +/// +/// ```no_run +/// use pmoaudio::TimerNode; +/// +/// #[tokio::main] +/// async fn main() { +/// let (mut timer, _tx) = TimerNode::new(10); +/// let handle = timer.get_position_handle(); +/// let handle_clone = handle.clone(); +/// +/// // Utiliser depuis plusieurs tasks +/// tokio::spawn(async move { +/// loop { +/// let pos = handle_clone.position_sec().await; +/// println!("Position: {:.2}s", pos); +/// tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; +/// } +/// }); +/// } +/// ``` +#[derive(Clone)] +pub struct TimerHandle { + elapsed_samples: Arc>, + current_sample_rate: Arc>, +} + +impl TimerHandle { + /// Retourne la position actuelle en secondes + pub async fn position_sec(&self) -> f64 { + let elapsed = *self.elapsed_samples.read().await; + let sample_rate = *self.current_sample_rate.read().await; + elapsed as f64 / sample_rate as f64 + } + + /// Retourne le nombre total d'échantillons écoulés + pub async fn elapsed_samples(&self) -> u64 { + *self.elapsed_samples.read().await + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_timer_node_position_calculation() { + let (mut node, tx) = TimerNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + let handle = node.get_position_handle(); + + // Spawn le node + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer 3 chunks de 1000 samples à 48000 Hz + for i in 0..3 { + let chunk = AudioChunk::new(i, vec![0.0; 1000], vec![0.0; 1000], 48000); + tx.send(Arc::new(chunk)).await.unwrap(); + } + + // Attendre que les chunks soient traités + for _ in 0..3 { + out_rx.recv().await.unwrap(); + } + + // Vérifier la position + let position = handle.position_sec().await; + let expected = 3000.0 / 48000.0; // 3 chunks * 1000 samples / 48000 Hz + assert!((position - expected).abs() < 0.0001); + + let elapsed = handle.elapsed_samples().await; + assert_eq!(elapsed, 3000); + } + + #[tokio::test] + async fn test_timer_node_passthrough() { + let (mut node, tx) = TimerNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Envoyer un chunk + let chunk = AudioChunk::new(42, vec![1.0, 2.0, 3.0], vec![4.0, 5.0, 6.0], 48000); + let chunk_arc = Arc::new(chunk); + tx.send(chunk_arc.clone()).await.unwrap(); + + // Recevoir le chunk + let received = out_rx.recv().await.unwrap(); + + // Vérifier que c'est le même Arc (pas de clone des données) + assert!(Arc::ptr_eq(&chunk_arc, &received)); + assert_eq!(received.order, 42); + } + + #[tokio::test] + async fn test_timer_node_sample_rate_change() { + let (mut node, tx) = TimerNode::new(10); + let (out_tx, mut out_rx) = mpsc::channel(10); + + node.add_subscriber(out_tx); + + let handle = node.get_position_handle(); + + tokio::spawn(async move { + node.run().await.unwrap(); + }); + + // Chunk à 48000 Hz + let chunk1 = AudioChunk::new(0, vec![0.0; 48000], vec![0.0; 48000], 48000); + tx.send(Arc::new(chunk1)).await.unwrap(); + out_rx.recv().await.unwrap(); + + // Après 48000 samples à 48000 Hz = 1 seconde + let pos1 = handle.position_sec().await; + assert!((pos1 - 1.0).abs() < 0.0001); + + // Chunk à 96000 Hz + let chunk2 = AudioChunk::new(1, vec![0.0; 96000], vec![0.0; 96000], 96000); + tx.send(Arc::new(chunk2)).await.unwrap(); + out_rx.recv().await.unwrap(); + + // Position calculée avec le nouveau sample rate + let pos2 = handle.position_sec().await; + let expected = (48000.0 + 96000.0) / 96000.0; + assert!((pos2 - expected).abs() < 0.0001); + } +} diff --git a/pmoaudio/tests/integration_test.rs b/pmoaudio/tests/integration_test.rs new file mode 100644 index 00000000..2d199a2f --- /dev/null +++ b/pmoaudio/tests/integration_test.rs @@ -0,0 +1,152 @@ +//! Tests d'intégration pour le pipeline audio complet + +use pmoaudio::{BufferNode, DecoderNode, DspNode, SinkNode, SourceNode, TimerNode}; + +#[tokio::test] +async fn test_complete_pipeline() { + // Créer un pipeline complet : Source → Decoder → DSP → Buffer → Timer → Sink + + let (mut decoder, decoder_tx) = DecoderNode::new(10); + let (mut dsp, dsp_tx) = DspNode::new(10, 0.5); // Gain de 0.5 + let (mut buffer, buffer_tx) = BufferNode::new(50, 10); + let (mut timer, timer_tx) = TimerNode::new(10); + let (sink, sink_tx) = SinkNode::new("Integration Test".to_string(), 10); + + // Connecter le pipeline + decoder.add_subscriber(dsp_tx); + dsp.add_subscriber(buffer_tx); + buffer.add_next_subscriber(timer_tx); + timer.add_subscriber(sink_tx); + + let timer_handle = timer.get_position_handle(); + + // Spawn tous les nodes + tokio::spawn(async move { decoder.run_passthrough().await.unwrap() }); + tokio::spawn(async move { dsp.run().await.unwrap() }); + tokio::spawn(async move { buffer.run().await.unwrap() }); + tokio::spawn(async move { timer.run().await.unwrap() }); + + let sink_handle = tokio::spawn(async move { sink.run_with_stats().await.unwrap() }); + + // Générer des chunks + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(decoder_tx); + source.generate_chunks(10, 4800, 48000, 440.0).await.unwrap(); + }); + + // Attendre la fin + let stats = sink_handle.await.unwrap(); + + // Vérifier les résultats + assert_eq!(stats.chunks_received, 10); + assert_eq!(stats.total_samples, 48000); + + // Vérifier que le gain a été appliqué (peak devrait être ~0.5) + assert!(stats.peak_left < 0.51 && stats.peak_left > 0.49); + + // Vérifier la position + let position = timer_handle.position_sec().await; + assert!((position - 1.0).abs() < 0.01); // ~1 seconde +} + +#[tokio::test] +async fn test_multiroom_buffering() { + // Tester le BufferNode avec plusieurs abonnés avec offsets + + let (buffer, buffer_tx) = BufferNode::new(50, 20); + + let (sink1, sink1_tx) = SinkNode::new("Room 1".to_string(), 20); + let (sink2, sink2_tx) = SinkNode::new("Room 2".to_string(), 20); + let (sink3, sink3_tx) = SinkNode::new("Room 3".to_string(), 20); + + buffer.add_subscriber_with_offset(sink1_tx, 0).await; + buffer.add_subscriber_with_offset(sink2_tx, 3).await; + buffer.add_subscriber_with_offset(sink3_tx, 6).await; + + tokio::spawn(async move { buffer.run().await.unwrap() }); + + let sink1_handle = tokio::spawn(async move { sink1.run_with_stats().await.unwrap() }); + let sink2_handle = tokio::spawn(async move { sink2.run_with_stats().await.unwrap() }); + let sink3_handle = tokio::spawn(async move { sink3.run_with_stats().await.unwrap() }); + + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(buffer_tx); + source.generate_chunks(20, 1000, 48000, 440.0).await.unwrap(); + }); + + let stats1 = sink1_handle.await.unwrap(); + let stats2 = sink2_handle.await.unwrap(); + let stats3 = sink3_handle.await.unwrap(); + + // Room 1 devrait avoir tous les chunks + assert_eq!(stats1.chunks_received, 20); + + // Room 2 devrait avoir 3 chunks de moins + assert_eq!(stats2.chunks_received, 17); + + // Room 3 devrait avoir 6 chunks de moins + assert_eq!(stats3.chunks_received, 14); +} + +#[tokio::test] +async fn test_timer_accuracy() { + // Tester la précision du TimerNode + + let (mut timer, timer_tx) = TimerNode::new(10); + let (sink, sink_tx) = SinkNode::new("Timer Test".to_string(), 10); + + timer.add_subscriber(sink_tx); + let timer_handle = timer.get_position_handle(); + + tokio::spawn(async move { timer.run().await.unwrap() }); + let sink_handle = tokio::spawn(async move { sink.run_silent().await.unwrap() }); + + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(timer_tx); + + // 48000 samples à 48kHz = 1 seconde + source.generate_chunks(1, 48000, 48000, 440.0).await.unwrap(); + }); + + sink_handle.await.unwrap(); + + let position = timer_handle.position_sec().await; + let samples = timer_handle.elapsed_samples().await; + + assert_eq!(samples, 48000); + assert!((position - 1.0).abs() < 0.0001); +} + +#[tokio::test] +async fn test_arc_sharing() { + // Vérifier que les chunks sont bien partagés via Arc sans copie + + let (mut timer, timer_tx) = TimerNode::new(10); + let (sink1, sink1_tx) = SinkNode::new("Sink1".to_string(), 10); + let (sink2, sink2_tx) = SinkNode::new("Sink2".to_string(), 10); + + timer.add_subscriber(sink1_tx); + timer.add_subscriber(sink2_tx); + + tokio::spawn(async move { timer.run().await.unwrap() }); + + let sink1_handle = tokio::spawn(async move { sink1.run_with_stats().await.unwrap() }); + let sink2_handle = tokio::spawn(async move { sink2.run_with_stats().await.unwrap() }); + + tokio::spawn(async move { + let mut source = SourceNode::new(); + source.add_subscriber(timer_tx); + source.generate_silence(5, 1000, 48000).await.unwrap(); + }); + + let stats1 = sink1_handle.await.unwrap(); + let stats2 = sink2_handle.await.unwrap(); + + // Les deux sinks devraient avoir reçu les mêmes chunks + assert_eq!(stats1.chunks_received, 5); + assert_eq!(stats2.chunks_received, 5); + assert_eq!(stats1.total_samples, stats2.total_samples); +}