Files
pmomusic/pmoparadise/RADIO_PARADISE_STREAM_SOURCE.md
Claude 17dec3351e fix: Use while loop instead of if for robust cache size guarantee
Problem:
- With `if >= CACHE_SIZE`, only ONE element removed per call
- If cache ever had >10 elements (abnormal state), would stay oversized
- Example: 12 elements → if removes 1 → 11 elements → add 1 → 12 elements 

Solution:
- Use `while >= CACHE_SIZE` to remove ALL excess elements
- Example: 12 elements → while removes 2 → 10 elements → add 1 → 10 elements 
- Guarantees exactly ≤10 elements regardless of initial state

Changes:
- mark_block_downloaded(): changed `if` to `while`
- Updated comment to reflect "tous les éléments excédentaires"
- Documentation updated with robustness guarantee
2025-11-05 05:35:09 +00:00

8.6 KiB

RadioParadiseStreamSource - Documentation Technique

Vue d'ensemble

RadioParadiseStreamSource est un nœud source pour pmoaudio qui télécharge et décode les blocs FLAC de Radio Paradise en temps réel, avec gestion automatique des transitions entre pistes (TrackBoundary).

Architecture

Pattern Node

Suit l'architecture séparée logique/pipeline de pmoaudio :

RadioParadiseStreamSource (wrapper)
    └── Node<RadioParadiseStreamSourceLogic>
            └── RadioParadiseStreamSourceLogic (logique métier)

RadioParadiseStreamSourceLogic

Responsabilités :

  • File d'attente : VecDeque<EventId> pour les blocks à télécharger
  • Cache anti-redondance : VecDeque<EventId> pour 10 blocs récents (FIFO)
  • Téléchargement : Fetch bloc FLAC (bitrate=4 uniquement)
  • Décodage : Stream FLAC via pmoflac::decode_audio_stream
  • Timing : Calcul précis pour insertion TrackBoundary

Flux d'exécution

┌─────────────────────────────────────────────────────────┐
│ 1. Attente block ID (timeout 3s)                       │
│    └─> VecDeque::pop_front()                           │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 2. Vérification cache                                   │
│    └─> VecDeque::contains(&event_id)                   │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 3. Téléchargement métadonnées                          │
│    └─> client.get_block(event_id)                      │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 4. Téléchargement FLAC (bitrate=4)                     │
│    └─> client.download_block_file(&block, 4)           │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 5. Décodage streaming                                   │
│    └─> pmoflac::decode_audio_stream(reader)            │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 6. Découpage en chunks                                  │
│    └─> pcm_to_audio_chunk(pcm, sr, bps)                │
└─────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────┐
│ 7. Insertion TrackBoundary (timing sample-based)       │
│    └─> elapsed_ms = (total_samples * 1000) / sr        │
└─────────────────────────────────────────────────────────┘

Timing TrackBoundary

Algorithme

let elapsed_ms = (total_samples * 1000) / sample_rate as u64;

if elapsed_ms >= song.elapsed {
    // Envoyer TrackBoundary AVANT le chunk (même order)
    send_track_boundary(*order, song, block).await;
}

Exemple concret

Bloc FLAC contenant 3 chansons :

  • Song 0 : elapsed = 0ms
  • Song 1 : elapsed = 180000ms (3min)
  • Song 2 : elapsed = 420000ms (7min)

Timeline :

0ms              180000ms            420000ms
│                │                   │
Song 0           TrackBoundary       TrackBoundary
                 └─> Song 1          └─> Song 2

SyncMarker Order

Règle : TrackBoundary a le même order que le chunk suivant.

// TrackBoundary order = 42
AudioSegment::new_sync(42, SyncMarker::TrackBoundary { ... })

// Chunk suivant order = 42
AudioSegment::new_audio(42, AudioChunk::I16(...))

Gestion du cache

Stratégie FIFO simple

const RECENT_BLOCKS_CACHE_SIZE: usize = 10;

fn mark_block_downloaded(&mut self, event_id: EventId) {
    // Retirer tous les éléments excédentaires (garantit <= CACHE_SIZE)
    while self.recent_blocks.len() >= RECENT_BLOCKS_CACHE_SIZE {
        self.recent_blocks.pop_front();
    }

    // Puis ajouter le nouveau bloc
    self.recent_blocks.push_back(event_id);
}

Avantages VecDeque :

  • Ordre FIFO garanti (le plus ancien est toujours retiré)
  • Simple et prévisible
  • Robuste : while garantit exactement 10 éléments max, même en cas d'état anormal
  • Ne dépasse jamais la capacité pré-allouée (retire avant d'ajouter)
  • Pour 10 éléments, contains() en O(n) reste très performant

Support FLAC

Formats supportés

  • 16-bit : AudioChunk::I16
  • 24-bit : AudioChunk::I24

Conversion PCM

match bits_per_sample {
    16 => {
        let samples: Vec<i16> = pcm_data
            .chunks_exact(2)
            .map(|chunk| i16::from_le_bytes([chunk[0], chunk[1]]))
            .collect();
        AudioChunk::I16(...)
    }
    24 => {
        let samples: Vec<I24> = pcm_data
            .chunks_exact(3)
            .map(|chunk| {
                let value = i32::from_le_bytes([chunk[0], chunk[1], chunk[2], 0]) >> 8;
                I24::from_i32(value)
            })
            .collect();
        AudioChunk::I24(...)
    }
}

Métadonnées

TrackMetadata

Champs extraits de Song :

  • title : Titre de la chanson
  • artist : Artiste
  • album : Album (optionnel)
  • year : Année (optionnel)
  • cover_url : URL de la pochette (async via tokio::spawn)

Gestion asynchrone du cover

tokio::spawn(async move {
    if let Ok(mut meta) = metadata_clone.write().await {
        let _ = meta.set_cover_url(Some(cover_url)).await;
    }
});

API Publique

Création

pub fn new(client: RadioParadiseClient, chunk_duration_ms: u32) -> Self

Configuration

pub fn push_block_id(&mut self, event_id: EventId)

Ajoute un block ID à télécharger dans la file d'attente.

Exécution

async fn run(self: Box<Self>, stop_token: CancellationToken) -> Result<(), AudioError>

Hérite de AudioPipelineNode.

Exemple d'utilisation

Voir examples/radio_paradise_stream.rs pour :

  • Utilisation basique
  • Intégration avec nowplaying stream
  • Connexion à un sink

Constantes

const BLOCK_ID_TIMEOUT_SECS: u64 = 3;           // Timeout attente nouveau block
const RECENT_BLOCKS_CACHE_SIZE: usize = 10;     // Taille cache anti-redondance

Dépendances

  • pmoaudio : Pipeline audio, types AudioChunk/AudioSegment
  • pmoflac : Décodage FLAC streaming
  • pmometadata : Métadonnées pistes
  • futures-util : StreamExt pour le décodage
  • tokio : Runtime async
  • tokio-util : StreamReader, CancellationToken

Feature gate

[features]
pmoaudio = ["dep:pmoaudio", "dep:pmoflac", "dep:pmometadata", "dep:futures-util"]

Activer avec : cargo build -p pmoparadise --features pmoaudio